GH-2293: MessageConverter Bean Name

Add test to consume full `ReceiverRecord`.

Resolves #2295
This commit is contained in:
Gary Russell
2022-03-12 17:17:03 -05:00
committed by Oleg Zhurakousky
parent 9ccf9ce48b
commit 6d8be51437
4 changed files with 105 additions and 13 deletions

View File

@@ -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<Object, Object> 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<KafkaConsumerProperties> 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;

View File

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

View File

@@ -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<String, String> consumer;
private static Consumer<String, String> consumer1;
private static Consumer<String, String> consumer2;
@BeforeAll
public static void setUp() {
@@ -58,8 +68,10 @@ public class ReactorKafkaBinderIntegrationTests {
embeddedKafka);
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
DefaultKafkaConsumerFactory<String, String> 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<String, Object> senderProps = KafkaTestUtils.producerProps(embeddedKafka);
DefaultKafkaProducerFactory<Integer, String> pf = new DefaultKafkaProducerFactory<>(senderProps);
try {
KafkaTemplate<Integer, String> template = new KafkaTemplate<>(pf, true);
template.setDefaultTopic("words");
template.sendDefault("foobar");
template.send("words", "foobar");
template.send("words1", "BAZQUX");
ConsumerRecord<String, String> cr = KafkaTestUtils.getSingleRecord(consumer, "uppercased-words");
assertThat(cr.value().equals("FOOBAR")).isTrue();
ConsumerRecord<String, String> 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<String>, Flux<String>> uppercase() {
return s -> s.map(String::toUpperCase);
}
@Bean
public Function<Flux<ReceiverRecord<byte[], byte[]>>, Flux<String>> lowercase() {
return s -> s.map(rec -> new String(rec.value()).toLowerCase());
}
}
}

View File

@@ -7,9 +7,9 @@
<logger name="org.apache.kafka" level="WARN"/>
<logger name="reactor.kafka" level="DEBUG"/>
<logger name="org.springframework.integration.kafka" level="INFO"/>
<logger name="org.springframework.kafka" level="INFO"/>
<logger name="org.springframework.kafka" level="DEBUG"/>
<logger name="org.springframework.cloud.stream" level="INFO" />
<logger name="org.springframework.integration.channel" level="INFO" />
<logger name="org.springframework.integration.channel" level="DEBUG" />
<logger name="kafka.server.ReplicaFetcherThread" level="ERROR"/>
<logger name="kafka.server.LogDirFailureChannel" level="FATAL"/>
<logger name="kafka.server.BrokerMetadataCheckpoint" level="ERROR"/>