From fd3ef3e8490893db35c1d2af59c7383b2890d6a2 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 16 May 2022 09:23:10 -0400 Subject: [PATCH] GH-2297: Polish Concurrency Test --- .../binder/reactorkafka/ReactorKafkaBinderTests.java | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java index 925669f42..d6c302988 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderTests.java @@ -125,7 +125,8 @@ public class ReactorKafkaBinderTests { binder.setApplicationContext(mock(GenericApplicationContext.class)); CountDownLatch subscriptionLatch = new CountDownLatch(1); - CountDownLatch messageLatch = new CountDownLatch(4); + CountDownLatch messageLatch1 = new CountDownLatch(4); + CountDownLatch messageLatch2 = new CountDownLatch(10); Set partitions = new HashSet<>(); FluxMessageChannel inbound = new FluxMessageChannel(); @@ -133,14 +134,15 @@ public class ReactorKafkaBinderTests { @Override public void onSubscribe(Subscription s) { - s.request(6); + s.request(10); subscriptionLatch.countDown(); } @Override public void onNext(Message msg) { partitions.add(msg.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION, Integer.class)); - messageLatch.countDown(); + messageLatch1.countDown(); + messageLatch2.countDown(); } @Override @@ -169,11 +171,12 @@ public class ReactorKafkaBinderTests { kt.send("testC1", 1, null, "bar").get(10, TimeUnit.SECONDS); kt.send("testC1", 0, null, "baz").get(10, TimeUnit.SECONDS); kt.send("testC1", 1, null, "qux").get(10, TimeUnit.SECONDS); + assertThat(messageLatch1.await(10, TimeUnit.SECONDS)).isTrue(); consumer.stop(); consumer.start(); kt.send("testC1", 0, null, "fiz").get(10, TimeUnit.SECONDS); kt.send("testC1", 1, null, "buz").get(10, TimeUnit.SECONDS); - assertThat(messageLatch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(messageLatch2.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(partitions).hasSize(2); consumer.unbind(); pf.destroy();