GH-130: Support KafkaNull Payloads
Fixes #130 Requires spring-kafka 1.0.3. Upgrade to `spring-kafka-1.0.3.RELEASE`
This commit is contained in:
committed by
Artem Bilan
parent
cda52fe966
commit
a4ab7e54b1
@@ -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<K, V> 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) {
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user