From cbfc6b0e73d0e35d275854ce549d8c24ead49e2c Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 7 Sep 2016 16:49:35 -0400 Subject: [PATCH] GH-141: Support Batch Payloads Resolves: https://github.com/spring-projects/spring-integration-kafka/issues/141 Enum Javadocs Polishing - PR Comments --- ...afkaMessageDrivenChannelAdapterParser.java | 1 + .../KafkaMessageDrivenChannelAdapter.java | 140 ++++++++++++++---- .../main/resources/META-INF/spring.schemas | 4 +- ...0.xsd => spring-integration-kafka-2.1.xsd} | 20 +++ ...rivenChannelAdapterParserTests-context.xml | 28 ++++ ...essageDrivenChannelAdapterParserTests.java | 25 +++- .../inbound/MessageDrivenAdapterTests.java | 66 ++++++++- 7 files changed, 249 insertions(+), 35 deletions(-) rename spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/{spring-integration-kafka-2.0.xsd => spring-integration-kafka-2.1.xsd} (89%) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParser.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParser.java index 202ade8ba6..3bf247cc18 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParser.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParser.java @@ -47,6 +47,7 @@ public class KafkaMessageDrivenChannelAdapterParser extends AbstractChannelAdapt else { parserContext.getReaderContext().error("The 'listener-container' attribute is required.", element); } + builder.addConstructorArgValue(element.getAttribute("mode")); builder.addPropertyReference("outputChannel", channelName); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-timeout"); 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 625f8a5f0f..3859780505 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 @@ -16,17 +16,23 @@ package org.springframework.integration.kafka.inbound; +import java.util.List; + import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.integration.context.OrderlyShutdownCapable; import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.AcknowledgingMessageListener; +import org.springframework.kafka.listener.BatchAcknowledgingMessageListener; +import org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter; import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter; +import org.springframework.kafka.listener.adapter.FilteringBatchAcknowledgingMessageListenerAdapter; import org.springframework.kafka.listener.adapter.RecordFilterStrategy; import org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter; import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter; import org.springframework.kafka.support.Acknowledgment; +import org.springframework.kafka.support.converter.BatchMessageConverter; import org.springframework.kafka.support.converter.MessageConverter; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.messaging.Message; @@ -49,7 +55,11 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo private final AbstractMessageListenerContainer messageListenerContainer; - private final RecordMessagingMessageListenerAdapter listener = new IntegrationMessageListener(); + private final RecordMessagingMessageListenerAdapter recordListener = new IntegrationRecordMessageListener(); + + private final BatchMessagingMessageListenerAdapter batchListener = new IntegrationBatchMessageListener(); + + private final ListenerMode mode; private RecordFilterStrategy recordFilterStrategy; @@ -61,23 +71,47 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo private boolean filterInRetry; + /** + * Construct an instance with mode {@link ListenerMode#record}. + * @param messageListenerContainer the container. + */ public KafkaMessageDrivenChannelAdapter(AbstractMessageListenerContainer messageListenerContainer) { + this(messageListenerContainer, ListenerMode.record); + } + + /** + * Construct an instance with the provided mode. + * @param messageListenerContainer the container. + * @param mode the mode. + * @since 1.2 + */ + public KafkaMessageDrivenChannelAdapter(AbstractMessageListenerContainer messageListenerContainer, + ListenerMode mode) { Assert.notNull(messageListenerContainer, "messageListenerContainer is required"); Assert.isNull(messageListenerContainer.getContainerProperties().getMessageListener(), "Container must not already have a listener"); this.messageListenerContainer = messageListenerContainer; this.messageListenerContainer.setAutoStartup(false); + this.mode = mode; } /** - * Set the message converter; must be a {@link RecordMessageConverter}. + * Set the message converter; must be a {@link RecordMessageConverter} or + * {@link BatchMessageConverter} depending on mode. * @param messageConverter the converter. - * @deprecated in favor of {@link #setRecordMessageConverter(RecordMessageConverter)}. */ - @Deprecated public void setMessageConverter(MessageConverter messageConverter) { - Assert.isInstanceOf(RecordMessageConverter.class, messageConverter); - this.listener.setMessageConverter((RecordMessageConverter) messageConverter); + if (messageConverter instanceof RecordMessageConverter) { + this.recordListener.setMessageConverter((RecordMessageConverter) messageConverter); + } + else if (messageConverter instanceof BatchMessageConverter) { + this.batchListener.setBatchMessageConverter((BatchMessageConverter) messageConverter); + } + else { + throw new IllegalArgumentException( + "Message converter must be a 'RecordMessageConverter' or 'BatchMessageConverter'"); + } + } /** @@ -86,12 +120,21 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo * @since 2.1 */ public void setRecordMessageConverter(RecordMessageConverter messageConverter) { - this.listener.setMessageConverter(messageConverter); + this.recordListener.setMessageConverter(messageConverter); + } + + /** + * Set the message converter to use with a batch-based consumer. + * @param messageConverter the converter. + * @since 2.1 + */ + public void setBatchMessageConverter(BatchMessageConverter messageConverter) { + this.batchListener.setBatchMessageConverter(messageConverter); } /** * Specify a {@link RecordFilterStrategy} to wrap - * {@link KafkaMessageDrivenChannelAdapter.IntegrationMessageListener} into + * {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into * {@link FilteringAcknowledgingMessageListenerAdapter}. * @param recordFilterStrategy the {@link RecordFilterStrategy} to use. * @since 2.0.1 @@ -113,12 +156,14 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo /** * Specify a {@link RetryTemplate} instance to wrap - * {@link KafkaMessageDrivenChannelAdapter.IntegrationMessageListener} into + * {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into * {@link RetryingAcknowledgingMessageListenerAdapter}. * @param retryTemplate the {@link RetryTemplate} to use. * @since 2.0.1 */ public void setRetryTemplate(RetryTemplate retryTemplate) { + Assert.isTrue(retryTemplate == null || this.mode.equals(ListenerMode.record), + "Retry is not supported with mode=batch"); this.retryTemplate = retryTemplate; } @@ -152,28 +197,38 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo protected void onInit() { super.onInit(); - AcknowledgingMessageListener listener = this.listener; + if (this.mode.equals(ListenerMode.record)) { + AcknowledgingMessageListener listener = this.recordListener; - boolean filterInRetry = this.filterInRetry && this.retryTemplate != null && this.recordFilterStrategy != null; + boolean filterInRetry = this.filterInRetry && this.retryTemplate != null && this.recordFilterStrategy != null; - if (filterInRetry) { - listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, - this.ackDiscarded); - listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, - this.recoveryCallback); - } - else { - if (this.retryTemplate != null) { - listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, - this.recoveryCallback); - } - if (this.recordFilterStrategy != null) { + if (filterInRetry) { listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, this.ackDiscarded); + listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, + this.recoveryCallback); } + else { + if (this.retryTemplate != null) { + listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, + this.recoveryCallback); + } + if (this.recordFilterStrategy != null) { + listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + this.ackDiscarded); + } + } + this.messageListenerContainer.getContainerProperties().setMessageListener(listener); } + else { + BatchAcknowledgingMessageListener listener = this.batchListener; - this.messageListenerContainer.getContainerProperties().setMessageListener(listener); + if (this.recordFilterStrategy != null) { + listener = new FilteringBatchAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + this.ackDiscarded); + } + this.messageListenerContainer.getContainerProperties().setMessageListener(listener); + } } @Override @@ -202,9 +257,28 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo return getPhase(); } - private class IntegrationMessageListener extends RecordMessagingMessageListenerAdapter { + /** + * The listener mode for the container, record or batch. + * @since 1.2 + * + */ + public enum ListenerMode { - IntegrationMessageListener() { + /** + * Each {@link Message} will be converted from a single {@code ConsumerRecord}. + */ + record, + + /** + * Each {@link Message} will be converted from the {@code ConsumerRecords} + * returned by a poll. + */ + batch + } + + private class IntegrationRecordMessageListener extends RecordMessagingMessageListenerAdapter { + + IntegrationRecordMessageListener() { super(null, null); } @@ -216,4 +290,18 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } + private class IntegrationBatchMessageListener extends BatchMessagingMessageListenerAdapter { + + IntegrationBatchMessageListener() { + super(null, null); + } + + @Override + public void onMessage(List> records, Acknowledgment acknowledgment) { + Message message = toMessagingMessage(records, acknowledgment); + sendMessage(message); + } + + } + } diff --git a/spring-integration-kafka/src/main/resources/META-INF/spring.schemas b/spring-integration-kafka/src/main/resources/META-INF/spring.schemas index f189965bcc..0806da78d5 100644 --- a/spring-integration-kafka/src/main/resources/META-INF/spring.schemas +++ b/spring-integration-kafka/src/main/resources/META-INF/spring.schemas @@ -1,2 +1,2 @@ -http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka-2.0.xsd=org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd -http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka.xsd=org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd +http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka-2.1.xsd=org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd +http\://www.springframework.org/schema/integration/kafka/spring-integration-kafka.xsd=org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd similarity index 89% rename from spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd rename to spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd index 5294d79a96..4fafa2f868 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.0.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd @@ -155,6 +155,8 @@ An 'org.springframework.kafka.support.converter.MessageConverter' bean reference. + if mode = 'record' must be a 'RecordMessageConverter'; if mode = 'batch' must be + a `BatchMessageConverter`. Defaults to the default implementation for each mode. @@ -163,7 +165,25 @@ + + + + 'record' or 'batch' - default 'record' - one converted ConsumerRecord per message, when + 'batch' then the payload is a collection of converted ConsumerRecords. + + + + + + + + + + + + + diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests-context.xml b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests-context.xml index 0607bf9e66..1da36150c5 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests-context.xml +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests-context.xml @@ -23,6 +23,17 @@ message-converter="messageConverter" error-channel="errorChannel" /> + + @@ -42,4 +53,21 @@ + + + + + + + + + + + + + + + + + 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 781d07f006..4d5dcca6c6 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 @@ -29,6 +29,7 @@ import org.springframework.integration.channel.NullChannel; import org.springframework.integration.channel.PublishSubscribeChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter; +import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.ListenerMode; import org.springframework.integration.test.util.TestUtils; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.listener.KafkaMessageListenerContainer; @@ -57,6 +58,9 @@ public class KafkaMessageDrivenChannelAdapterParserTests { @Autowired private KafkaMessageDrivenChannelAdapter kafkaListener; + @Autowired + private KafkaMessageDrivenChannelAdapter kafkaBatchListener; + @Test public void testKafkaMessageDrivenChannelAdapterParser() throws Exception { assertThat(this.kafkaListener.isAutoStartup()).isFalse(); @@ -68,6 +72,23 @@ public class KafkaMessageDrivenChannelAdapterParserTests { TestUtils.getPropertyValue(this.kafkaListener, "messageListenerContainer", KafkaMessageListenerContainer.class); assertThat(container).isNotNull(); + assertThat(TestUtils.getPropertyValue(kafkaListener, "mode", ListenerMode.class)) + .isEqualTo(ListenerMode.record); + } + + @Test + public void testKafkaBatchMessageDrivenChannelAdapterParser() throws Exception { + assertThat(this.kafkaBatchListener.isAutoStartup()).isFalse(); + assertThat(this.kafkaBatchListener.isRunning()).isFalse(); + assertThat(this.kafkaBatchListener.getPhase()).isEqualTo(100); + assertThat(TestUtils.getPropertyValue(this.kafkaBatchListener, "outputChannel")).isSameAs(this.nullChannel); + assertThat(TestUtils.getPropertyValue(this.kafkaBatchListener, "errorChannel")).isSameAs(this.errorChannel); + KafkaMessageListenerContainer container = + TestUtils.getPropertyValue(this.kafkaBatchListener, "messageListenerContainer", + KafkaMessageListenerContainer.class); + assertThat(container).isNotNull(); + assertThat(TestUtils.getPropertyValue(kafkaBatchListener, "mode", ListenerMode.class)) + .isEqualTo(ListenerMode.batch); } @Test @@ -91,7 +112,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests { Object delegate = TestUtils.getPropertyValue(messageListener, "delegate"); - assertThat(delegate.getClass().getName()).contains("$IntegrationMessageListener"); + assertThat(delegate.getClass().getName()).contains("$IntegrationRecordMessageListener"); adapter.setRecordFilterStrategy(null); adapter.setRetryTemplate(new RetryTemplate()); @@ -102,7 +123,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests { delegate = TestUtils.getPropertyValue(messageListener, "delegate"); - assertThat(delegate.getClass().getName()).contains("$IntegrationMessageListener"); + assertThat(delegate.getClass().getName()).contains("$IntegrationRecordMessageListener"); adapter.setRecordFilterStrategy(mock(RecordFilterStrategy.class)); adapter.afterPropertiesSet(); 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 a04f6cae91..eb52e8dd54 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 @@ -19,6 +19,7 @@ package org.springframework.integration.kafka.inbound; import static org.assertj.core.api.Assertions.assertThat; import java.lang.reflect.Type; +import java.util.List; import java.util.Map; import org.apache.kafka.clients.consumer.ConsumerConfig; @@ -27,6 +28,7 @@ import org.junit.ClassRule; import org.junit.Test; import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.ListenerMode; import org.springframework.integration.support.MessageBuilder; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -37,8 +39,10 @@ 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.BatchMessagingMessageConverter; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.kafka.test.utils.ContainerTestUtils; import org.springframework.kafka.test.utils.KafkaTestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; @@ -53,11 +57,13 @@ public class MessageDrivenAdapterTests { private static String topic1 = "testTopic1"; + private static String topic2 = "testTopic2"; + @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, topic1); + public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, topic1, topic2); @Test - public void testInbound() throws Exception { + public void testInboundRecord() throws Exception { Map props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props); @@ -78,11 +84,12 @@ public class MessageDrivenAdapterTests { }); adapter.start(); + ContainerTestUtils.waitForAssignment(container, 2); Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); ProducerFactory pf = new DefaultKafkaProducerFactory(senderProps); KafkaTemplate template = new KafkaTemplate<>(pf); - template.setDefaultTopic("testTopic1"); + template.setDefaultTopic(topic1); template.sendDefault(1, "foo"); Message received = out.receive(10000); @@ -90,7 +97,7 @@ public class MessageDrivenAdapterTests { MessageHeaders 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_TOPIC)).isEqualTo(topic1); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0); assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L); assertThat(headers.get("testHeader")).isEqualTo("testValue"); @@ -103,7 +110,7 @@ public class MessageDrivenAdapterTests { 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_TOPIC)).isEqualTo(topic1); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0); assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(1L); assertThat(headers.get("testHeader")).isEqualTo("testValue"); @@ -111,4 +118,53 @@ public class MessageDrivenAdapterTests { adapter.stop(); } + @Test + public void testInboundBatch() throws Exception { + Map props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props); + ContainerProperties containerProps = new ContainerProperties(topic2); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + KafkaMessageDrivenChannelAdapter adapter = new KafkaMessageDrivenChannelAdapter<>(container, + ListenerMode.batch); + QueueChannel out = new QueueChannel(); + adapter.setOutputChannel(out); + adapter.afterPropertiesSet(); + adapter.setBatchMessageConverter(new BatchMessagingMessageConverter() { + + @Override + public Message toMessage(List> records, Acknowledgment acknowledgment, Type type) { + Message message = super.toMessage(records, acknowledgment, type); + return MessageBuilder.fromMessage(message).setHeader("testHeader", "testValue").build(); + } + + }); + adapter.start(); + ContainerTestUtils.waitForAssignment(container, 2); + + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + ProducerFactory pf = new DefaultKafkaProducerFactory(senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf); + template.setDefaultTopic(topic2); + template.sendDefault(1, "foo"); + template.sendDefault(1, "bar"); + + Message received = out.receive(10000); + assertThat(received).isNotNull(); + Object payload = received.getPayload(); + assertThat(payload).isInstanceOf(List.class); + List list = (List) payload; + assertThat(list.size()).isGreaterThan(0); + + MessageHeaders headers = received.getHeaders(); + assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY).toString()).contains("[1"); + assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC).toString()).contains(topic2); + assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID).toString()).contains("0"); + assertThat(headers.get(KafkaHeaders.OFFSET).toString()).contains("[0"); + assertThat(headers.get("testHeader")).isEqualTo("testValue"); + + adapter.stop(); + } + }