From ee3279072dbc7652896486c6b9b9b210ed076afa Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 22 May 2020 16:11:18 -0400 Subject: [PATCH] Fixing a random test failure on CI --- ...serializationErrorHandlerByKafkaTests.java | 25 +++++++++++-------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java index eda3273f4..a70ebb9ee 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java @@ -66,14 +66,15 @@ import static org.mockito.Mockito.verify; @RunWith(SpringRunner.class) @ContextConfiguration @DirtiesContext -@Ignore public abstract class DeserializationErrorHandlerByKafkaTests { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, - "DeserializationErrorHandlerByKafkaTests-In", + "abc-DeserializationErrorHandlerByKafkaTests-In", + "xyz-DeserializationErrorHandlerByKafkaTests-In", "DeserializationErrorHandlerByKafkaTests-out", - "error.DeserializationErrorHandlerByKafkaTests-In.group", + "error.abc-DeserializationErrorHandlerByKafkaTests-In.group", + "error.xyz-DeserializationErrorHandlerByKafkaTests-In.group", "error.word1.groupx", "error.word2.groupx"); @@ -99,7 +100,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "DeserializationErrorHandlerByKafkaTests-out"); + embeddedKafka.consumeFromEmbeddedTopics(consumer, "DeserializationErrorHandlerByKafkaTests-out", "DeserializationErrorHandlerByKafkaTests-out"); } @AfterClass @@ -111,6 +112,8 @@ public abstract class DeserializationErrorHandlerByKafkaTests { } @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.destination=abc-DeserializationErrorHandlerByKafkaTests-In", + "spring.cloud.stream.bindings.output.destination=DeserializationErrorHandlerByKafkaTests-Out", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq", "spring.cloud.stream.bindings.input.group=group", "spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler=sendToDlq", @@ -125,7 +128,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("DeserializationErrorHandlerByKafkaTests-In"); + template.setDefaultTopic("abc-DeserializationErrorHandlerByKafkaTests-In"); template.sendDefault(1, null, "foobar"); Map consumerProps = KafkaTestUtils.consumerProps("foobar", @@ -134,10 +137,10 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); Consumer consumer1 = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.DeserializationErrorHandlerByKafkaTests-In.group"); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.abc-DeserializationErrorHandlerByKafkaTests-In.group"); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer1, - "error.DeserializationErrorHandlerByKafkaTests-In.group"); + "error.abc-DeserializationErrorHandlerByKafkaTests-In.group"); assertThat(cr.value()).isEqualTo("foobar"); assertThat(cr.partition()).isEqualTo(0); // custom partition function @@ -150,6 +153,8 @@ public abstract class DeserializationErrorHandlerByKafkaTests { } @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.destination=xyz-DeserializationErrorHandlerByKafkaTests-In", + "spring.cloud.stream.bindings.output.destination=DeserializationErrorHandlerByKafkaTests-Out", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq", "spring.cloud.stream.bindings.input.group=group", "spring.cloud.stream.kafka.streams.bindings.input.consumer.deserializationExceptionHandler=sendToDlq", @@ -164,7 +169,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("DeserializationErrorHandlerByKafkaTests-In"); + template.setDefaultTopic("xyz-DeserializationErrorHandlerByKafkaTests-In"); template.sendDefault(1, null, "foobar"); Map consumerProps = KafkaTestUtils.consumerProps("foobar", @@ -173,10 +178,10 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); Consumer consumer1 = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.DeserializationErrorHandlerByKafkaTests-In.group"); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.xyz-DeserializationErrorHandlerByKafkaTests-In.group"); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer1, - "error.DeserializationErrorHandlerByKafkaTests-In.group"); + "error.xyz-DeserializationErrorHandlerByKafkaTests-In.group"); assertThat(cr.value()).isEqualTo("foobar"); assertThat(cr.partition()).isEqualTo(0); // custom partition function