From fb3973785f4170c7f617ee62261b1b7446fac1e0 Mon Sep 17 00:00:00 2001 From: Christophe Bornet Date: Wed, 16 Nov 2022 21:59:05 +0100 Subject: [PATCH] Add useKeyOrderedProcessing option to DefaultReactivePulsarMessageListenerContainer (#205) --- .../DefaultReactivePulsarMessageListenerContainer.java | 10 ++++++++-- .../reactive/ReactivePulsarContainerProperties.java | 10 ++++++++++ ...ultReactivePulsarMessageListenerContainerTests.java | 5 ++++- 3 files changed, 22 insertions(+), 3 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainer.java index 5415e9e6..c8ffa41e 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainer.java @@ -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 .messageHandler(((ReactivePulsarOneByOneMessageHandler) messageHandler)::received) .handlingTimeout(containerProperties.getHandlingTimeout()); if (containerProperties.getConcurrency() > 0) { - pipeline = messagePipelineBuilder.concurrent().concurrency(containerProperties.getConcurrency()) - .maxInflight(containerProperties.getMaxInFlight()).build(); + ConcurrentOneByOneMessagePipelineBuilder concurrentPipelineBuilder = messagePipelineBuilder + .concurrent().concurrency(containerProperties.getConcurrency()) + .maxInflight(containerProperties.getMaxInFlight()); + if (containerProperties.isUseKeyOrderedProcessing()) { + concurrentPipelineBuilder.useKeyOrderedProcessing(); + } + pipeline = concurrentPipelineBuilder.build(); } else { pipeline = pipelineBuilder.build(); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/ReactivePulsarContainerProperties.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/ReactivePulsarContainerProperties.java index 84125a81..e8e1c7a0 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/ReactivePulsarContainerProperties.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/reactive/ReactivePulsarContainerProperties.java @@ -52,6 +52,8 @@ public class ReactivePulsarContainerProperties { private int maxInFlight = 0; + private boolean useKeyOrderedProcessing = false; + public ReactivePulsarMessageHandler getMessageHandler() { return this.messageHandler; } @@ -136,4 +138,12 @@ public class ReactivePulsarContainerProperties { this.maxInFlight = maxInFlight; } + public boolean isUseKeyOrderedProcessing() { + return this.useKeyOrderedProcessing; + } + + public void setUseKeyOrderedProcessing(boolean useKeyOrderedProcessing) { + this.useKeyOrderedProcessing = useKeyOrderedProcessing; + } + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainerTests.java index 033b2068..0c3f3e49 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/reactive/DefaultReactivePulsarMessageListenerContainerTests.java @@ -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 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();