GH-147: Add Payload Type to Message Driven Adapter

Resolves: #147

Send Conversion Errors to the Error Channel
This commit is contained in:
Gary Russell
2016-09-30 14:23:44 -04:00
committed by Artem Bilan
parent 2c1b77c8fa
commit 171dc9832a
7 changed files with 235 additions and 7 deletions

View File

@@ -53,6 +53,7 @@ public class KafkaMessageDrivenChannelAdapterParser extends AbstractChannelAdapt
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-timeout");
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "error-channel");
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "message-converter");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "payload-type");
return builder.getBeanDefinition();
}

View File

@@ -33,9 +33,11 @@ import org.springframework.kafka.listener.adapter.RecordMessagingMessageListener
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.ConversionException;
import org.springframework.kafka.support.converter.MessageConverter;
import org.springframework.kafka.support.converter.RecordMessageConverter;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.retry.RecoveryCallback;
import org.springframework.retry.support.RetryTemplate;
import org.springframework.util.Assert;
@@ -193,6 +195,17 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
this.filterInRetry = filterInRetry;
}
/**
* When using a type-aware message converter (such as {@code StringJsonMessageConverter},
* set the payload type the converter should create. Defaults to {@link Object}.
* @param payloadType the type.
* @since 2.1.1
*/
public void setPayloadType(Class<?> payloadType) {
this.recordListener.setFallbackType(payloadType);
this.batchListener.setFallbackType(payloadType);
}
@Override
protected void onInit() {
super.onInit();
@@ -200,7 +213,8 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
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,
@@ -284,8 +298,23 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
@Override
public void onMessage(ConsumerRecord<K, V> record, Acknowledgment acknowledgment) {
Message<?> message = toMessagingMessage(record, acknowledgment);
sendMessage(message);
Message<?> message = null;
try {
message = toMessagingMessage(record, acknowledgment);
}
catch (RuntimeException e) {
Exception exception = new ConversionException("Failed to convert to message for: " + record, e);
if (getErrorChannel() != null) {
getMessagingTemplate().send(getErrorChannel(), new ErrorMessage(exception));
}
}
if (message != null) {
sendMessage(message);
}
else {
KafkaMessageDrivenChannelAdapter.this.logger.debug("Converter returned a null message for: "
+ record);
}
}
}
@@ -298,8 +327,23 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
@Override
public void onMessage(List<ConsumerRecord<K, V>> records, Acknowledgment acknowledgment) {
Message<?> message = toMessagingMessage(records, acknowledgment);
sendMessage(message);
Message<?> message = null;
try {
message = toMessagingMessage(records, acknowledgment);
}
catch (RuntimeException e) {
Exception exception = new ConversionException("Failed to convert to message for: " + records, e);
if (getErrorChannel() != null) {
getMessagingTemplate().send(getErrorChannel(), new ErrorMessage(exception));
}
}
if (message != null) {
sendMessage(message);
}
else {
KafkaMessageDrivenChannelAdapter.this.logger.debug("Converter returned a null message for: "
+ records);
}
}
}

View File

@@ -188,6 +188,14 @@
<xsd:union memberTypes="listenerMode xsd:string"/>
</xsd:simpleType>
</xsd:attribute>
<xsd:attribute name="payload-type" type="xsd:string">
<xsd:annotation>
<xsd:documentation>
Set the payload type to convert to when using a type-aware message converter such as the
StringJsonMessageConverter. Fully qualified class name; defaults to 'java.lang.Object'.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
</xsd:complexType>
</xsd:element>

View File

@@ -21,6 +21,7 @@
send-timeout="5000"
channel="nullChannel"
message-converter="messageConverter"
payload-type="java.lang.String"
error-channel="errorChannel" />
<int-kafka:message-driven-channel-adapter

View File

@@ -74,6 +74,10 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
assertThat(container).isNotNull();
assertThat(TestUtils.getPropertyValue(kafkaListener, "mode", ListenerMode.class))
.isEqualTo(ListenerMode.record);
assertThat(TestUtils.getPropertyValue(this.kafkaListener, "recordListener.fallbackType"))
.isEqualTo(String.class);
assertThat(TestUtils.getPropertyValue(this.kafkaListener, "batchListener.fallbackType"))
.isEqualTo(String.class);
}
@Test

View File

@@ -24,6 +24,7 @@ import java.util.Map;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.junit.ClassRule;
import org.junit.Test;
@@ -39,13 +40,18 @@ 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.BatchMessageConverter;
import org.springframework.kafka.support.converter.BatchMessagingMessageConverter;
import org.springframework.kafka.support.converter.ConversionException;
import org.springframework.kafka.support.converter.MessagingMessageConverter;
import org.springframework.kafka.support.converter.RecordMessageConverter;
import org.springframework.kafka.support.converter.StringJsonMessageConverter;
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;
import org.springframework.messaging.PollableChannel;
/**
* @author Gary Russell
@@ -59,8 +65,10 @@ public class MessageDrivenAdapterTests {
private static String topic2 = "testTopic2";
private static String topic3 = "testTopic3";
@ClassRule
public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, topic1, topic2);
public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, topic1, topic2, topic3);
@Test
public void testInboundRecord() throws Exception {
@@ -115,12 +123,33 @@ public class MessageDrivenAdapterTests {
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(1L);
assertThat(headers.get("testHeader")).isEqualTo("testValue");
adapter.setMessageConverter(new RecordMessageConverter() {
@Override
public Message<?> toMessage(ConsumerRecord<?, ?> record, Acknowledgment acknowledgment, Type payloadType) {
throw new RuntimeException("testError");
}
@Override
public ProducerRecord<?, ?> fromMessage(Message<?> message, String defaultTopic) {
return null;
}
});
PollableChannel errors = new QueueChannel();
adapter.setErrorChannel(errors);
template.sendDefault(1, "bar");
Message<?> error = errors.receive(10000);
assertThat(error).isNotNull();
assertThat(error.getPayload()).isInstanceOf(ConversionException.class);
assertThat(((ConversionException) error.getPayload()).getMessage())
.contains("Failed to convert to message for: ConsumerRecord(topic = testTopic1");
adapter.stop();
}
@Test
public void testInboundBatch() throws Exception {
Map<String, Object> props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka);
Map<String, Object> props = KafkaTestUtils.consumerProps("test2", "true", embeddedKafka);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
DefaultKafkaConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<Integer, String>(props);
ContainerProperties containerProps = new ContainerProperties(topic2);
@@ -164,7 +193,131 @@ public class MessageDrivenAdapterTests {
assertThat(headers.get(KafkaHeaders.OFFSET).toString()).contains("[0");
assertThat(headers.get("testHeader")).isEqualTo("testValue");
adapter.setMessageConverter(new BatchMessageConverter() {
@Override
public Message<?> toMessage(List<ConsumerRecord<?, ?>> records, Acknowledgment acknowledgment,
Type payloadType) {
throw new RuntimeException("testError");
}
@Override
public List<ProducerRecord<?, ?>> fromMessage(Message<?> message, String defaultTopic) {
return null;
}
});
PollableChannel errors = new QueueChannel();
adapter.setErrorChannel(errors);
template.sendDefault(1, "bar");
Message<?> error = errors.receive(10000);
assertThat(error).isNotNull();
assertThat(error.getPayload()).isInstanceOf(ConversionException.class);
assertThat(((ConversionException) error.getPayload()).getMessage())
.contains("Failed to convert to message for: [ConsumerRecord(topic = testTopic2");
adapter.stop();
}
@Test
public void testInboundJson() throws Exception {
Map<String, Object> props = KafkaTestUtils.consumerProps("test3", "true", embeddedKafka);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
DefaultKafkaConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<Integer, String>(props);
ContainerProperties containerProps = new ContainerProperties(topic3);
KafkaMessageListenerContainer<Integer, String> container =
new KafkaMessageListenerContainer<>(cf, containerProps);
KafkaMessageDrivenChannelAdapter<Integer, String> adapter = new KafkaMessageDrivenChannelAdapter<>(container);
adapter.setRecordMessageConverter(new StringJsonMessageConverter());
QueueChannel out = new QueueChannel();
adapter.setOutputChannel(out);
adapter.afterPropertiesSet();
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(topic3);
template.sendDefault(1, "{\"bar\":\"baz\"}");
Message<?> received = out.receive(10000);
assertThat(received).isNotNull();
MessageHeaders headers = received.getHeaders();
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic3);
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L);
assertThat(received.getPayload()).isInstanceOf(Map.class);
adapter.setPayloadType(Foo.class);
template.sendDefault(1, "{\"bar\":\"baz\"}");
received = out.receive(10000);
assertThat(received).isNotNull();
headers = received.getHeaders();
assertThat(headers.get(KafkaHeaders.RECEIVED_MESSAGE_KEY)).isEqualTo(1);
assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic3);
assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0);
assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(1L);
assertThat(received.getPayload()).isInstanceOf(Foo.class);
assertThat(received.getPayload()).isEqualTo(new Foo("baz"));
adapter.stop();
}
public static class Foo {
private String bar;
public Foo() {
}
public Foo(String bar) {
this.bar = bar;
}
protected String getBar() {
return this.bar;
}
protected void setBar(String bar) {
this.bar = bar;
}
@Override
public int hashCode() {
final int prime = 31;
int result = 1;
result = prime * result + ((this.bar == null) ? 0 : this.bar.hashCode());
return result;
}
@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null) {
return false;
}
if (getClass() != obj.getClass()) {
return false;
}
Foo other = (Foo) obj;
if (this.bar == null) {
if (other.bar != null) {
return false;
}
}
else if (!this.bar.equals(other.bar)) {
return false;
}
return true;
}
}
}

View File

@@ -0,0 +1,17 @@
<?xml version="1.0" encoding="UTF-8"?>
<Configuration>
<Appenders>
<Console name="STDOUT" target="SYSTEM_OUT">
<PatternLayout pattern="%d %5p %c [%t] : %m%n" />
</Console>
</Appenders>
<Loggers>
<Logger name="kafka" level="warn"/>
<Logger name="org.apache.kafka" level="warn"/>
<Logger name="org.springframework.kafka" level="debug"/>
<Logger name="org.springframework.integration" level="debug"/>
<Root level="warn">
<AppenderRef ref="STDOUT" />
</Root>
</Loggers>
</Configuration>