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
This commit is contained in:
committed by
Artem Bilan
parent
35be1cb4fb
commit
3cb57e96eb
@@ -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<Object> 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<String, Object> 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<String, Object> getRawHeaders() {
|
||||
return super.getRawHeaders();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user