Add more settings to Consumer in PulsarReactiveProperties (#226)
This commit is contained in:
committed by
GitHub
parent
a332cd9329
commit
d9018ea51c
@@ -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<String, String> 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<String, String> 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);
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user