diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index add287ac..1db59dcc 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -277,7 +277,8 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS DeadLetterPolicy deadLetterPolicy = DeadLetterPolicy.builder().maxRedeliverCount(1) .deadLetterTopic("dpmlct-016-dlq-topic").build(); Map config = Map.of("topicNames", Collections.singleton("dpmlct-016"), "subscriptionName", - "dpmlct-sb-016", "ackTimeoutMillis", 1, "deadLetterPolicy", deadLetterPolicy); + "dpmlct-sb-016", "negativeAckRedeliveryDelayMicros", TimeUnit.SECONDS.toMicros(1), "deadLetterPolicy", + deadLetterPolicy); PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); @@ -332,7 +333,8 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS DeadLetterPolicy deadLetterPolicy = DeadLetterPolicy.builder().maxRedeliverCount(5).deadLetterTopic("dlq-topic") .build(); Map config = Map.of("topicNames", Collections.singleton("dpmlct-017"), "subscriptionName", - "dpmlct-sb-016", "ackTimeoutMillis", 1, "deadLetterPolicy", deadLetterPolicy); + "dpmlct-sb-016", "negativeAckRedeliveryDelayMicros", TimeUnit.SECONDS.toMicros(1), "deadLetterPolicy", + deadLetterPolicy); PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()) .build(); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java index 05b255b0..22c67204 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java @@ -298,7 +298,7 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @ContextConfiguration(classes = AckTimeoutkRedeliveryBackoffTest.AckTimeoutRedeliveryConfig.class) class AckTimeoutkRedeliveryBackoffTest { - static CountDownLatch ackTimeoutRedeliveryBackoffLatch = new CountDownLatch(5); + static CountDownLatch ackTimeoutRedeliveryBackoffLatch = new CountDownLatch(3); @Test void pulsarListenerWithAckTimeoutRedeliveryBackoff() throws Exception { @@ -310,19 +310,29 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @Configuration static class AckTimeoutRedeliveryConfig { + /** + * The following PulsarListener is for testing the ack timeout settings. We + * set an ack timeout of 1 second in the listener and then will not ack the + * incoming message unless the tracking CountDownLatch count goes down to + * zero. This means that there were enough ack timeouts and corresponding + * redeliveries of the same message. Note that we are doing a manual ack + * because we want to ack the message only when the latch count becomes zero. + */ @PulsarListener(id = "withAckTimeoutRedeliveryBackoff", subscriptionName = "withAckTimeoutRedeliveryBackoffSubscription", topics = "withAckTimeoutRedeliveryBackoff-test-topic", - ackTimeoutRedeliveryBackoff = "ackTimeoutRedeliveryBackoff", + ackTimeoutRedeliveryBackoff = "ackTimeoutRedeliveryBackoff", ackMode = AckMode.MANUAL, subscriptionType = SubscriptionType.Shared, properties = { "ackTimeoutMillis=1000" }) - void listen(String ignored) { + void listen(String ignored, Acknowledgement acknowledgement) { ackTimeoutRedeliveryBackoffLatch.countDown(); - throw new RuntimeException(); + if (ackTimeoutRedeliveryBackoffLatch.getCount() == 0) { + acknowledgement.acknowledge(); + } } @Bean public RedeliveryBackoff ackTimeoutRedeliveryBackoff() { - return MultiplierRedeliveryBackoff.builder().minDelayMs(1000).maxDelayMs(5 * 1000).multiplier(2) + return MultiplierRedeliveryBackoff.builder().minDelayMs(1000).maxDelayMs(3 * 1000).multiplier(2) .build(); } @@ -394,7 +404,8 @@ public class PulsarListenerTests implements PulsarTestContainerSupport { @PulsarListener(id = "deadLetterPolicyListener", subscriptionName = "deadLetterPolicySubscription", topics = "dlpt-topic-1", deadLetterPolicy = "deadLetterPolicy", - subscriptionType = SubscriptionType.Shared, properties = { "ackTimeoutMillis=1000" }) + subscriptionType = SubscriptionType.Shared, + properties = { "negativeAckRedeliveryDelayMicros=1000000" }) void listen(String msg) { latch.countDown(); throw new RuntimeException("fail " + msg);