diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index 3160693a4..79218652f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -47,7 +47,6 @@ import org.springframework.cloud.stream.binding.StreamListenerResultAdapter; import org.springframework.cloud.stream.config.BinderProperties; import org.springframework.cloud.stream.config.BindingServiceConfiguration; import org.springframework.cloud.stream.config.BindingServiceProperties; -import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Conditional; @@ -56,6 +55,7 @@ import org.springframework.core.env.Environment; import org.springframework.core.env.MapPropertySource; import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.core.CleanupConfig; +import org.springframework.messaging.converter.CompositeMessageConverter; import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; @@ -266,17 +266,17 @@ public class KafkaStreamsBinderSupportAutoConfiguration { @Bean public KafkaStreamsMessageConversionDelegate messageConversionDelegate( - CompositeMessageConverterFactory compositeMessageConverterFactory, + CompositeMessageConverter compositeMessageConverter, SendToDlqAndContinue sendToDlqAndContinue, KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, KafkaStreamsBinderConfigurationProperties binderConfigurationProperties) { - return new KafkaStreamsMessageConversionDelegate(compositeMessageConverterFactory, sendToDlqAndContinue, + return new KafkaStreamsMessageConversionDelegate(compositeMessageConverter, sendToDlqAndContinue, KafkaStreamsBindingInformationCatalogue, binderConfigurationProperties); } @Bean public CompositeNonNativeSerde compositeNonNativeSerde( - CompositeMessageConverterFactory compositeMessageConverterFactory) { + CompositeMessageConverter compositeMessageConverterFactory) { return new CompositeNonNativeSerde(compositeMessageConverterFactory); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java index 052ca3870..4a4b3b1b2 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java @@ -33,9 +33,9 @@ import org.apache.kafka.streams.processor.Processor; import org.apache.kafka.streams.processor.ProcessorContext; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; -import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.CompositeMessageConverter; import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; @@ -57,7 +57,7 @@ public class KafkaStreamsMessageConversionDelegate { private static final ThreadLocal> keyValueThreadLocal = new ThreadLocal<>(); - private final CompositeMessageConverterFactory compositeMessageConverterFactory; + private final CompositeMessageConverter compositeMessageConverter; private final SendToDlqAndContinue sendToDlqAndContinue; @@ -66,11 +66,11 @@ public class KafkaStreamsMessageConversionDelegate { private final KafkaStreamsBinderConfigurationProperties kstreamBinderConfigurationProperties; KafkaStreamsMessageConversionDelegate( - CompositeMessageConverterFactory compositeMessageConverterFactory, + CompositeMessageConverter compositeMessageConverter, SendToDlqAndContinue sendToDlqAndContinue, KafkaStreamsBindingInformationCatalogue kstreamBindingInformationCatalogue, KafkaStreamsBinderConfigurationProperties kstreamBinderConfigurationProperties) { - this.compositeMessageConverterFactory = compositeMessageConverterFactory; + this.compositeMessageConverter = compositeMessageConverter; this.sendToDlqAndContinue = sendToDlqAndContinue; this.kstreamBindingInformationCatalogue = kstreamBindingInformationCatalogue; this.kstreamBinderConfigurationProperties = kstreamBinderConfigurationProperties; @@ -85,8 +85,7 @@ public class KafkaStreamsMessageConversionDelegate { public KStream serializeOnOutbound(KStream outboundBindTarget) { String contentType = this.kstreamBindingInformationCatalogue .getContentType(outboundBindTarget); - MessageConverter messageConverter = this.compositeMessageConverterFactory - .getMessageConverterForAllRegistered(); + MessageConverter messageConverter = this.compositeMessageConverter; final PerRecordContentTypeHolder perRecordContentTypeHolder = new PerRecordContentTypeHolder(); final KStream kStreamWithEnrichedHeaders = outboundBindTarget.mapValues((v) -> { @@ -148,8 +147,7 @@ public class KafkaStreamsMessageConversionDelegate { @SuppressWarnings({ "unchecked", "rawtypes" }) public KStream deserializeOnInbound(Class valueClass, KStream bindingTarget) { - MessageConverter messageConverter = this.compositeMessageConverterFactory - .getMessageConverterForAllRegistered(); + MessageConverter messageConverter = this.compositeMessageConverter; final PerRecordContentTypeHolder perRecordContentTypeHolder = new PerRecordContentTypeHolder(); resolvePerRecordContentType(bindingTarget, perRecordContentTypeHolder); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java index 3d59cdce9..d6a0ce117 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerde.java @@ -24,9 +24,9 @@ import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serializer; -import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.CompositeMessageConverter; import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.Assert; @@ -35,7 +35,7 @@ import org.springframework.util.MimeTypeUtils; /** * A {@link Serde} implementation that wraps the list of {@link MessageConverter}s from - * {@link CompositeMessageConverterFactory}. + * {@link CompositeMessageConverter}. * * The primary motivation for this class is to provide an avro based {@link Serde} that is * compatible with the schema registry that Spring Cloud Stream provides. When using the @@ -90,11 +90,11 @@ public class CompositeNonNativeSerde implements Serde { private final CompositeNonNativeSerializer compositeNonNativeSerializer; public CompositeNonNativeSerde( - CompositeMessageConverterFactory compositeMessageConverterFactory) { + CompositeMessageConverter compositeMessageConverter) { this.compositeNonNativeDeserializer = new CompositeNonNativeDeserializer<>( - compositeMessageConverterFactory); + compositeMessageConverter); this.compositeNonNativeSerializer = new CompositeNonNativeSerializer<>( - compositeMessageConverterFactory); + compositeMessageConverter); } @Override @@ -150,9 +150,8 @@ public class CompositeNonNativeSerde implements Serde { private Class valueClass; CompositeNonNativeDeserializer( - CompositeMessageConverterFactory compositeMessageConverterFactory) { - this.messageConverter = compositeMessageConverterFactory - .getMessageConverterForAllRegistered(); + CompositeMessageConverter compositeMessageConverter) { + this.messageConverter = compositeMessageConverter; } @Override @@ -197,9 +196,8 @@ public class CompositeNonNativeSerde implements Serde { private MimeType mimeType; CompositeNonNativeSerializer( - CompositeMessageConverterFactory compositeMessageConverterFactory) { - this.messageConverter = compositeMessageConverterFactory - .getMessageConverterForAllRegistered(); + CompositeMessageConverter compositeMessageConverter) { + this.messageConverter = compositeMessageConverter; } @Override diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java index 57dcde65c..9ba84b187 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java @@ -55,7 +55,7 @@ public class CompositeNonNativeSerdeTest { CompositeMessageConverterFactory compositeMessageConverterFactory = new CompositeMessageConverterFactory( messageConverters, new ObjectMapper()); CompositeNonNativeSerde compositeNonNativeSerde = new CompositeNonNativeSerde( - compositeMessageConverterFactory); + compositeMessageConverterFactory.getMessageConverterForAllRegistered()); Map configs = new HashMap<>(); configs.put("valueClass", Sensor.class);