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 3bf247cc18..0db6167370 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 @@ -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(); } 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 3859780505..7d82dda7a3 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 @@ -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 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 extends MessageProducerSuppo 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, @@ -284,8 +298,23 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo @Override public void onMessage(ConsumerRecord 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 extends MessageProducerSuppo @Override public void onMessage(List> 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); + } } } diff --git a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd index a9400888f6..ff2f1ebfd0 100644 --- a/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd +++ b/spring-integration-kafka/src/main/resources/org/springframework/integration/kafka/config/spring-integration-kafka-2.1.xsd @@ -188,6 +188,14 @@ + + + + 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'. + + + 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 1da36150c5..101789a50d 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 @@ -21,6 +21,7 @@ send-timeout="5000" channel="nullChannel" message-converter="messageConverter" + payload-type="java.lang.String" error-channel="errorChannel" /> 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 props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka); + Map props = KafkaTestUtils.consumerProps("test2", "true", embeddedKafka); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(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> records, Acknowledgment acknowledgment, + Type payloadType) { + throw new RuntimeException("testError"); + } + + @Override + public List> 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 props = KafkaTestUtils.consumerProps("test3", "true", embeddedKafka); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props); + ContainerProperties containerProps = new ContainerProperties(topic3); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + KafkaMessageDrivenChannelAdapter adapter = new KafkaMessageDrivenChannelAdapter<>(container); + adapter.setRecordMessageConverter(new StringJsonMessageConverter()); + QueueChannel out = new QueueChannel(); + adapter.setOutputChannel(out); + adapter.afterPropertiesSet(); + adapter.start(); + ContainerTestUtils.waitForAssignment(container, 2); + + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + ProducerFactory pf = new DefaultKafkaProducerFactory(senderProps); + KafkaTemplate 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; + } + + } + } diff --git a/spring-integration-kafka/src/test/resources/log4j2-test.xml b/spring-integration-kafka/src/test/resources/log4j2-test.xml new file mode 100644 index 0000000000..a2b7417bca --- /dev/null +++ b/spring-integration-kafka/src/test/resources/log4j2-test.xml @@ -0,0 +1,17 @@ + + + + + + + + + + + + + + + + +