From 3cb57e96ebe56b74f9551ea549b66b8cce6c1e23 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Tue, 2 Jun 2015 13:24:01 -0400 Subject: [PATCH] INTEXT-170 Use a custom MessageHeaders class - revert to using a custom KafkaMessageHeaders class directly; - avoid the use of MutableMessageHeaders and the subsequent assignment of ids and timestamps within a GenericMessageTemplate --- .../KafkaMessageDrivenChannelAdapter.java | 39 +++++++++++-------- ...essageDrivenChannelAdapterParserTests.java | 6 +-- 2 files changed, 25 insertions(+), 20 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java index a3dc4c5886..e044e30887 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java @@ -33,10 +33,8 @@ import org.springframework.integration.support.DefaultMessageBuilderFactory; import org.springframework.integration.support.MessageBuilderFactory; import org.springframework.integration.support.MutableMessageBuilderFactory; import org.springframework.messaging.Message; -import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.MessageBuilder; -import org.springframework.messaging.support.MessageHeaderAccessor; import org.springframework.util.Assert; /** @@ -194,32 +192,41 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSupport imp private Message toMessage(Object key, Object payload, KafkaMessageMetadata metadata, Acknowledgment acknowledgment) { - final MessageHeaderAccessor headerAccessor = new MessageHeaderAccessor(); - - headerAccessor.setHeader(KafkaHeaders.MESSAGE_KEY, key); - headerAccessor.setHeader(KafkaHeaders.TOPIC, metadata.getPartition().getTopic()); - headerAccessor.setHeader(KafkaHeaders.PARTITION_ID, metadata.getPartition().getId()); - headerAccessor.setHeader(KafkaHeaders.OFFSET, metadata.getOffset()); - headerAccessor.setHeader(KafkaHeaders.NEXT_OFFSET, metadata.getNextOffset()); - - // pre-set the message id header if set to not generate - headerAccessor.setLeaveMutable(!(this.generateMessageId || this.generateTimestamp)); + KafkaMessageHeaders kafkaMessageHeaders = new KafkaMessageHeaders(generateMessageId, generateTimestamp); + Map rawHeaders = kafkaMessageHeaders.getRawHeaders(); + rawHeaders.put(KafkaHeaders.MESSAGE_KEY, key); + rawHeaders.put(KafkaHeaders.TOPIC, metadata.getPartition().getTopic()); + rawHeaders.put(KafkaHeaders.PARTITION_ID, metadata.getPartition().getId()); + rawHeaders.put(KafkaHeaders.OFFSET, metadata.getOffset()); + rawHeaders.put(KafkaHeaders.NEXT_OFFSET, metadata.getNextOffset()); if (!this.autoCommitOffset) { - headerAccessor.setHeader(KafkaHeaders.ACKNOWLEDGMENT, acknowledgment); + rawHeaders.put(KafkaHeaders.ACKNOWLEDGMENT, acknowledgment); } if (this.useMessageBuilderFactory) { return getMessageBuilderFactory() .withPayload(payload) - .copyHeaders(headerAccessor.toMessageHeaders()) + .copyHeaders(kafkaMessageHeaders) .build(); } else { - return MessageBuilder.createMessage(payload, headerAccessor.getMessageHeaders()); + return MessageBuilder.createMessage(payload, kafkaMessageHeaders); + } + } + + @SuppressWarnings("serial") + private static class KafkaMessageHeaders extends MessageHeaders { + + public KafkaMessageHeaders(boolean generateId, boolean generateTimestamp) { + super(null, generateId ? null : ID_VALUE_NONE, generateTimestamp ? null : -1L); + } + + @Override + public Map getRawHeaders() { + return super.getRawHeaders(); } } - } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java index 469396ea68..4e61eb2e83 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java @@ -135,7 +135,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests { public void doWith(Method method) throws IllegalArgumentException, IllegalAccessException { method.setAccessible(true); toMessage.set(method); -; } + } }, new MethodFilter() { @@ -167,9 +167,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests { m = getAMessageFrom(this.withOverrideIdTS, toMessage.get()); assertNotNull(m.getHeaders().getId()); -//TODO org.springframework.messaging.support.MessageBuilder doesn't support simple way -// to provide TIMESTAMP generation option. -// assertNotNull(m.getHeaders().getTimestamp()); + assertNotNull(m.getHeaders().getTimestamp()); assertNull(m.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT)); assertRest(m); }