From 34c4512dd8bb3e134fdee961c609d896d0f42707 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 31 Aug 2022 16:36:29 -0400 Subject: [PATCH] Batch acknowledgment test changes --- .../pulsar/core/ConsumerAcknowledgmentTests.java | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java index e9b99d76..92697395 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java @@ -139,6 +139,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { @Test void testBatchAckButSomeRecordsFail() throws Exception { + final Object lock = new Object(); Map config = new HashMap<>(); final Set strings = new HashSet<>(); strings.add("cons-ack-tests-013"); @@ -151,9 +152,11 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); CountDownLatch latch = new CountDownLatch(10); pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> { - latch.countDown(); - if (latch.getCount() % 2 == 0) { - throw new RuntimeException("fail"); + synchronized (lock) { + latch.countDown(); + if (latch.getCount() % 2 == 0) { + throw new RuntimeException("fail"); + } } }); pulsarContainerProperties.setSchema(Schema.STRING); @@ -174,9 +177,9 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { Thread.sleep(1_000); // Half of the message get acknowledged, and the other half gets negatively // acknowledged. - await().atMost(Duration.ofSeconds(30)) + await().atMost(Duration.ofSeconds(10)) .untilAsserted(() -> verify(containerConsumer, times(5)).acknowledge(any(Message.class))); - await().atMost(Duration.ofSeconds(30)) + await().atMost(Duration.ofSeconds(10)) .untilAsserted(() -> verify(containerConsumer, times(5)).negativeAcknowledge(any(Message.class))); container.stop(); pulsarClient.close();