diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java index e532348f..b0163544 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveProperties.java @@ -30,14 +30,13 @@ import org.apache.pulsar.client.api.CompressionType; import org.apache.pulsar.client.api.ConsumerCryptoFailureAction; import org.apache.pulsar.client.api.DeadLetterPolicy; import org.apache.pulsar.client.api.HashingScheme; -import org.apache.pulsar.client.api.KeySharedMode; -import org.apache.pulsar.client.api.KeySharedPolicy; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.ProducerAccessMode; import org.apache.pulsar.client.api.ProducerCryptoFailureAction; import org.apache.pulsar.client.api.Range; import org.apache.pulsar.client.api.RegexSubscriptionMode; import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionMode; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.reactive.client.api.ImmutableReactiveMessageConsumerSpec; @@ -570,16 +569,16 @@ public class PulsarReactiveProperties { */ private SubscriptionType subscriptionType = SubscriptionType.Exclusive; - /* - * KeyShared mode of KeyShared subscription. - */ - private KeySharedMode keySharedMode; - /** * Map of properties to add to the subscription. */ private SortedMap subscriptionProperties = new TreeMap<>(); + /** + * Subscription mode to be used when subscribing to the topic. + */ + private SubscriptionMode subscriptionMode = SubscriptionMode.Durable; + /** * Number of messages that can be accumulated before the consumer calls "receive". */ @@ -616,7 +615,7 @@ public class PulsarReactiveProperties { /** * Whether the retry letter topic is enabled. */ - private Boolean retryLetterTopicEnable; + private Boolean retryLetterTopicEnable = false; /** * Maximum number of messages that a consumer can be pushed at once from a broker @@ -687,7 +686,7 @@ public class PulsarReactiveProperties { */ private Boolean autoUpdatePartitions = true; - private Duration autoUpdatePartitionsInterval; + private Duration autoUpdatePartitionsInterval = Duration.ofMinutes(1); /** * Whether to replicate subscription state. @@ -743,14 +742,6 @@ public class PulsarReactiveProperties { this.subscriptionType = subscriptionType; } - public KeySharedMode getKeySharedMode() { - return this.keySharedMode; - } - - public void setKeySharedMode(KeySharedMode keySharedMode) { - this.keySharedMode = keySharedMode; - } - public SortedMap getSubscriptionProperties() { return this.subscriptionProperties; } @@ -759,6 +750,14 @@ public class PulsarReactiveProperties { this.subscriptionProperties = subscriptionProperties; } + public SubscriptionMode getSubscriptionMode() { + return this.subscriptionMode; + } + + public void setSubscriptionMode(SubscriptionMode subscriptionMode) { + this.subscriptionMode = subscriptionMode; + } + public Integer getReceiverQueueSize() { return this.receiverQueueSize; } @@ -969,11 +968,8 @@ public class PulsarReactiveProperties { map.from(this::getTopicsPattern).to(spec::setTopicsPattern); map.from(this::getSubscriptionName).to(spec::setSubscriptionName); map.from(this::getSubscriptionType).to(spec::setSubscriptionType); - map.from(this::getKeySharedMode).as((mode) -> switch (mode) { - case STICKY -> KeySharedPolicy.stickyHashRange(); - case AUTO_SPLIT -> KeySharedPolicy.autoSplitHashRange(); - }).to(spec::setKeySharedPolicy); map.from(this::getSubscriptionProperties).to(spec::setSubscriptionProperties); + map.from(this::getSubscriptionMode).to(spec::setSubscriptionMode); map.from(this::getReceiverQueueSize).to(spec::setReceiverQueueSize); map.from(this::getAcknowledgementsGroupTime).to(spec::setAcknowledgementsGroupTime); map.from(this::getAcknowledgeAsynchronously).to(spec::setAcknowledgeAsynchronously); @@ -996,6 +992,7 @@ public class PulsarReactiveProperties { map.from(this::getProperties).to(spec::setProperties); map.from(this::getReadCompacted).to(spec::setReadCompacted); map.from(this::getBatchIndexAckEnabled).to(spec::setBatchIndexAckEnabled); + map.from(this::getSubscriptionInitialPosition).to(spec::setSubscriptionInitialPosition); map.from(this::getTopicsPatternAutoDiscoveryPeriod).to(spec::setTopicsPatternAutoDiscoveryPeriod); map.from(this::getTopicsPatternSubscriptionMode).to(spec::setTopicsPatternSubscriptionMode); map.from(this::getAutoUpdatePartitions).to(spec::setAutoUpdatePartitions); diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java index c8c99fa2..f0d08217 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java @@ -26,11 +26,12 @@ import java.util.Map; import org.apache.pulsar.client.api.CompressionType; import org.apache.pulsar.client.api.ConsumerCryptoFailureAction; import org.apache.pulsar.client.api.HashingScheme; -import org.apache.pulsar.client.api.KeySharedMode; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.ProducerAccessMode; import org.apache.pulsar.client.api.ProducerCryptoFailureAction; import org.apache.pulsar.client.api.RegexSubscriptionMode; +import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionMode; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; @@ -130,6 +131,7 @@ public class PulsarReactivePropertiesTests { props.put("spring.pulsar.reactive.consumer.topics-pattern", "my-pattern"); props.put("spring.pulsar.reactive.consumer.subscription-name", "my-subscription"); props.put("spring.pulsar.reactive.consumer.subscription-type", "Shared"); + props.put("spring.pulsar.reactive.consumer.subscription-mode", "NonDurable"); props.put("spring.pulsar.reactive.consumer.subscription-properties[my-sub-prop]", "my-sub-prop-value"); props.put("spring.pulsar.reactive.consumer.receiver-queue-size", "1"); props.put("spring.pulsar.reactive.consumer.acknowledgements-group-time", "2s"); @@ -149,6 +151,7 @@ public class PulsarReactivePropertiesTests { props.put("spring.pulsar.reactive.consumer.properties[my-prop]", "my-prop-value"); props.put("spring.pulsar.reactive.consumer.read-compacted", "true"); props.put("spring.pulsar.reactive.consumer.batch-index-ack-enabled", "true"); + props.put("spring.pulsar.reactive.consumer.subscription-initial-position", "Earliest"); props.put("spring.pulsar.reactive.consumer.topics-pattern-auto-discovery-period", "9s"); props.put("spring.pulsar.reactive.consumer.topics-pattern-subscription-mode", "AllTopics"); props.put("spring.pulsar.reactive.consumer.auto-update-partitions", "false"); @@ -165,6 +168,7 @@ public class PulsarReactivePropertiesTests { assertThat(consumerSpec.getTopicsPattern().toString()).isEqualTo("my-pattern"); assertThat(consumerSpec.getSubscriptionName()).isEqualTo("my-subscription"); assertThat(consumerSpec.getSubscriptionType()).isEqualTo(SubscriptionType.Shared); + assertThat(consumerSpec.getSubscriptionMode()).isEqualTo(SubscriptionMode.NonDurable); assertThat(consumerSpec.getSubscriptionProperties()).hasSize(1).containsEntry("my-sub-prop", "my-sub-prop-value"); assertThat(consumerSpec.getReceiverQueueSize()).isEqualTo(1); @@ -185,6 +189,7 @@ public class PulsarReactivePropertiesTests { assertThat(consumerSpec.getProperties()).hasSize(1).containsEntry("my-prop", "my-prop-value"); assertThat(consumerSpec.getReadCompacted()).isTrue(); assertThat(consumerSpec.getBatchIndexAckEnabled()).isTrue(); + assertThat(consumerSpec.getSubscriptionInitialPosition()).isEqualTo(SubscriptionInitialPosition.Earliest); assertThat(consumerSpec.getTopicsPatternAutoDiscoveryPeriod()).isEqualTo(Duration.ofSeconds(9)); assertThat(consumerSpec.getTopicsPatternSubscriptionMode()).isEqualTo(RegexSubscriptionMode.AllTopics); assertThat(consumerSpec.getAutoUpdatePartitions()).isFalse(); @@ -195,15 +200,6 @@ public class PulsarReactivePropertiesTests { assertThat(consumerSpec.getExpireTimeOfIncompleteChunkedMessage()).isEqualTo(Duration.ofSeconds(12)); } - @ParameterizedTest - @EnumSource(KeySharedMode.class) - void keySharedModeProperty(KeySharedMode keySharedMode) { - bind("spring.pulsar.reactive.consumer.key-shared-mode", keySharedMode.name()); - ReactiveMessageConsumerSpec consumerSpec = properties.buildReactiveMessageConsumerSpec(); - - assertThat(consumerSpec.getKeySharedPolicy().getKeySharedMode()).isEqualTo(keySharedMode); - } - @ParameterizedTest @EnumSource(value = SchedulerType.class, names = "immediate", mode = Mode.EXCLUDE) void acknowledgeScheduler(SchedulerType acknowledgeSchedulerType) {