From 600e88ec56c087ea25113771757e2b27399ca9b8 Mon Sep 17 00:00:00 2001 From: Cameron Mayfield Date: Tue, 25 Jun 2019 07:12:37 -0700 Subject: [PATCH] GH-272: Add KafkaMDrivenChAdapterSpec.payloadType Fixes https://github.com/spring-projects/spring-integration-kafka/issues/272 Add `payloadType` option into `KafkaMessageDrivenChannelAdapterSpec` --- .../KafkaMessageDrivenChannelAdapterSpec.java | 13 +++++ .../inbound/MessageDrivenAdapterTests.java | 47 ++++++++++++++++--- 2 files changed, 53 insertions(+), 7 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java index 4a03b9e990..949951be76 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java @@ -46,6 +46,7 @@ import org.springframework.util.Assert; * * @author Artem Bilan * @author Gary Russell + * @author Cameron Mayfield * * @since 3.0 */ @@ -141,6 +142,18 @@ public class KafkaMessageDrivenChannelAdapterSpec props = KafkaTestUtils.consumerProps("test3", "true", embeddedKafka); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); @@ -428,17 +432,46 @@ public class MessageDrivenAdapterTests { assertThat(headers.get("foo")).isEqualTo("bar"); assertThat(received.getPayload()).isInstanceOf(Map.class); - adapter.setPayloadType(Foo.class); + adapter.stop(); + } + + @Test + public void testInboundJsonWithPayload() { + Map props = KafkaTestUtils.consumerProps("test6", "true", embeddedKafka); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); + ContainerProperties containerProps = new ContainerProperties(topic6); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + + KafkaMessageDrivenChannelAdapter adapter = Kafka.messageDrivenChannelAdapter(container, ListenerMode.record) + .recordMessageConverter(new StringJsonMessageConverter()) + .payloadType(Foo.class) + .get(); + 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(topic6); + Headers kHeaders = new RecordHeaders(); + MessageHeaders siHeaders = new MessageHeaders(Collections.singletonMap("foo", "bar")); + new DefaultKafkaHeaderMapper().fromHeaders(siHeaders, kHeaders); + template.sendDefault(1, "{\"bar\":\"baz\"}"); - received = out.receive(10000); + Message received = out.receive(10000); assertThat(received).isNotNull(); - headers = received.getHeaders(); + 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_TOPIC)).isEqualTo(topic6); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0); - assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(1L); + 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");