diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index a3c19ce8d..4d38188e7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -1,5 +1,5 @@ /* - * Copyright 2017-2018 the original author or authors. + * Copyright 2017-2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -21,6 +21,7 @@ import java.util.Map; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; @@ -120,9 +121,15 @@ class KStreamBinder extends this.kafkaTopicProvisioner.provisionProducerDestination(name, extendedProducerProperties); Serde keySerde = this.keyValueSerdeResolver - .getOuboundKeySerde(properties.getExtension()); - Serde valueSerde = this.keyValueSerdeResolver.getOutboundValueSerde(properties, - properties.getExtension()); + .getOuboundKeySerde(properties.getExtension(), kafkaStreamsBindingInformationCatalogue.getOutboundKStreamResolvable()); + Serde valueSerde; + if (properties.isUseNativeEncoding()) { + valueSerde = this.keyValueSerdeResolver.getOutboundValueSerde(properties, + properties.getExtension(), kafkaStreamsBindingInformationCatalogue.getOutboundKStreamResolvable()); + } + else { + valueSerde = Serdes.ByteArray(); + } to(properties.isUseNativeEncoding(), name, outboundBindTarget, (Serde) keySerde, (Serde) valueSerde); return new DefaultBinding<>(name, null, outboundBindTarget, null); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java index b9bb99134..52eb740da 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java @@ -27,6 +27,7 @@ import org.apache.kafka.streams.kstream.KStream; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.core.ResolvableType; import org.springframework.kafka.config.StreamsBuilderFactoryBean; /** @@ -46,6 +47,8 @@ class KafkaStreamsBindingInformationCatalogue { private final Set streamsBuilderFactoryBeans = new HashSet<>(); + private ResolvableType outboundKStreamResolvable; + /** * For a given bounded {@link KStream}, retrieve it's corresponding destination on the * broker. @@ -122,4 +125,11 @@ class KafkaStreamsBindingInformationCatalogue { return this.streamsBuilderFactoryBeans; } + public void setOutboundKStreamResolvable(ResolvableType outboundResolvable) { + this.outboundKStreamResolvable = outboundResolvable; + } + + public ResolvableType getOutboundKStreamResolvable() { + return outboundKStreamResolvable; + } } 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 45178a707..e6180e7bc 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 @@ -30,6 +30,7 @@ import java.util.function.Function; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; @@ -94,6 +95,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { private Set origInputs = new TreeSet<>(); private Set origOutputs = new TreeSet<>(); + private ResolvableType outboundResolvableType; + public KafkaStreamsFunctionProcessor(BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, @@ -132,18 +135,20 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { resolvableTypeMap.put(next, resolvableType.getGeneric(0)); origInputs.remove(next); + ResolvableType iterableResType = resolvableType; for (int i = 1; i < inputCount; i++) { if (iterator.hasNext()) { - ResolvableType generic = resolvableType.getGeneric(1); - if (generic.getRawClass() != null && - (generic.getRawClass().equals(Function.class) || - generic.getRawClass().equals(Consumer.class))) { + iterableResType = iterableResType.getGeneric(1); + if (iterableResType.getRawClass() != null && + (iterableResType.getRawClass().equals(Function.class) || + iterableResType.getRawClass().equals(Consumer.class))) { final String next1 = iterator.next(); - resolvableTypeMap.put(next1, generic.getGeneric(0)); + resolvableTypeMap.put(next1, iterableResType.getGeneric(0)); origInputs.remove(next1); } } } + outboundResolvableType = iterableResType.getGeneric(1); return resolvableTypeMap; } @@ -184,6 +189,8 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { i++; } if (result != null) { + kafkaStreamsBindingInformationCatalogue.setOutboundKStreamResolvable( + outboundResolvableType != null ? outboundResolvableType : resolvableType.getGeneric(1)); final Set outputs = new TreeSet<>(origOutputs); final Iterator iterator = outputs.iterator(); @@ -251,9 +258,17 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(input); //get state store spec - Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties); - Serde valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( - bindingProperties.getConsumer(), extendedConsumerProperties); + + Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties, stringResolvableTypeMap.get(input)); + Serde valueSerde; + + if (bindingServiceProperties.getConsumerProperties(input).isUseNativeDecoding()) { + valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( + bindingProperties.getConsumer(), extendedConsumerProperties, stringResolvableTypeMap.get(input)); + } + else { + valueSerde = Serdes.ByteArray(); + } final KafkaConsumerProperties.StartOffset startOffset = extendedConsumerProperties.getStartOffset(); Topology.AutoOffsetReset autoOffsetReset = null; 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 6fc1bb652..56dc608a4 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 @@ -28,6 +28,7 @@ import java.util.Properties; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; @@ -64,6 +65,7 @@ import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.core.MethodParameter; +import org.springframework.core.ResolvableType; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBean; @@ -190,6 +192,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator Assert.isTrue(methodAnnotatedOutboundNames.length == 1, "Result does not match with the number of declared outbounds"); } + kafkaStreamsBindingInformationCatalogue.setOutboundKStreamResolvable(ResolvableType.forMethodReturnType(method)); if (result.getClass().isArray()) { Object[] outboundKStreams = (Object[]) result; int i = 0; @@ -268,10 +271,19 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator .getExtendedConsumerProperties(inboundName); // get state store spec KafkaStreamsStateStoreProperties spec = buildStateStoreSpec(method); + Serde keySerde = this.keyValueSerdeResolver - .getInboundKeySerde(extendedConsumerProperties); - Serde valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( - bindingProperties.getConsumer(), extendedConsumerProperties); + .getInboundKeySerde(extendedConsumerProperties, ResolvableType.forMethodParameter(methodParameter)); + Serde valueSerde; + + if (bindingServiceProperties.getConsumerProperties(inboundName).isUseNativeDecoding()) { + valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( + bindingProperties.getConsumer(), extendedConsumerProperties, ResolvableType.forMethodParameter(methodParameter)); + } + else { + //keySerde = Serdes.ByteArray(); + valueSerde = Serdes.ByteArray(); + } final KafkaConsumerProperties.StartOffset startOffset = extendedConsumerProperties .getStartOffset(); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java index d2e13b31e..08c2648c9 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KeyValueSerdeResolver.java @@ -21,12 +21,17 @@ import java.util.Map; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.streams.kstream.GlobalKTable; +import org.apache.kafka.streams.kstream.KStream; +import org.apache.kafka.streams.kstream.KTable; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsProducerProperties; +import org.springframework.core.ResolvableType; +import org.springframework.kafka.support.serializer.JsonSerde; import org.springframework.util.StringUtils; /** @@ -78,6 +83,13 @@ public class KeyValueSerdeResolver { return getKeySerde(keySerdeString); } + public Serde getInboundKeySerde( + KafkaStreamsConsumerProperties extendedConsumerProperties, ResolvableType resolvableType) { + String keySerdeString = extendedConsumerProperties.getKeySerde(); + + return getKeySerde(keySerdeString, resolvableType); + } + /** * Provide the {@link Serde} for inbound value. * @param consumerProperties {@link ConsumerProperties} on binding @@ -105,6 +117,27 @@ public class KeyValueSerdeResolver { return valueSerde; } + public Serde getInboundValueSerde(ConsumerProperties consumerProperties, + KafkaStreamsConsumerProperties extendedConsumerProperties, + ResolvableType resolvableType) { + Serde valueSerde; + + String valueSerdeString = extendedConsumerProperties.getValueSerde(); + try { + if (consumerProperties != null && consumerProperties.isUseNativeDecoding()) { + valueSerde = getValueSerde(valueSerdeString, resolvableType); + } + else { + valueSerde = Serdes.ByteArray(); + } + valueSerde.configure(this.streamConfigGlobalProperties, false); + } + catch (ClassNotFoundException ex) { + throw new IllegalStateException("Serde class not found: ", ex); + } + return valueSerde; + } + /** * Provide the {@link Serde} for outbound key. * @param properties binding level extended {@link KafkaStreamsProducerProperties} @@ -114,6 +147,11 @@ public class KeyValueSerdeResolver { return getKeySerde(properties.getKeySerde()); } + public Serde getOuboundKeySerde(KafkaStreamsProducerProperties properties, ResolvableType resolvableType) { + return getKeySerde(properties.getKeySerde(), resolvableType); + } + + /** * Provide the {@link Serde} for outbound value. * @param producerProperties {@link ProducerProperties} on binding @@ -140,6 +178,25 @@ public class KeyValueSerdeResolver { return valueSerde; } + public Serde getOutboundValueSerde(ProducerProperties producerProperties, + KafkaStreamsProducerProperties kafkaStreamsProducerProperties, ResolvableType resolvableType) { + Serde valueSerde; + try { + if (producerProperties.isUseNativeEncoding()) { + valueSerde = getValueSerde( + kafkaStreamsProducerProperties.getValueSerde(), resolvableType); + } + else { + valueSerde = Serdes.ByteArray(); + } + valueSerde.configure(this.streamConfigGlobalProperties, false); + } + catch (ClassNotFoundException ex) { + throw new IllegalStateException("Serde class not found: ", ex); + } + return valueSerde; + } + /** * Provide the {@link Serde} for state store. * @param keySerdeString serde class used for key @@ -186,6 +243,76 @@ public class KeyValueSerdeResolver { return keySerde; } + private Serde getKeySerde(String keySerdeString, ResolvableType resolvableType) { + Serde keySerde = null; + try { + if (StringUtils.hasText(keySerdeString)) { + keySerde = Utils.newInstance(keySerdeString, Serde.class); + } + else { + if (resolvableType != null && + (isResolvalbeKafkaStreamsType(resolvableType) || isResolvableKStreamArrayType(resolvableType))) { + ResolvableType generic = resolvableType.isArray() ? resolvableType.getComponentType().getGeneric(0) : resolvableType.getGeneric(0); + keySerde = getSerde(keySerde, generic); + } + if (keySerde == null) { + keySerde = this.binderConfigurationProperties.getConfiguration() + .containsKey("default.key.serde") + ? Utils.newInstance(this.binderConfigurationProperties + .getConfiguration().get("default.key.serde"), + Serde.class) + : Serdes.ByteArray(); + } + } + keySerde.configure(this.streamConfigGlobalProperties, true); + } + catch (ClassNotFoundException ex) { + throw new IllegalStateException("Serde class not found: ", ex); + } + return keySerde; + } + + private boolean isResolvableKStreamArrayType(ResolvableType resolvableType) { + return resolvableType.isArray() && + KStream.class.isAssignableFrom(resolvableType.getComponentType().getRawClass()); + } + + private boolean isResolvalbeKafkaStreamsType(ResolvableType resolvableType) { + return resolvableType.getRawClass() != null && (KStream.class.isAssignableFrom(resolvableType.getRawClass()) || KTable.class.isAssignableFrom(resolvableType.getRawClass()) || + GlobalKTable.class.isAssignableFrom(resolvableType.getRawClass())); + } + + private Serde getSerde(Serde keySerde, ResolvableType generic) { + if (generic.getRawClass() != null) { + if (Integer.class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.Integer(); + } + else if (Long.class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.Long(); + } + else if (Short.class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.Short(); + } + else if (Double.class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.Double(); + } + else if (Float.class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.Float(); + } + else if (byte[].class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.ByteArray(); + } + else if (String.class.isAssignableFrom(generic.getRawClass())) { + keySerde = Serdes.String(); + } + else { + keySerde = new JsonSerde(generic.getRawClass()); + } + } + return keySerde; + } + + private Serde getValueSerde(String valueSerdeString) throws ClassNotFoundException { Serde valueSerde; @@ -203,4 +330,32 @@ public class KeyValueSerdeResolver { return valueSerde; } + @SuppressWarnings("unchecked") + private Serde getValueSerde(String valueSerdeString, ResolvableType resolvableType) + throws ClassNotFoundException { + Serde valueSerde = null; + if (StringUtils.hasText(valueSerdeString)) { + valueSerde = Utils.newInstance(valueSerdeString, Serde.class); + } + else { + + if (resolvableType != null && ((isResolvalbeKafkaStreamsType(resolvableType)) || + (isResolvableKStreamArrayType(resolvableType)))) { + ResolvableType generic = resolvableType.isArray() ? resolvableType.getComponentType().getGeneric(1) : resolvableType.getGeneric(1); + valueSerde = getSerde(valueSerde, generic); + } + + if (valueSerde == null) { + + valueSerde = this.binderConfigurationProperties.getConfiguration() + .containsKey("default.value.serde") + ? Utils.newInstance(this.binderConfigurationProperties + .getConfiguration().get("default.value.serde"), + Serde.class) + : Serdes.ByteArray(); + } + } + return valueSerde; + } + } 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 55fe6222e..995b7c50a 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 @@ -94,9 +94,9 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { "=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.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 d6e0f35aa..b7095e722 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 @@ -93,7 +93,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { "=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.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 bad3df63f..594e46975 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 @@ -49,7 +49,7 @@ import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.serializer.JsonDeserializer; -import org.springframework.kafka.support.serializer.JsonSerde; +import org.springframework.kafka.support.serializer.JsonSerializer; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.kafka.test.utils.KafkaTestUtils; @@ -76,26 +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.kafka.streams.bindings.input.consumer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" + - "=org.springframework.cloud.stream.binder.kafka.streams.function" + - ".StreamToGlobalKTableFunctionTests$OrderSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.valueSerde" + - "=org.springframework.cloud.stream.binder.kafka.streams.function" + - ".StreamToGlobalKTableFunctionTests$CustomerSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.valueSerde" + - "=org.springframework.cloud.stream.binder.kafka.streams.function" + - ".StreamToGlobalKTableFunctionTests$ProductSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" + - "=org.springframework.cloud.stream.binder.kafka.streams." + - "function.StreamToGlobalKTableFunctionTests$EnrichedOrderSerde", "--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" + @@ -107,9 +87,8 @@ public class StreamToGlobalKTableFunctionTests { "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + embeddedKafka.getZookeeperConnectionString())) { Map senderPropsCustomer = KafkaTestUtils.producerProps(embeddedKafka); senderPropsCustomer.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - CustomerSerde customerSerde = new CustomerSerde(); senderPropsCustomer.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, - customerSerde.serializer().getClass()); + JsonSerializer.class); DefaultKafkaProducerFactory pfCustomer = new DefaultKafkaProducerFactory<>(senderPropsCustomer); @@ -123,8 +102,7 @@ public class StreamToGlobalKTableFunctionTests { Map senderPropsProduct = KafkaTestUtils.producerProps(embeddedKafka); senderPropsProduct.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - ProductSerde productSerde = new ProductSerde(); - senderPropsProduct.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, productSerde.serializer().getClass()); + senderPropsProduct.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); DefaultKafkaProducerFactory pfProduct = new DefaultKafkaProducerFactory<>(senderPropsProduct); @@ -139,8 +117,7 @@ public class StreamToGlobalKTableFunctionTests { Map senderPropsOrder = KafkaTestUtils.producerProps(embeddedKafka); senderPropsOrder.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - OrderSerde orderSerde = new OrderSerde(); - senderPropsOrder.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, orderSerde.serializer().getClass()); + senderPropsOrder.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); DefaultKafkaProducerFactory pfOrder = new DefaultKafkaProducerFactory<>(senderPropsOrder); KafkaTemplate orderTemplate = new KafkaTemplate<>(pfOrder, true); @@ -157,9 +134,8 @@ public class StreamToGlobalKTableFunctionTests { embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); - EnrichedOrderSerde enrichedOrderSerde = new EnrichedOrderSerde(); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, - enrichedOrderSerde.deserializer().getClass()); + JsonDeserializer.class); consumerProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "org.springframework.cloud.stream.binder.kafka.streams." + "function.StreamToGlobalKTableFunctionTests.EnrichedOrder"); @@ -334,15 +310,4 @@ public class StreamToGlobalKTableFunctionTests { } } - public static class OrderSerde extends JsonSerde { - } - - public static class CustomerSerde extends JsonSerde { - } - - public static class ProductSerde extends JsonSerde { - } - - public static class EnrichedOrderSerde extends JsonSerde { - } } 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 5744b88d6..d48e83d9c 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,18 +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.kafka.streams.bindings.input-1.consumer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.valueSerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.valueSerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", "--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" + @@ -236,18 +224,6 @@ public class StreamToTableJoinFunctionTests { "--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-1.consumer.startOffset=earliest", - "--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" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.valueSerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" + - "=org.apache.kafka.common.serialization.Serdes$LongSerde", "--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" + 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 0094c694e..18d348fce 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 @@ -109,7 +109,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { "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", - "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=" + "spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) // @checkstyle:on public static class DeserializationByKafkaAndDlqTests @@ -148,10 +148,11 @@ public abstract class DeserializationErrorHandlerByKafkaTests { // @checkstyle:off @SpringBootTest(properties = { "spring.cloud.stream.bindings.input.destination=word1,word2", - "spring.cloud.stream.kafka.streams.default.consumer.applicationId=deser-kafka-dlq-multi-input", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq-multi-input", + //"spring.cloud.stream.kafka.streams.default.consumer.applicationId=deser-kafka-dlq-multi-input", "spring.cloud.stream.bindings.input.group=groupx", "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", - "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=" + "spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) // @checkstyle:on public static class DeserializationByKafkaAndDlqTestsWithMultipleInputs 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 e8586eb9b..5a101ecc2 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 @@ -111,8 +111,8 @@ public abstract class DeserializtionErrorHandlerByBinderTests { + "=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.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", +// "spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id" + "=deserializationByBinderAndDlqTests", @@ -160,8 +160,8 @@ public abstract class DeserializtionErrorHandlerByBinderTests { + "=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.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", +// "spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id" + "=deserializationByBinderAndDlqTestsWithMultipleInputs", 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 a40c260eb..2e0fbc4b7 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 @@ -207,14 +207,14 @@ public class KafkaStreamsBinderHealthIndicatorTests { + "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.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.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=" @@ -237,23 +237,23 @@ public class KafkaStreamsBinderHealthIndicatorTests { + "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.keySerde=" - + "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.output.producer.keySerde=" +// + "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", 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 82ec11d62..6465427a9 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,7 +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.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 b229c0c00..36d67a5e8 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 @@ -97,13 +97,13 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { + "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.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.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=" 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 7c396a255..32b80f58d 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 @@ -114,7 +114,7 @@ public class KafkaStreamsBinderWordCountIntegrationTests { + "=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.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=" 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 eccfd9f72..8c00ac8e3 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 @@ -100,9 +100,9 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { "--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.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" 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 0d2480367..0ace06389 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,9 +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.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=" @@ -100,12 +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.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 33dc0aa1e..e4abf7650 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 @@ -97,11 +97,11 @@ public class KafkastreamsBinderPojoInputStringOutputIntegrationTests { + "=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.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.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/StreamToGlobalKTableJoinIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/StreamToGlobalKTableJoinIntegrationTests.java index 1b4eaa5c2..3d726f90e 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 @@ -48,7 +48,7 @@ import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.serializer.JsonDeserializer; -import org.springframework.kafka.support.serializer.JsonSerde; +import org.springframework.kafka.support.serializer.JsonSerializer; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.rule.EmbeddedKafkaRule; import org.springframework.kafka.test.utils.KafkaTestUtils; @@ -81,26 +81,26 @@ 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.kafka.streams.bindings.input.consumer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" - + "=org.springframework.cloud.stream.binder.kafka.streams.integration" - + ".StreamToGlobalKTableJoinIntegrationTests$OrderSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.valueSerde" - + "=org.springframework.cloud.stream.binder.kafka.streams.integration." - + "StreamToGlobalKTableJoinIntegrationTests$CustomerSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.valueSerde" - + "=org.springframework.cloud.stream.binder.kafka.streams." - + "integration.StreamToGlobalKTableJoinIntegrationTests$ProductSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" - + "=org.springframework.cloud.stream.binder.kafka.streams.integration" - + ".StreamToGlobalKTableJoinIntegrationTests$EnrichedOrderSerde", +// "--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" +// + "=org.springframework.cloud.stream.binder.kafka.streams.integration" +// + ".StreamToGlobalKTableJoinIntegrationTests$OrderSerde", +// "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.valueSerde" +// + "=org.springframework.cloud.stream.binder.kafka.streams.integration." +// + "StreamToGlobalKTableJoinIntegrationTests$CustomerSerde", +// "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.valueSerde" +// + "=org.springframework.cloud.stream.binder.kafka.streams." +// + "integration.StreamToGlobalKTableJoinIntegrationTests$ProductSerde", +// "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" +// + "=org.springframework.cloud.stream.binder.kafka.streams.integration" +// + ".StreamToGlobalKTableJoinIntegrationTests$EnrichedOrderSerde", "--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" @@ -116,9 +116,8 @@ public class StreamToGlobalKTableJoinIntegrationTests { .producerProps(embeddedKafka); senderPropsCustomer.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - CustomerSerde customerSerde = new CustomerSerde(); senderPropsCustomer.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, - customerSerde.serializer().getClass()); + JsonSerializer.class); DefaultKafkaProducerFactory pfCustomer = new DefaultKafkaProducerFactory<>( senderPropsCustomer); @@ -135,9 +134,8 @@ public class StreamToGlobalKTableJoinIntegrationTests { .producerProps(embeddedKafka); senderPropsProduct.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - ProductSerde productSerde = new ProductSerde(); senderPropsProduct.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, - productSerde.serializer().getClass()); + JsonSerializer.class); DefaultKafkaProducerFactory pfProduct = new DefaultKafkaProducerFactory<>( senderPropsProduct); @@ -155,9 +153,8 @@ public class StreamToGlobalKTableJoinIntegrationTests { .producerProps(embeddedKafka); senderPropsOrder.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); - OrderSerde orderSerde = new OrderSerde(); senderPropsOrder.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, - orderSerde.serializer().getClass()); + JsonSerializer.class); DefaultKafkaProducerFactory pfOrder = new DefaultKafkaProducerFactory<>( senderPropsOrder); @@ -176,9 +173,8 @@ public class StreamToGlobalKTableJoinIntegrationTests { consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); - EnrichedOrderSerde enrichedOrderSerde = new EnrichedOrderSerde(); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, - enrichedOrderSerde.deserializer().getClass()); + JsonDeserializer.class); consumerProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "org.springframework.cloud.stream.binder.kafka.streams.integration." + "StreamToGlobalKTableJoinIntegrationTests.EnrichedOrder"); @@ -368,20 +364,20 @@ public class StreamToGlobalKTableJoinIntegrationTests { } - public static class OrderSerde extends JsonSerde { - - } - - public static class CustomerSerde extends JsonSerde { - - } - - public static class ProductSerde extends JsonSerde { - - } - - public static class EnrichedOrderSerde extends JsonSerde { - - } +// public static class OrderSerde extends JsonSerde { +// +// } +// +// public static class CustomerSerde extends JsonSerde { +// +// } +// +// public static class ProductSerde extends JsonSerde { +// +// } +// +// public static class EnrichedOrderSerde extends JsonSerde { +// +// } } 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 3d5a4888a..e9ef13249 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 @@ -96,18 +96,18 @@ public class StreamToTableJoinIntegrationTests { "--spring.cloud.stream.bindings.input-x.destination=user-regions-1", "--spring.cloud.stream.bindings.output.destination=output-topic-1", - "--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" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.valueSerde" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--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" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$StringSerde", +// "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.valueSerde" +// + "=org.apache.kafka.common.serialization.Serdes$StringSerde", +// "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$StringSerde", +// "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", "--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" @@ -247,18 +247,18 @@ public class StreamToTableJoinIntegrationTests { "--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" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.valueSerde" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" - + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" - + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--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" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", +// "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$StringSerde", +// "--spring.cloud.stream.kafka.streams.bindings.inputX.consumer.valueSerde" +// + "=org.apache.kafka.common.serialization.Serdes$StringSerde", +// "--spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" +// + "=org.apache.kafka.common.serialization.Serdes$StringSerde", +// "--spring.cloud.stream.kafka.streams.bindings.output.producer.valueSerde" +// + "=org.apache.kafka.common.serialization.Serdes$LongSerde", "--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" 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 5c763e72b..1f3ea183d 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 @@ -103,9 +103,9 @@ public class WordCountMultipleBranchesIntegrationTests { + "=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.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"