GH-3790: Use new header constants for Kafka headers
Fixes https://github.com/spring-projects/spring-integration/issues/3790 Some `KafkaHeaders` constants have been removed and replaced with new more meaningful * Fix removed constants everywhere in the code and docs in favor of newly introduced, which replaces old
This commit is contained in:
@@ -416,12 +416,12 @@ public class KafkaInboundGateway<K, V, R> extends MessagingGatewaySupport
|
||||
}
|
||||
builder.setHeader(KafkaHeaders.TOPIC, requestHeaders.get(KafkaHeaders.REPLY_TOPIC));
|
||||
}
|
||||
if (replyHeaders.get(KafkaHeaders.PARTITION_ID) == null &&
|
||||
if (replyHeaders.get(KafkaHeaders.PARTITION) == null &&
|
||||
requestHeaders.get(KafkaHeaders.REPLY_PARTITION) != null) {
|
||||
if (builder == null) {
|
||||
builder = getMessageBuilderFactory().fromMessage(reply);
|
||||
}
|
||||
builder.setHeader(KafkaHeaders.PARTITION_ID, requestHeaders.get(KafkaHeaders.REPLY_PARTITION));
|
||||
builder.setHeader(KafkaHeaders.PARTITION, requestHeaders.get(KafkaHeaders.REPLY_PARTITION));
|
||||
}
|
||||
if (builder != null) {
|
||||
return builder.build();
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2021 the original author or authors.
|
||||
* Copyright 2013-2022 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -174,7 +174,7 @@ public class KafkaProducerMessageHandler<K, V> extends AbstractReplyProducingMes
|
||||
if (this.isGateway) {
|
||||
setAsync(true);
|
||||
updateNotPropagatedHeaders(
|
||||
new String[]{KafkaHeaders.TOPIC, KafkaHeaders.PARTITION_ID, KafkaHeaders.MESSAGE_KEY}, false);
|
||||
new String[]{KafkaHeaders.TOPIC, KafkaHeaders.PARTITION, KafkaHeaders.KEY}, false);
|
||||
}
|
||||
if (JacksonPresent.isJackson2Present()) {
|
||||
this.headerMapper = new DefaultKafkaHeaderMapper();
|
||||
@@ -565,11 +565,11 @@ public class KafkaProducerMessageHandler<K, V> extends AbstractReplyProducingMes
|
||||
|
||||
Integer partitionId = this.partitionIdExpression != null ?
|
||||
this.partitionIdExpression.getValue(this.evaluationContext, message, Integer.class)
|
||||
: messageHeaders.get(KafkaHeaders.PARTITION_ID, Integer.class);
|
||||
: messageHeaders.get(KafkaHeaders.PARTITION, Integer.class);
|
||||
|
||||
Object messageKey = this.messageKeyExpression != null
|
||||
? this.messageKeyExpression.getValue(this.evaluationContext, message)
|
||||
: messageHeaders.get(KafkaHeaders.MESSAGE_KEY);
|
||||
: messageHeaders.get(KafkaHeaders.KEY);
|
||||
|
||||
Long timestamp = this.timestampExpression != null
|
||||
? this.timestampExpression.getValue(this.evaluationContext, message, Long.class)
|
||||
|
||||
@@ -196,8 +196,8 @@ public class KafkaDslTests {
|
||||
Acknowledgment acknowledgment = headers.get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class);
|
||||
acknowledgment.acknowledge();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(TEST_TOPIC1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(i + 1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(i + 1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo((long) i);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048633L);
|
||||
@@ -213,8 +213,8 @@ public class KafkaDslTests {
|
||||
Acknowledgment acknowledgment = headers.get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class);
|
||||
acknowledgment.acknowledge();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(TEST_TOPIC2);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(i + 1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(i + 1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo((long) i);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048644L);
|
||||
@@ -313,7 +313,7 @@ public class KafkaDslTests {
|
||||
.onPartitionsAssignedSeekCallback((map, callback) ->
|
||||
ContextConfiguration.this.onPartitionsAssignedCalledLatch.countDown()))
|
||||
.filter(Message.class, m ->
|
||||
m.getHeaders().get(KafkaHeaders.RECEIVED_MESSAGE_KEY, Integer.class) < 101,
|
||||
m.getHeaders().get(KafkaHeaders.RECEIVED_KEY, Integer.class) < 101,
|
||||
f -> f.throwExceptionOnRejection(true))
|
||||
.<String, String>transform(String::toUpperCase)
|
||||
.channel(c -> c.queue("listeningFromKafkaResults1"))
|
||||
@@ -336,7 +336,7 @@ public class KafkaDslTests {
|
||||
.messageDrivenChannelAdapter(kafkaListenerContainerFactory().createContainer(TEST_TOPIC2),
|
||||
KafkaMessageDrivenChannelAdapter.ListenerMode.record))
|
||||
.filter(Message.class, m ->
|
||||
m.getHeaders().get(KafkaHeaders.RECEIVED_MESSAGE_KEY, Integer.class) < 101,
|
||||
m.getHeaders().get(KafkaHeaders.RECEIVED_KEY, Integer.class) < 101,
|
||||
f -> f.throwExceptionOnRejection(true))
|
||||
.<String, String>transform(String::toUpperCase)
|
||||
.channel(c -> c.queue("listeningFromKafkaResults2"))
|
||||
|
||||
@@ -166,9 +166,9 @@ class InboundGatewayTests {
|
||||
assertThat(received).isNotNull();
|
||||
|
||||
MessageHeaders headers = received.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048607L);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
@@ -254,9 +254,9 @@ class InboundGatewayTests {
|
||||
MessageHeaders headers = failed.getHeaders();
|
||||
reply.send(MessageBuilder.withPayload("ERROR").copyHeaders(headers).build());
|
||||
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic3);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048607L);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
@@ -338,9 +338,9 @@ class InboundGatewayTests {
|
||||
MessageHeaders headers = failed.getHeaders();
|
||||
reply.send(MessageBuilder.withPayload("ERROR").copyHeaders(headers).build());
|
||||
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic5);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048607L);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
|
||||
@@ -167,9 +167,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(received).isNotNull();
|
||||
|
||||
MessageHeaders headers = received.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048607L);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
@@ -183,9 +183,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(received.getPayload()).isInstanceOf(KafkaNull.class);
|
||||
|
||||
headers = received.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(1L);
|
||||
assertThat((Long) headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isGreaterThan(0L);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
@@ -254,7 +254,7 @@ class MessageDrivenAdapterTests {
|
||||
ProducerFactory<Integer, String> pf = new DefaultKafkaProducerFactory<>(senderProps);
|
||||
KafkaTemplate<Integer, String> template = new KafkaTemplate<>(pf);
|
||||
template.setDefaultTopic(topic4);
|
||||
Message<?> msg = MessageBuilder.withPayload("foo").setHeader(KafkaHeaders.MESSAGE_KEY, 1).build();
|
||||
Message<?> msg = MessageBuilder.withPayload("foo").setHeader(KafkaHeaders.KEY, 1).build();
|
||||
NullChannel component = new NullChannel();
|
||||
component.setBeanName("myNullChannel");
|
||||
msg = MessageHistory.write(msg, component);
|
||||
@@ -268,9 +268,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(originalMessage).isNotNull();
|
||||
assertThat(originalMessage.getHeaders().get(IntegrationMessageHeaderAccessor.SOURCE_DATA)).isNull();
|
||||
headers = originalMessage.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic4);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat(StaticMessageHeaderAccessor.getDeliveryAttempt(originalMessage).get()).isEqualTo(2);
|
||||
|
||||
@@ -381,9 +381,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(originalMessage.getHeaders().get(IntegrationMessageHeaderAccessor.SOURCE_DATA))
|
||||
.isSameAs(headers.get(KafkaHeaders.RAW_DATA));
|
||||
headers = originalMessage.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic5);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat(StaticMessageHeaderAccessor.getDeliveryAttempt(originalMessage).get()).isEqualTo(1);
|
||||
|
||||
@@ -438,9 +438,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(list.size()).isGreaterThan(0);
|
||||
|
||||
MessageHeaders headers = received.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(Arrays.asList(1, 1));
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(Arrays.asList(1, 1));
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(Arrays.asList("testTopic2", "testTopic2"));
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(Arrays.asList(0, 0));
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(Arrays.asList(0, 0));
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(Arrays.asList(0L, 1L));
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE))
|
||||
.isEqualTo(Arrays.asList("CREATE_TIME", "CREATE_TIME"));
|
||||
@@ -507,9 +507,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(received).isNotNull();
|
||||
|
||||
MessageHeaders headers = received.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic3);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isEqualTo(1487694048607L);
|
||||
@@ -554,9 +554,9 @@ class MessageDrivenAdapterTests {
|
||||
assertThat(received).isNotNull();
|
||||
|
||||
MessageHeaders headers = received.getHeaders();
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_KEY)).isEqualTo(1);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic6);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION)).isEqualTo(0);
|
||||
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
|
||||
assertThat((Long) headers.get(KafkaHeaders.RECEIVED_TIMESTAMP)).isGreaterThan(0L);
|
||||
assertThat(headers.get(KafkaHeaders.TIMESTAMP_TYPE)).isEqualTo("CREATE_TIME");
|
||||
|
||||
@@ -165,8 +165,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
Message<?> message = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic1)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.build();
|
||||
handler.handleMessage(message);
|
||||
|
||||
@@ -177,7 +177,7 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
message = MessageBuilder.withPayload("bar")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic1)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 0)
|
||||
.setHeader(KafkaHeaders.PARTITION, 0)
|
||||
.build();
|
||||
handler.handleMessage(message);
|
||||
record = KafkaTestUtils.getSingleRecord(consumer, topic1);
|
||||
@@ -197,8 +197,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
message = MessageBuilder.withPayload(KafkaNull.INSTANCE)
|
||||
.setHeader(KafkaHeaders.TOPIC, topic1)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, "1")
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, "1")
|
||||
.build();
|
||||
handler.handleMessage(message);
|
||||
|
||||
@@ -221,8 +221,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
Message<?> message = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic2)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.setHeader(KafkaHeaders.TIMESTAMP, 1487694048607L)
|
||||
.setHeader("baz", "qux")
|
||||
.build();
|
||||
@@ -252,8 +252,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
Message<?> message = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic3)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.build();
|
||||
|
||||
handler.setTimestampExpression(new ValueExpression<>(1487694048633L));
|
||||
@@ -293,8 +293,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
Message<?> message = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic4)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.build();
|
||||
handler.handleMessage(message);
|
||||
|
||||
@@ -326,7 +326,7 @@ class KafkaProducerMessageHandlerTests {
|
||||
handler.afterPropertiesSet();
|
||||
message = MessageBuilder.withPayload("bar")
|
||||
.setHeader(KafkaHeaders.TOPIC, "foo")
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 0)
|
||||
.setHeader(KafkaHeaders.PARTITION, 0)
|
||||
.build();
|
||||
handler.handleMessage(message);
|
||||
|
||||
@@ -409,8 +409,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
if (payload == null) {
|
||||
message = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic5)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.build();
|
||||
}
|
||||
else {
|
||||
@@ -435,8 +435,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
final Message<?> messageToHandle1 = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic5)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.setHeader(KafkaHeaders.REPLY_TOPIC, "bad")
|
||||
.build();
|
||||
|
||||
@@ -447,8 +447,8 @@ class KafkaProducerMessageHandlerTests {
|
||||
|
||||
final Message<?> messageToHandle2 = MessageBuilder.withPayload("foo")
|
||||
.setHeader(KafkaHeaders.TOPIC, topic5)
|
||||
.setHeader(KafkaHeaders.MESSAGE_KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION_ID, 1)
|
||||
.setHeader(KafkaHeaders.KEY, 2)
|
||||
.setHeader(KafkaHeaders.PARTITION, 1)
|
||||
.setHeader(KafkaHeaders.REPLY_PARTITION, 999)
|
||||
.build();
|
||||
|
||||
|
||||
@@ -151,8 +151,8 @@ class KafkaDslKotlinTests {
|
||||
val acknowledgment = headers[KafkaHeaders.ACKNOWLEDGMENT] as Acknowledgment
|
||||
acknowledgment.acknowledge()
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_TOPIC]).isEqualTo(TEST_TOPIC1)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_MESSAGE_KEY]).isEqualTo(i + 1)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_PARTITION_ID]).isEqualTo(0)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_KEY]).isEqualTo(i + 1)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_PARTITION]).isEqualTo(0)
|
||||
assertThat(headers[KafkaHeaders.OFFSET]).isEqualTo(i.toLong())
|
||||
assertThat(headers[KafkaHeaders.TIMESTAMP_TYPE]).isEqualTo("CREATE_TIME")
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_TIMESTAMP]).isEqualTo(1487694048633L)
|
||||
@@ -168,8 +168,8 @@ class KafkaDslKotlinTests {
|
||||
val acknowledgment = headers[KafkaHeaders.ACKNOWLEDGMENT] as Acknowledgment
|
||||
acknowledgment.acknowledge()
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_TOPIC]).isEqualTo(TEST_TOPIC2)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_MESSAGE_KEY]).isEqualTo(i + 1)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_PARTITION_ID]).isEqualTo(0)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_KEY]).isEqualTo(i + 1)
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_PARTITION]).isEqualTo(0)
|
||||
assertThat(headers[KafkaHeaders.OFFSET]).isEqualTo(i.toLong())
|
||||
assertThat(headers[KafkaHeaders.TIMESTAMP_TYPE]).isEqualTo("CREATE_TIME")
|
||||
assertThat(headers[KafkaHeaders.RECEIVED_TIMESTAMP]).isEqualTo(1487694048644L)
|
||||
@@ -240,7 +240,7 @@ class KafkaDslKotlinTests {
|
||||
.errorChannel(errorChannel())
|
||||
.retryTemplate(RetryTemplate())
|
||||
.filterInRetry(true)) {
|
||||
filter<Message<*>>({ m -> (m.headers[KafkaHeaders.RECEIVED_MESSAGE_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) }
|
||||
filter<Message<*>>({ m -> (m.headers[KafkaHeaders.RECEIVED_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) }
|
||||
transform<String> { it.uppercase() }
|
||||
channel { queue("listeningFromKafkaResults1") }
|
||||
}
|
||||
@@ -254,7 +254,7 @@ class KafkaDslKotlinTests {
|
||||
.errorChannel(errorChannel())
|
||||
.retryTemplate(RetryTemplate())
|
||||
.filterInRetry(true)) {
|
||||
filter<Message<*>>({ m -> (m.headers[KafkaHeaders.RECEIVED_MESSAGE_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) }
|
||||
filter<Message<*>>({ m -> (m.headers[KafkaHeaders.RECEIVED_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) }
|
||||
transform<String> { it.uppercase() }
|
||||
channel { queue("listeningFromKafkaResults2") }
|
||||
}
|
||||
|
||||
@@ -209,7 +209,7 @@ Also, the `mode` attribute is available.
|
||||
It can accept values of `record` or `batch` (default: `record`).
|
||||
For `record` mode, each message payload is converted from a single `ConsumerRecord`.
|
||||
For `batch` mode, the payload is a list of objects that are converted from all the `ConsumerRecord` instances returned by the consumer poll.
|
||||
As with the batched `@KafkaListener`, the `KafkaHeaders.RECEIVED_MESSAGE_KEY`, `KafkaHeaders.RECEIVED_PARTITION_ID`, `KafkaHeaders.RECEIVED_TOPIC`, and `KafkaHeaders.OFFSET` headers are also lists, with positions corresponding to the position in the payload.
|
||||
As with the batched `@KafkaListener`, the `KafkaHeaders.RECEIVED_KEY`, `KafkaHeaders.RECEIVED_PARTITION`, `KafkaHeaders.RECEIVED_TOPIC`, and `KafkaHeaders.OFFSET` headers are also lists, with positions corresponding to the position in the payload.
|
||||
|
||||
Received messages have certain headers populated.
|
||||
See the https://docs.spring.io/spring-kafka/api/org/springframework/kafka/support/KafkaHeaders.html[`KafkaHeaders` class] for more information.
|
||||
@@ -779,7 +779,7 @@ The following example shows how to do so:
|
||||
[source, java]
|
||||
----
|
||||
@ServiceActivator(inputChannel = "fromSomeKafkaInboundEndpoint")
|
||||
public void in(@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key,
|
||||
public void in(@Header(KafkaHeaders.RECEIVED_KEY) String key,
|
||||
@Payload(required = false) Customer customer) {
|
||||
// customer is null if a tombstone record
|
||||
...
|
||||
|
||||
Reference in New Issue
Block a user