From a4ab7e54b1f46ed3a1bcb49433cd48cab26952c3 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 11 Aug 2016 17:31:32 -0400 Subject: [PATCH] GH-130: Support KafkaNull Payloads Fixes #130 Requires spring-kafka 1.0.3. Upgrade to `spring-kafka-1.0.3.RELEASE` --- .../outbound/KafkaProducerMessageHandler.java | 13 +++++++++---- .../kafka/inbound/MessageDrivenAdapterTests.java | 14 ++++++++++++++ .../outbound/KafkaProducerMessageHandlerTests.java | 13 +++++++++++++ 3 files changed, 36 insertions(+), 4 deletions(-) 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 d902b34ac7..ba0d09a4c1 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 @@ -26,6 +26,7 @@ import org.springframework.integration.expression.ExpressionUtils; import org.springframework.integration.handler.AbstractMessageHandler; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.kafka.support.KafkaNull; import org.springframework.messaging.Message; import org.springframework.util.Assert; import org.springframework.util.StringUtils; @@ -127,20 +128,24 @@ public class KafkaProducerMessageHandler extends AbstractMessageHandler { ListenableFuture future; + V payload = (V) message.getPayload(); + if (payload instanceof KafkaNull) { + payload = null; + } if (partitionId == null) { if (messageKey == null) { - future = this.kafkaTemplate.send(topic, (V) message.getPayload()); + future = this.kafkaTemplate.send(topic, payload); } else { - future = this.kafkaTemplate.send(topic, (K) messageKey, (V) message.getPayload()); + future = this.kafkaTemplate.send(topic, (K) messageKey, payload); } } else { if (messageKey == null) { - future = this.kafkaTemplate.send(topic, partitionId, (V) message.getPayload()); + future = this.kafkaTemplate.send(topic, partitionId, payload); } else { - future = this.kafkaTemplate.send(topic, partitionId, (K) messageKey, (V) message.getPayload()); + future = this.kafkaTemplate.send(topic, partitionId, (K) messageKey, payload); } } if (this.sync) { diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java index f0ca998a9f..b6e54c48a1 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java @@ -36,6 +36,7 @@ import org.springframework.kafka.listener.KafkaMessageListenerContainer; import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.kafka.support.KafkaNull; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.test.rule.KafkaEmbedded; import org.springframework.kafka.test.utils.KafkaTestUtils; @@ -94,6 +95,19 @@ public class MessageDrivenAdapterTests { assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L); assertThat(headers.get("testHeader")).isEqualTo("testValue"); + template.sendDefault(1, null); + + received = out.receive(10000); + assertThat(received).isNotNull(); + assertThat(received.getPayload()).isInstanceOf(KafkaNull.class); + + headers = received.getHeaders(); + assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1); + assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo("testTopic1"); + assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0); + assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(1L); + assertThat(headers.get("testHeader")).isEqualTo("testValue"); + adapter.stop(); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java index 486b172897..5e42087f05 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/outbound/KafkaProducerMessageHandlerTests.java @@ -36,6 +36,7 @@ import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.kafka.support.KafkaNull; import org.springframework.kafka.test.rule.KafkaEmbedded; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.Message; @@ -102,6 +103,18 @@ public class KafkaProducerMessageHandlerTests { record = KafkaTestUtils.getSingleRecord(consumer, topic1); assertThat(record).has(key((Integer) null)); assertThat(record).has(value("baz")); + + message = MessageBuilder.withPayload(KafkaNull.INSTANCE) + .setHeader(KafkaHeaders.TOPIC, topic1) + .setHeader(KafkaHeaders.MESSAGE_KEY, 2) + .setHeader(KafkaHeaders.PARTITION_ID, 1) + .build(); + handler.handleMessage(message); + + record = KafkaTestUtils.getSingleRecord(consumer, topic1); + assertThat(record).has(key(2)); + assertThat(record).has(partition(1)); + assertThat(record.value()).isNull(); } }