From acc8b1cb929c6f0166e50657c6708f1c7a8b4ef1 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 30 Aug 2023 12:37:47 -0400 Subject: [PATCH] KafkaNull Test Changes There was a regression introduced in Spring Cloud Function where consumers of type Consumer> receive null values when tombstone records are given as KafkaNull. See this issue for more details: https://github.com/spring-cloud/spring-cloud-function/issues/1060 Regression is addressed in Spring Cloud Function and making the corresponding test changes in Spring Cloud Stream Kafka binder. --- .../binder/kafka/integration/KafkaNullConverterTest.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaNullConverterTest.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaNullConverterTest.java index dc7a862b8..166bdf45d 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaNullConverterTest.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaNullConverterTest.java @@ -33,6 +33,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.KafkaNull; import org.springframework.kafka.test.context.EmbeddedKafka; +import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.handler.annotation.Payload; import org.springframework.messaging.support.GenericMessage; @@ -100,9 +101,13 @@ public class KafkaNullConverterTest { } @Bean - public Consumer inputListen() { + public Consumer> inputListen() { return in -> { - this.inputPayload = in; + Object v = in.getPayload(); + String className = v.getClass().getName(); + if (className.equals("org.springframework.kafka.support.KafkaNull")) { + this.inputPayload = null; + } countDownLatchInput.countDown(); }; }