Add useKeyOrderedProcessing option to DefaultReactivePulsarMessageListenerContainer (#205)
This commit is contained in:
committed by
GitHub
parent
e91e8e33d5
commit
fb3973785f
@@ -24,6 +24,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumer;
|
||||
import org.apache.pulsar.reactive.client.api.ReactiveMessagePipeline;
|
||||
import org.apache.pulsar.reactive.client.api.ReactiveMessagePipelineBuilder;
|
||||
import org.apache.pulsar.reactive.client.api.ReactiveMessagePipelineBuilder.ConcurrentOneByOneMessagePipelineBuilder;
|
||||
import org.apache.pulsar.reactive.client.internal.api.ApiImplementationFactory;
|
||||
|
||||
import org.springframework.core.log.LogAccessor;
|
||||
@@ -179,8 +180,13 @@ public non-sealed class DefaultReactivePulsarMessageListenerContainer<T>
|
||||
.messageHandler(((ReactivePulsarOneByOneMessageHandler<T>) messageHandler)::received)
|
||||
.handlingTimeout(containerProperties.getHandlingTimeout());
|
||||
if (containerProperties.getConcurrency() > 0) {
|
||||
pipeline = messagePipelineBuilder.concurrent().concurrency(containerProperties.getConcurrency())
|
||||
.maxInflight(containerProperties.getMaxInFlight()).build();
|
||||
ConcurrentOneByOneMessagePipelineBuilder<T> concurrentPipelineBuilder = messagePipelineBuilder
|
||||
.concurrent().concurrency(containerProperties.getConcurrency())
|
||||
.maxInflight(containerProperties.getMaxInFlight());
|
||||
if (containerProperties.isUseKeyOrderedProcessing()) {
|
||||
concurrentPipelineBuilder.useKeyOrderedProcessing();
|
||||
}
|
||||
pipeline = concurrentPipelineBuilder.build();
|
||||
}
|
||||
else {
|
||||
pipeline = pipelineBuilder.build();
|
||||
|
||||
@@ -52,6 +52,8 @@ public class ReactivePulsarContainerProperties<T> {
|
||||
|
||||
private int maxInFlight = 0;
|
||||
|
||||
private boolean useKeyOrderedProcessing = false;
|
||||
|
||||
public ReactivePulsarMessageHandler getMessageHandler() {
|
||||
return this.messageHandler;
|
||||
}
|
||||
@@ -136,4 +138,12 @@ public class ReactivePulsarContainerProperties<T> {
|
||||
this.maxInFlight = maxInFlight;
|
||||
}
|
||||
|
||||
public boolean isUseKeyOrderedProcessing() {
|
||||
return this.useKeyOrderedProcessing;
|
||||
}
|
||||
|
||||
public void setUseKeyOrderedProcessing(boolean useKeyOrderedProcessing) {
|
||||
this.useKeyOrderedProcessing = useKeyOrderedProcessing;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -29,6 +29,7 @@ import org.apache.pulsar.client.api.PulsarClient;
|
||||
import org.apache.pulsar.client.api.Schema;
|
||||
import org.apache.pulsar.client.api.SubscriptionType;
|
||||
import org.apache.pulsar.reactive.client.adapter.AdaptedReactivePulsarClientFactory;
|
||||
import org.apache.pulsar.reactive.client.adapter.DefaultMessageGroupingFunction;
|
||||
import org.apache.pulsar.reactive.client.api.MessageResult;
|
||||
import org.apache.pulsar.reactive.client.api.MutableReactiveMessageConsumerSpec;
|
||||
import org.apache.pulsar.reactive.client.api.MutableReactiveMessageSenderSpec;
|
||||
@@ -143,6 +144,7 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo
|
||||
pulsarContainerProperties.setTopics(List.of(topic));
|
||||
pulsarContainerProperties.setSubscriptionName(subscriptionName);
|
||||
pulsarContainerProperties.setConcurrency(5);
|
||||
pulsarContainerProperties.setUseKeyOrderedProcessing(true);
|
||||
pulsarContainerProperties.setMaxInFlight(6);
|
||||
pulsarContainerProperties.setHandlingTimeout(Duration.ofMillis(7));
|
||||
DefaultReactivePulsarMessageListenerContainer<String> container = new DefaultReactivePulsarMessageListenerContainer<>(
|
||||
@@ -158,7 +160,8 @@ class DefaultReactivePulsarMessageListenerContainerTests implements PulsarTestCo
|
||||
|
||||
assertThat(container).extracting("pipeline", InstanceOfAssertFactories.type(ReactiveMessagePipeline.class))
|
||||
.hasFieldOrPropertyWithValue("concurrency", 5).hasFieldOrPropertyWithValue("maxInflight", 6)
|
||||
.hasFieldOrPropertyWithValue("handlingTimeout", Duration.ofMillis(7));
|
||||
.hasFieldOrPropertyWithValue("handlingTimeout", Duration.ofMillis(7)).extracting("groupingFunction")
|
||||
.isInstanceOf(DefaultMessageGroupingFunction.class);
|
||||
|
||||
container.stop();
|
||||
pulsarClient.close();
|
||||
|
||||
Reference in New Issue
Block a user