diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java index 8c76de738..6b5e806fd 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java @@ -23,6 +23,7 @@ import org.apache.kafka.streams.kstream.GlobalKTable; import org.springframework.aop.framework.ProxyFactory; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binding.AbstractBindingTargetFactory; +import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.util.Assert; @@ -47,8 +48,12 @@ public class GlobalKTableBoundElementFactory @Override public GlobalKTable createInput(String name) { - ConsumerProperties consumerProperties = this.bindingServiceProperties - .getConsumerProperties(name); + BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(name); + ConsumerProperties consumerProperties = bindingProperties.getConsumer(); + if (consumerProperties == null) { + consumerProperties = this.bindingServiceProperties.getConsumerProperties(name); + consumerProperties.setUseNativeDecoding(true); + } // Always set multiplex to true in the kafka streams binder consumerProperties.setMultiplex(true); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java index 3a9e47351..9376c53f1 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java @@ -22,6 +22,7 @@ import org.apache.kafka.streams.kstream.KStream; import org.springframework.aop.framework.ProxyFactory; import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binding.AbstractBindingTargetFactory; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; @@ -52,8 +53,12 @@ class KStreamBoundElementFactory extends AbstractBindingTargetFactory { @Override public KStream createInput(String name) { - ConsumerProperties consumerProperties = this.bindingServiceProperties - .getConsumerProperties(name); + BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(name); + ConsumerProperties consumerProperties = bindingProperties.getConsumer(); + if (consumerProperties == null) { + consumerProperties = this.bindingServiceProperties.getConsumerProperties(name); + consumerProperties.setUseNativeDecoding(true); + } // Always set multiplex to true in the kafka streams binder consumerProperties.setMultiplex(true); return createProxyForKStream(name); @@ -62,6 +67,13 @@ class KStreamBoundElementFactory extends AbstractBindingTargetFactory { @Override @SuppressWarnings("unchecked") public KStream createOutput(final String name) { + + BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(name); + ProducerProperties producerProperties = bindingProperties.getProducer(); + if (producerProperties == null) { + producerProperties = this.bindingServiceProperties.getProducerProperties(name); + producerProperties.setUseNativeEncoding(true); + } return createProxyForKStream(name); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java index dbb958872..f9f3fef02 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java @@ -23,6 +23,7 @@ import org.apache.kafka.streams.kstream.KTable; import org.springframework.aop.framework.ProxyFactory; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binding.AbstractBindingTargetFactory; +import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.util.Assert; @@ -45,8 +46,12 @@ class KTableBoundElementFactory extends AbstractBindingTargetFactory { @Override public KTable createInput(String name) { - ConsumerProperties consumerProperties = this.bindingServiceProperties - .getConsumerProperties(name); + BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(name); + ConsumerProperties consumerProperties = bindingProperties.getConsumer(); + if (consumerProperties == null) { + consumerProperties = this.bindingServiceProperties.getConsumerProperties(name); + consumerProperties.setUseNativeDecoding(true); + } // Always set multiplex to true in the kafka streams binder consumerProperties.setMultiplex(true); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index baa7697d8..45178a707 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -51,7 +51,6 @@ import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.core.FluxedConsumer; import org.springframework.cloud.function.core.FluxedFunction; -import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; @@ -240,7 +239,6 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { Assert.isInstanceOf(String.class, input, "Annotation value must be a String"); Object targetBean = applicationContext.getBean(input); BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(input); - enableNativeDecodingForKTableAlways(parameterType, bindingProperties); //Retrieve the StreamsConfig created for this method if available. //Otherwise, create the StreamsBuilderFactory and get the underlying config. if (!this.methodStreamsBuilderFactoryBeanMap.containsKey(functionName)) { @@ -437,16 +435,6 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { return stream; } - private void enableNativeDecodingForKTableAlways(Class parameterType, BindingProperties bindingProperties) { - if (parameterType.isAssignableFrom(KTable.class) || parameterType.isAssignableFrom(GlobalKTable.class)) { - if (bindingProperties.getConsumer() == null) { - bindingProperties.setConsumer(new ConsumerProperties()); - } - //No framework level message conversion provided for KTable/GlobalKTable, its done by the broker. - bindingProperties.getConsumer().setUseNativeDecoding(true); - } - } - @SuppressWarnings({"unchecked"}) private void buildStreamsBuilderAndRetrieveConfig(String functionName, ApplicationContext applicationContext, String inboundName) { diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index 924afa4fd..6fc1bb652 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -49,7 +49,6 @@ import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsStateStore; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; @@ -254,7 +253,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator .getBean((String) targetReferenceValue); BindingProperties bindingProperties = this.bindingServiceProperties .getBindingProperties(inboundName); - enableNativeDecodingForKTableAlways(parameterType, bindingProperties); // Retrieve the StreamsConfig created for this method if available. // Otherwise, create the StreamsBuilderFactory and get the underlying // config. @@ -503,19 +501,6 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator return stream; } - private void enableNativeDecodingForKTableAlways(Class parameterType, - BindingProperties bindingProperties) { - if (parameterType.isAssignableFrom(KTable.class) - || parameterType.isAssignableFrom(GlobalKTable.class)) { - if (bindingProperties.getConsumer() == null) { - bindingProperties.setConsumer(new ConsumerProperties()); - } - // No framework level message conversion provided for KTable/GlobalKTable, its - // done by the broker. - bindingProperties.getConsumer().setUseNativeDecoding(true); - } - } - @SuppressWarnings({"unchecked"}) private void buildStreamsBuilderAndRetrieveConfig(Method method, ApplicationContext applicationContext, String inboundName) { diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java index 0766bf25f..55fe6222e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java @@ -87,16 +87,16 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output1.destination=counts", - "--spring.cloud.stream.bindings.output1.contentType=application/json", "--spring.cloud.stream.bindings.output2.destination=foo", - "--spring.cloud.stream.bindings.output2.contentType=application/json", "--spring.cloud.stream.bindings.output3.destination=bar", - "--spring.cloud.stream.bindings.output3.contentType=application/json", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output1.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output2.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output3.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java index e566d9187..d6e0f35aa 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountFunctionTests.java @@ -87,13 +87,13 @@ public class KafkaStreamsBinderWordCountFunctionTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.bindings.output.contentType=application/json", "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { receiveAndValidate(context); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java index 4fa4caf94..bad3df63f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java @@ -76,10 +76,6 @@ public class StreamToGlobalKTableFunctionTests { "--spring.cloud.stream.bindings.input-x.destination=customers", "--spring.cloud.stream.bindings.input-y.destination=products", "--spring.cloud.stream.bindings.output.destination=enriched-order", - "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-x.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-y.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.keySerde" + "=org.apache.kafka.common.serialization.Serdes$LongSerde", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" + diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java index ff7af6d00..5744b88d6 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToTableJoinFunctionTests.java @@ -88,9 +88,6 @@ public class StreamToTableJoinFunctionTests { "--spring.cloud.stream.bindings.input-1.destination=user-clicks-1", "--spring.cloud.stream.bindings.input-2.destination=user-regions-1", "--spring.cloud.stream.bindings.output.destination=output-topic-1", - "--spring.cloud.stream.bindings.input-1.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-2.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.keySerde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.valueSerde" + diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java index 0f3e2e2c8..0094c694e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java @@ -106,8 +106,6 @@ public abstract class DeserializationErrorHandlerByKafkaTests { // @checkstyle:off @SpringBootTest(properties = { - "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq", "spring.cloud.stream.bindings.input.group=group", "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", @@ -149,8 +147,6 @@ public abstract class DeserializationErrorHandlerByKafkaTests { // @checkstyle:off @SpringBootTest(properties = { - "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", "spring.cloud.stream.bindings.input.destination=word1,word2", "spring.cloud.stream.kafka.streams.default.consumer.applicationId=deser-kafka-dlq-multi-input", "spring.cloud.stream.bindings.input.group=groupx", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java index 620c29e69..e8586eb9b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java @@ -101,7 +101,10 @@ public abstract class DeserializtionErrorHandlerByBinderTests { consumer.close(); } - @SpringBootTest(properties = { "spring.cloud.stream.bindings.input.destination=foos", + @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=false", + "spring.cloud.stream.bindings.output.producer.useNativeEncoding=false", + "spring.cloud.stream.bindings.input.destination=foos", "spring.cloud.stream.bindings.output.destination=counts-id", "spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" @@ -148,6 +151,8 @@ public abstract class DeserializtionErrorHandlerByBinderTests { } @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=false", + "spring.cloud.stream.bindings.output.producer.useNativeEncoding=false", "spring.cloud.stream.bindings.input.destination=foos1,foos2", "spring.cloud.stream.bindings.output.destination=counts-id", "spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java index 2d526a029..a40c260eb 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderHealthIndicatorTests.java @@ -209,6 +209,12 @@ public class KafkaStreamsBinderHealthIndicatorTests { + "org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderHealthIndicatorTests.Product", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderHealthIndicatorTests.Product", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=" + "ApplicationHealthTest-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" @@ -235,6 +241,20 @@ public class KafkaStreamsBinderHealthIndicatorTests { + "org.apache.kafka.common.serialization.Serdes$IntegerSerde", "--spring.cloud.stream.kafka.streams.bindings.output2.producer.keySerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderHealthIndicatorTests.Product", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderHealthIndicatorTests.Product", + + "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output2.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output2.producer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderHealthIndicatorTests.Product", + "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderHealthIndicatorTests.Product", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=" + "ApplicationHealthTest-xyz", "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.applicationId=" @@ -302,7 +322,7 @@ public class KafkaStreamsBinderHealthIndicatorTests { } - static class Product { + public static class Product { Integer id; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderMultipleInputTopicsTest.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderMultipleInputTopicsTest.java index 9be47a52b..82ec11d62 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderMultipleInputTopicsTest.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderMultipleInputTopicsTest.java @@ -108,6 +108,7 @@ public class KafkaStreamsBinderMultipleInputTopicsTest { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests.java index b7125889a..b229c0c00 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests.java @@ -18,10 +18,10 @@ package org.springframework.cloud.stream.binder.kafka.streams.integration; import java.util.Map; -import com.fasterxml.jackson.databind.ObjectMapper; 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.common.serialization.LongDeserializer; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Materialized; @@ -63,7 +63,7 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule .getEmbeddedKafka(); - private static Consumer consumer; + private static Consumer consumer; @BeforeClass public static void setUp() throws Exception { @@ -72,7 +72,8 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { // consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, // Deserializer.class.getName()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( + consumerProps.put("value.deserializer", LongDeserializer.class); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); consumer = cf.createConsumer(); embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "counts-id"); @@ -98,6 +99,11 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { + "org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=" + + "org.apache.kafka.common.serialization.Serdes$LongSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests.Product", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=" + "KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" @@ -120,13 +126,11 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { KafkaTemplate template = new KafkaTemplate<>(pf, true); template.setDefaultTopic("foos"); template.sendDefault("{\"id\":\"123\"}"); - ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, "counts-id"); assertThat(cr.key()).isEqualTo(123); - ObjectMapper om = new ObjectMapper(); - Long aLong = om.readValue(cr.value(), Long.class); - assertThat(aLong).isEqualTo(1L); + assertThat(cr.value()).isEqualTo(1L); } @EnableBinding(KafkaStreamsProcessor.class) @@ -142,12 +146,14 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { new JsonSerde<>(Product.class))) .windowedBy(TimeWindows.of(5000)) .count(Materialized.as("id-count-store-x")).toStream() - .map((key, value) -> new KeyValue<>(key.key().id, value)); + .map((key, value) -> { + return new KeyValue<>(key.key().id, value); + }); } } - static class Product { + public static class Product { Integer id; diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderWordCountIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderWordCountIntegrationTests.java index e4018695a..7c396a255 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderWordCountIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderWordCountIntegrationTests.java @@ -108,13 +108,13 @@ public class KafkaStreamsBinderWordCountIntegrationTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.bindings.output.contentType=application/json", "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.streams.binder.brokers=" @@ -136,13 +136,13 @@ public class KafkaStreamsBinderWordCountIntegrationTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.bindings.output.contentType=application/json", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=basic-word-count", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.bindings.input.consumer.concurrency=2", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java index 786fe974d..eccfd9f72 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java @@ -99,6 +99,11 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsInteractiveQueryIntegrationTests.Product", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=ProductCountApplication-abc", "--spring.cloud.stream.kafka.streams.binder.configuration.application.server" + "=" + embeddedKafka.getBrokersAsString(), diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java index e4e075310..127f83fd0 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java @@ -105,8 +105,7 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { } @SpringBootTest(properties = { - "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=NativeEncodingDecodingEnabledTests-abc" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) public static class NativeEncodingDecodingEnabledTests @@ -132,8 +131,11 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { } // @checkstyle:off - @SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = "spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" - + "=NativeEncodingDecodingEnabledTests-xyz") + @SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = { + "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=false", + "spring.cloud.stream.bindings.output.producer.useNativeEncoding=false", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + + "=NativeEncodingDecodingEnabledTests-xyz" }) // @checkstyle:on public static class NativeEncodingDecodingDisabledTests extends KafkaStreamsNativeEncodingDecodingTests { diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java index c60800810..0d2480367 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java @@ -67,6 +67,9 @@ public class KafkaStreamsStateStoreIntegrationTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsStateStoreIntegrationTests.Product", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=KafkaStreamsStateStoreIntegrationTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" @@ -97,6 +100,12 @@ public class KafkaStreamsStateStoreIntegrationTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsStateStoreIntegrationTests.Product", + "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsStateStoreIntegrationTests.Product", "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.applicationId" + "=KafkaStreamsStateStoreIntegrationTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkastreamsBinderPojoInputStringOutputIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkastreamsBinderPojoInputStringOutputIntegrationTests.java index d4a1e88c9..33dc0aa1e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkastreamsBinderPojoInputStringOutputIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkastreamsBinderPojoInputStringOutputIntegrationTests.java @@ -99,6 +99,9 @@ public class KafkastreamsBinderPojoInputStringOutputIntegrationTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.configuration.spring.json.value.default.type=" + + "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkastreamsBinderPojoInputStringOutputIntegrationTests.Product", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=ProductCountApplication-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java index 327487777..1b866e1df 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java @@ -102,6 +102,8 @@ public class PerRecordAvroContentTypeTests { try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=false", + "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=false", "--spring.cloud.stream.bindings.input.destination=sensors", "--spring.cloud.stream.bindings.output.destination=received-sensors", "--spring.cloud.stream.bindings.output.contentType=application/avro", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java index d26ee1b50..1b4eaa5c2 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java @@ -81,10 +81,6 @@ public class StreamToGlobalKTableJoinIntegrationTests { "--spring.cloud.stream.bindings.input-x.destination=customers", "--spring.cloud.stream.bindings.input-y.destination=products", "--spring.cloud.stream.bindings.output.destination=enriched-order", - "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-x.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-y.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.keySerde" + "=org.apache.kafka.common.serialization.Serdes$LongSerde", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java index 9c0ae02c2..3d5a4888a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToTableJoinIntegrationTests.java @@ -95,9 +95,7 @@ public class StreamToTableJoinIntegrationTests { "--spring.cloud.stream.bindings.input.destination=user-clicks-1", "--spring.cloud.stream.bindings.input-x.destination=user-regions-1", "--spring.cloud.stream.bindings.output.destination=output-topic-1", - "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-x.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.keySerde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" @@ -246,9 +244,7 @@ public class StreamToTableJoinIntegrationTests { "--spring.cloud.stream.bindings.input.destination=user-clicks-2", "--spring.cloud.stream.bindings.input-x.destination=user-regions-2", "--spring.cloud.stream.bindings.output.destination=output-topic-2", - "--spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.input-x.consumer.useNativeDecoding=true", - "--spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", + "--spring.cloud.stream.kafka.streams.binder.configuration.auto.offset.reset=latest", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.startOffset=earliest", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.keySerde" diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java index 7be37368a..5c763e72b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java @@ -96,16 +96,16 @@ public class WordCountMultipleBranchesIntegrationTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output1.destination=counts", - "--spring.cloud.stream.bindings.output1.contentType=application/json", "--spring.cloud.stream.bindings.output2.destination=foo", - "--spring.cloud.stream.bindings.output2.contentType=application/json", "--spring.cloud.stream.bindings.output3.destination=bar", - "--spring.cloud.stream.bindings.output3.contentType=application/json", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.output1.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output2.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", + "--spring.cloud.stream.kafka.streams.bindings.output3.producer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId"