GH-141: Support Batch Payloads

Resolves: https://github.com/spring-projects/spring-integration-kafka/issues/141

Enum Javadocs

Polishing - PR Comments
This commit is contained in:
Gary Russell
2016-09-07 16:49:35 -04:00
committed by Artem Bilan
parent 0cf61af620
commit cbfc6b0e73
7 changed files with 249 additions and 35 deletions

View File

@@ -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");

View File

@@ -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<K, V> extends MessageProducerSuppo
private final AbstractMessageListenerContainer<K, V> messageListenerContainer;
private final RecordMessagingMessageListenerAdapter<K, V> listener = new IntegrationMessageListener();
private final RecordMessagingMessageListenerAdapter<K, V> recordListener = new IntegrationRecordMessageListener();
private final BatchMessagingMessageListenerAdapter<K, V> batchListener = new IntegrationBatchMessageListener();
private final ListenerMode mode;
private RecordFilterStrategy<K, V> recordFilterStrategy;
@@ -61,23 +71,47 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
private boolean filterInRetry;
/**
* Construct an instance with mode {@link ListenerMode#record}.
* @param messageListenerContainer the container.
*/
public KafkaMessageDrivenChannelAdapter(AbstractMessageListenerContainer<K, V> 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<K, V> 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<K, V> 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<K, V> 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<K, V> extends MessageProducerSuppo
protected void onInit() {
super.onInit();
AcknowledgingMessageListener<K, V> listener = this.listener;
if (this.mode.equals(ListenerMode.record)) {
AcknowledgingMessageListener<K, V> 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<K, V> 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<K, V> extends MessageProducerSuppo
return getPhase();
}
private class IntegrationMessageListener extends RecordMessagingMessageListenerAdapter<K, V> {
/**
* 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<K, V> {
IntegrationRecordMessageListener() {
super(null, null);
}
@@ -216,4 +290,18 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
}
private class IntegrationBatchMessageListener extends BatchMessagingMessageListenerAdapter<K, V> {
IntegrationBatchMessageListener() {
super(null, null);
}
@Override
public void onMessage(List<ConsumerRecord<K, V>> records, Acknowledgment acknowledgment) {
Message<?> message = toMessagingMessage(records, acknowledgment);
sendMessage(message);
}
}
}

View File

@@ -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

View File

@@ -155,6 +155,8 @@
<xsd:annotation>
<xsd:documentation>
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.
</xsd:documentation>
<xsd:appinfo>
<tool:annotation kind="ref">
@@ -163,7 +165,25 @@
</xsd:appinfo>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="mode" default="record">
<xsd:annotation>
<xsd:documentation>
'record' or 'batch' - default 'record' - one converted ConsumerRecord per message, when
'batch' then the payload is a collection of converted ConsumerRecords.
</xsd:documentation>
</xsd:annotation>
<xsd:simpleType>
<xsd:union memberTypes="listenerMode xsd:string"/>
</xsd:simpleType>
</xsd:attribute>
</xsd:complexType>
</xsd:element>
<xsd:simpleType name="listenerMode">
<xsd:restriction base="xsd:token">
<xsd:enumeration value="record" />
<xsd:enumeration value="batch" />
</xsd:restriction>
</xsd:simpleType>
</xsd:schema>

View File

@@ -23,6 +23,17 @@
message-converter="messageConverter"
error-channel="errorChannel" />
<int-kafka:message-driven-channel-adapter
id="kafkaBatchListener"
listener-container="container2"
auto-startup="false"
phase="100"
send-timeout="5000"
channel="nullChannel"
mode="batch"
message-converter="messageConverter"
error-channel="errorChannel" />
<bean id="messageConverter" class="org.springframework.kafka.support.converter.MessagingMessageConverter"/>
<bean id="container1" class="org.springframework.kafka.listener.KafkaMessageListenerContainer">
@@ -42,4 +53,21 @@
</constructor-arg>
</bean>
<bean id="container2" class="org.springframework.kafka.listener.KafkaMessageListenerContainer">
<constructor-arg>
<bean class="org.springframework.kafka.core.DefaultKafkaConsumerFactory">
<constructor-arg>
<map>
<entry key="" value="" />
</map>
</constructor-arg>
</bean>
</constructor-arg>
<constructor-arg>
<bean class="org.springframework.kafka.listener.config.ContainerProperties">
<constructor-arg name="topics" value="foo" />
</bean>
</constructor-arg>
</bean>
</beans>

View File

@@ -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();

View File

@@ -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<String, Object> props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
DefaultKafkaConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<Integer, String>(props);
@@ -78,11 +84,12 @@ public class MessageDrivenAdapterTests {
});
adapter.start();
ContainerTestUtils.waitForAssignment(container, 2);
Map<String, Object> senderProps = KafkaTestUtils.producerProps(embeddedKafka);
ProducerFactory<Integer, String> pf = new DefaultKafkaProducerFactory<Integer, String>(senderProps);
KafkaTemplate<Integer, String> 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<String, Object> props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
DefaultKafkaConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<Integer, String>(props);
ContainerProperties containerProps = new ContainerProperties(topic2);
KafkaMessageListenerContainer<Integer, String> container =
new KafkaMessageListenerContainer<>(cf, containerProps);
KafkaMessageDrivenChannelAdapter<Integer, String> adapter = new KafkaMessageDrivenChannelAdapter<>(container,
ListenerMode.batch);
QueueChannel out = new QueueChannel();
adapter.setOutputChannel(out);
adapter.afterPropertiesSet();
adapter.setBatchMessageConverter(new BatchMessagingMessageConverter() {
@Override
public Message<?> toMessage(List<ConsumerRecord<?, ?>> 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<String, Object> senderProps = KafkaTestUtils.producerProps(embeddedKafka);
ProducerFactory<Integer, String> pf = new DefaultKafkaProducerFactory<Integer, String>(senderProps);
KafkaTemplate<Integer, String> 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();
}
}