From 22b74c7a337d418d59b19cfe16305e064c1e54b7 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 5 May 2022 21:51:16 -0400 Subject: [PATCH] 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 --- .../kafka/inbound/KafkaInboundGateway.java | 4 +-- .../outbound/KafkaProducerMessageHandler.java | 8 ++--- .../integration/kafka/dsl/KafkaDslTests.java | 12 +++---- .../kafka/inbound/InboundGatewayTests.java | 12 +++---- .../inbound/MessageDrivenAdapterTests.java | 30 ++++++++-------- .../KafkaProducerMessageHandlerTests.java | 36 +++++++++---------- .../kafka/dsl/kotlin/KafkaDslKotlinTests.kt | 12 +++---- src/reference/asciidoc/kafka.adoc | 4 +-- 8 files changed, 59 insertions(+), 59 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java index 2fe6072429..74d9081960 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java @@ -416,12 +416,12 @@ public class KafkaInboundGateway 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(); 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 e0d36b12fa..801b55cad1 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 @@ -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 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 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) diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java index e327e6744a..37cc74650e 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java @@ -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)) .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)) .transform(String::toUpperCase) .channel(c -> c.queue("listeningFromKafkaResults2")) diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java index 93cea395b9..aabf054309 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java @@ -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"); 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 cb83fae627..e06a01e7a0 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 @@ -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 pf = new DefaultKafkaProducerFactory<>(senderProps); KafkaTemplate 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"); 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 a8eb7bfa47..a4a9989e4d 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 @@ -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(); diff --git a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt index f1f9ddb25f..ec65d3e946 100644 --- a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt +++ b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/kotlin/KafkaDslKotlinTests.kt @@ -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>({ m -> (m.headers[KafkaHeaders.RECEIVED_MESSAGE_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) } + filter>({ m -> (m.headers[KafkaHeaders.RECEIVED_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) } transform { it.uppercase() } channel { queue("listeningFromKafkaResults1") } } @@ -254,7 +254,7 @@ class KafkaDslKotlinTests { .errorChannel(errorChannel()) .retryTemplate(RetryTemplate()) .filterInRetry(true)) { - filter>({ m -> (m.headers[KafkaHeaders.RECEIVED_MESSAGE_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) } + filter>({ m -> (m.headers[KafkaHeaders.RECEIVED_KEY] as Int) < 101 }) { throwExceptionOnRejection(true) } transform { it.uppercase() } channel { queue("listeningFromKafkaResults2") } } diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index 4f20c616b8..aa53bdd88a 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -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 ...