diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java index 7d1a4eb2c2..f6eaacf233 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandler.java @@ -626,7 +626,8 @@ public class KafkaProducerMessageHandler extends AbstractReplyProducingMes replyTopic = getSingleReplyTopic(); } else { - throw new IllegalStateException("No reply topic header and no default reply topic can be determined"); + throw new IllegalStateException("No reply topic header and no default reply topic can be determined; " + + "container's assigned partitions: " + this.replyTopicsAndPartitions); } } else { diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java index 61810899d8..13eac986e8 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java @@ -459,7 +459,9 @@ public class KafkaDslTests { @Override public void onPartitionsAssigned(Collection partitions) { - ContextConfiguration.this.replyContainerLatch.countDown(); + if (!partitions.isEmpty()) { + ContextConfiguration.this.replyContainerLatch.countDown(); + } } }); diff --git a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt index c71613dbe7..22ef8c2858 100644 --- a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt +++ b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt @@ -345,7 +345,9 @@ class KafkaDslKotlinTests { } override fun onPartitionsAssigned(partitions: Collection) { - this@ContextConfiguration.replyContainerLatch.countDown() + if (!partitions.isEmpty()) { + this@ContextConfiguration.replyContainerLatch.countDown() + } } })