From 6d8be514373ed9e5f32fcec8bbef7754cbf1fb0d Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Sat, 12 Mar 2022 17:17:03 -0500 Subject: [PATCH] GH-2293: MessageConverter Bean Name Add test to consume full `ReceiverRecord`. Resolves #2295 --- .../reactorkafka/ReactorKafkaBinder.java | 46 ++++++++++++- .../ReactorKafkaBinderConfiguration.java | 4 +- .../ReactorKafkaBinderIntegrationTests.java | 64 ++++++++++++++++--- .../src/test/resources/logback.xml | 4 +- 4 files changed, 105 insertions(+), 13 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java index 07be04529..798777dec 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java @@ -41,6 +41,7 @@ import reactor.kafka.sender.SenderOptions; import reactor.kafka.sender.SenderRecord; import reactor.kafka.sender.SenderResult; +import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; @@ -48,6 +49,7 @@ import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties.StandardHeaders; import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; @@ -57,6 +59,7 @@ import org.springframework.context.Lifecycle; import org.springframework.integration.core.MessageProducer; import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.handler.AbstractMessageHandler; +import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.converter.MessagingMessageConverter; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.messaging.Message; @@ -81,7 +84,7 @@ public class ReactorKafkaBinder private final KafkaBinderConfigurationProperties configurationProperties; - private final KafkaExtendedBindingProperties extendedBindingProperties = new KafkaExtendedBindingProperties(); + private KafkaExtendedBindingProperties extendedBindingProperties = new KafkaExtendedBindingProperties(); public ReactorKafkaBinder(KafkaBinderConfigurationProperties configurationProperties, KafkaTopicProvisioner provisioner) { @@ -123,7 +126,7 @@ public class ReactorKafkaBinder // this.consumerConfigCustomizer.configure(configs, bindingNameHolder.get(), destination); // } - RecordMessageConverter converter = new MessagingMessageConverter(); + RecordMessageConverter converter = getMessageConverter(properties); ReceiverOptions opts = ReceiverOptions.create(configs) .addAssignListener(parts -> System.out.println("Assigned: " + parts)) .subscription(Collections.singletonList(destination.getName())); @@ -155,6 +158,39 @@ public class ReactorKafkaBinder }; } + /* + * TODO: Copied (and modified) from Kafka binder - refactor to core + */ + private RecordMessageConverter getMessageConverter( + final ExtendedConsumerProperties extendedConsumerProperties) { + + RecordMessageConverter messageConverter; + if (extendedConsumerProperties.getExtension().getConverterBeanName() == null) { + MessagingMessageConverter mmc = new MessagingMessageConverter(); + StandardHeaders standardHeaders = extendedConsumerProperties.getExtension() + .getStandardHeaders(); + mmc.setGenerateMessageId(StandardHeaders.id.equals(standardHeaders) + || StandardHeaders.both.equals(standardHeaders)); + mmc.setGenerateTimestamp( + StandardHeaders.timestamp.equals(standardHeaders) + || StandardHeaders.both.equals(standardHeaders)); + mmc.setHeaderMapper(new DefaultKafkaHeaderMapper()); //TODO + messageConverter = mmc; + } + else { + try { + messageConverter = getApplicationContext().getBean( + extendedConsumerProperties.getExtension().getConverterBeanName(), + RecordMessageConverter.class); + } + catch (NoSuchBeanDefinitionException ex) { + throw new IllegalStateException( + "Converter bean not present in application context", ex); + } + } + return messageConverter; + } + /* * TODO: Copied from Kafka binder - refactor to core */ @@ -262,6 +298,12 @@ public class ReactorKafkaBinder return this.extendedBindingProperties.getExtendedPropertiesEntryClass(); } + public void setExtendedBindingProperties( + KafkaExtendedBindingProperties extendedBindingProperties) { + + this.extendedBindingProperties = extendedBindingProperties; + } + private static class ReactorMessageHandler extends AbstractMessageHandler implements Lifecycle { private final RecordMessageConverter converter; diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderConfiguration.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderConfiguration.java index f398b9e9e..bee014a3a 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderConfiguration.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderConfiguration.java @@ -52,9 +52,11 @@ public class ReactorKafkaBinderConfiguration { @Bean ReactorKafkaBinder reactorKafkaBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider) { + KafkaTopicProvisioner provisioningProvider, + KafkaExtendedBindingProperties extendedBindingProperties) { ReactorKafkaBinder reactorKafkaBinder = new ReactorKafkaBinder(configurationProperties, provisioningProvider); + reactorKafkaBinder.setExtendedBindingProperties(extendedBindingProperties); return reactorKafkaBinder; } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java index b548a25c4..76eb38084 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinderIntegrationTests.java @@ -16,28 +16,36 @@ package org.springframework.cloud.stream.binder.reactorkafka; +import java.lang.reflect.Type; import java.util.Map; import java.util.function.Function; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; +import reactor.kafka.receiver.ReceiverRecord; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import org.springframework.integration.support.MessageBuilder; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.support.Acknowledgment; +import org.springframework.kafka.support.converter.MessagingMessageConverter; +import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.condition.EmbeddedKafkaCondition; import org.springframework.kafka.test.context.EmbeddedKafka; import org.springframework.kafka.test.utils.KafkaTestUtils; +import org.springframework.messaging.Message; import static org.assertj.core.api.Assertions.assertThat; @@ -45,12 +53,14 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Soby Chacko * @author Gary Russell */ -@EmbeddedKafka(topics = "uppercased-words") +@EmbeddedKafka(topics = { "uppercased-words", "lowercased-words" }) public class ReactorKafkaBinderIntegrationTests { private static final EmbeddedKafkaBroker embeddedKafka = EmbeddedKafkaCondition.getBroker(); - private static Consumer consumer; + private static Consumer consumer1; + + private static Consumer consumer2; @BeforeAll public static void setUp() { @@ -58,8 +68,10 @@ public class ReactorKafkaBinderIntegrationTests { embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); - consumer = cf.createConsumer(); - embeddedKafka.consumeFromEmbeddedTopics(consumer, "uppercased-words"); + consumer1 = cf.createConsumer(); + embeddedKafka.consumeFromEmbeddedTopics(consumer1, "uppercased-words"); + consumer2 = cf.createConsumer("group2", null); + embeddedKafka.consumeFromEmbeddedTopics(consumer2, "lowercased-words"); } @Test @@ -70,20 +82,29 @@ public class ReactorKafkaBinderIntegrationTests { try (ConfigurableApplicationContext context = app.run( "--server.port=0", "--spring.jmx.enabled=false", + "--spring.cloud.function.definition=uppercase;lowercase", "--spring.cloud.stream.function.reactive.uppercase=true", + "--spring.cloud.stream.function.reactive.lowercase=true", + "--spring.cloud.stream.bindings.uppercase-in-0.group=grp1", "--spring.cloud.stream.bindings.uppercase-in-0.destination=words", "--spring.cloud.stream.bindings.uppercase-out-0.destination=uppercased-words", + "--spring.cloud.stream.bindings.lowercase-in-0.group=grp2", + "--spring.cloud.stream.bindings.lowercase-in-0.destination=words1", + "--spring.cloud.stream.bindings.lowercase-out-0.destination=lowercased-words", + "--spring.cloud.stream.kafka.bindings.lowercase-in-0.consumer.converterBeanName=fullRR", "--spring.cloud.stream.kafka.binder.brokers=" + embeddedKafka.getBrokersAsString())) { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); try { KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); - template.sendDefault("foobar"); + template.send("words", "foobar"); + template.send("words1", "BAZQUX"); - ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, "uppercased-words"); - assertThat(cr.value().equals("FOOBAR")).isTrue(); + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer1, "uppercased-words"); + assertThat(cr.value()).isEqualTo("FOOBAR"); + cr = KafkaTestUtils.getSingleRecord(consumer2, "lowercased-words"); + assertThat(cr.value()).isEqualTo("bazqux"); } finally { pf.destroy(); @@ -94,10 +115,37 @@ public class ReactorKafkaBinderIntegrationTests { @EnableAutoConfiguration public static class ReactiveKafkaApplication { + @Bean + RecordMessageConverter fullRR() { + return new RecordMessageConverter() { + + private final RecordMessageConverter converter = new MessagingMessageConverter(); + + @Override + public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, + Consumer consumer, Type payloadType) { + + return MessageBuilder.withPayload(record).build(); + } + + @Override + public ProducerRecord fromMessage(Message message, String defaultTopic) { + return this.converter.fromMessage(message, defaultTopic); + } + + }; + } + @Bean public Function, Flux> uppercase() { return s -> s.map(String::toUpperCase); } + @Bean + public Function>, Flux> lowercase() { + return s -> s.map(rec -> new String(rec.value()).toLowerCase()); + } + } + } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/resources/logback.xml b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/resources/logback.xml index 32c681652..5fdac6f28 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/resources/logback.xml +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/test/resources/logback.xml @@ -7,9 +7,9 @@ - + - +