diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java index 1969358b9..71655bc21 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java @@ -118,4 +118,9 @@ public class GlobalKTableBinder extends .getExtendedPropertiesEntryClass(); } + public void setKafkaStreamsExtendedBindingProperties( + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { + this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; + } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java index d08988bdc..b1ebedd7a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java @@ -25,6 +25,7 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -54,9 +55,13 @@ public class GlobalKTableBinderConfiguration { public GlobalKTableBinder GlobalKTableBinder( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, @Qualifier("kafkaStreamsDlqDispatchers") Map kafkaStreamsDlqDispatchers) { - return new GlobalKTableBinder(binderConfigurationProperties, + GlobalKTableBinder globalKTableBinder = new GlobalKTableBinder(binderConfigurationProperties, kafkaTopicProvisioner, kafkaStreamsDlqDispatchers); + globalKTableBinder.setKafkaStreamsExtendedBindingProperties( + kafkaStreamsExtendedBindingProperties); + return globalKTableBinder; } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java index 2e6e3dd14..1823acd4a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java @@ -122,4 +122,8 @@ class KTableBinder extends .getExtendedPropertiesEntryClass(); } + public void setKafkaStreamsExtendedBindingProperties( + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { + this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java index 691626f31..286aaab52 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java @@ -25,6 +25,7 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsExtendedBindingProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -54,10 +55,12 @@ public class KTableBinderConfiguration { public KTableBinder kTableBinder( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, @Qualifier("kafkaStreamsDlqDispatchers") Map kafkaStreamsDlqDispatchers) { - KTableBinder kStreamBinder = new KTableBinder(binderConfigurationProperties, + KTableBinder kTableBinder = new KTableBinder(binderConfigurationProperties, kafkaTopicProvisioner, kafkaStreamsDlqDispatchers); - return kStreamBinder; + kTableBinder.setKafkaStreamsExtendedBindingProperties(kafkaStreamsExtendedBindingProperties); + return kTableBinder; } } 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 6465427a9..9be47a52b 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,6 @@ 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 36d67a5e8..62c2c9a69 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,6 @@ 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.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 32b80f58d..e813fad60 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,6 @@ 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.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 8c00ac8e3..786fe974d 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,11 +99,6 @@ 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/KafkaStreamsStateStoreIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsStateStoreIntegrationTests.java index 0ace06389..ccf562430 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,6 @@ 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=" 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 e4abf7650..b6f117cca 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,6 @@ 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.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 3d726f90e..bf2bcbe20 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 @@ -41,8 +41,14 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.BinderFactory; +import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; +import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -75,32 +81,12 @@ public class StreamToGlobalKTableJoinIntegrationTests { SpringApplication app = new SpringApplication( StreamToGlobalKTableJoinIntegrationTests.OrderEnricherApplication.class); app.setWebApplicationType(WebApplicationType.NONE); - try (ConfigurableApplicationContext ignored = app.run("--server.port=0", + ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=orders", "--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.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" @@ -108,10 +94,44 @@ public class StreamToGlobalKTableJoinIntegrationTests { "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=StreamToGlobalKTableJoinIntegrationTests-abc", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.topic.properties.cleanup.policy=compact", + "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.topic.properties.cleanup.policy=compact", + "--spring.cloud.stream.kafka.streams.bindings.input-y.consumer.topic.properties.cleanup.policy=compact", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString())) { + + embeddedKafka.getZookeeperConnectionString()); + try { + // Testing certain ancillary configuration of GlobalKTable around topics creation. + // See this issue: https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/687 + + BinderFactory binderFactory = context.getBeanFactory() + .getBean(BinderFactory.class); + + Binder kStreamBinder = binderFactory + .getBinder("kstream", KStream.class); + + KafkaStreamsConsumerProperties input = (KafkaStreamsConsumerProperties) ((ExtendedPropertiesBinder) kStreamBinder) + .getExtendedConsumerProperties("input"); + String cleanupPolicy = input.getTopic().getProperties().get("cleanup.policy"); + + assertThat(cleanupPolicy).isEqualTo("compact"); + + Binder globalKTableBinder = binderFactory + .getBinder("globalktable", GlobalKTable.class); + + KafkaStreamsConsumerProperties inputX = (KafkaStreamsConsumerProperties) ((ExtendedPropertiesBinder) globalKTableBinder) + .getExtendedConsumerProperties("input-x"); + String cleanupPolicyX = inputX.getTopic().getProperties().get("cleanup.policy"); + + assertThat(cleanupPolicyX).isEqualTo("compact"); + + KafkaStreamsConsumerProperties inputY = (KafkaStreamsConsumerProperties) ((ExtendedPropertiesBinder) globalKTableBinder) + .getExtendedConsumerProperties("input-y"); + String cleanupPolicyY = inputY.getTopic().getProperties().get("cleanup.policy"); + + assertThat(cleanupPolicyY).isEqualTo("compact"); + Map senderPropsCustomer = KafkaTestUtils .producerProps(embeddedKafka); senderPropsCustomer.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, @@ -219,7 +239,9 @@ public class StreamToGlobalKTableJoinIntegrationTests { pfOrder.destroy(); consumer.close(); } - + finally { + context.close(); + } } interface CustomGlobalKTableProcessor extends KafkaStreamsProcessor { @@ -363,21 +385,4 @@ 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 { -// -// } - } 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 e9ef13249..65e1efba3 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 @@ -46,8 +46,14 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.BinderFactory; +import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; +import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; +import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsConsumerProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -90,24 +96,11 @@ public class StreamToTableJoinIntegrationTests { consumer = cf.createConsumer(); embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "output-topic-1"); - try (ConfigurableApplicationContext ignored = app.run("--server.port=0", + ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", "--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.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" @@ -115,10 +108,25 @@ public class StreamToTableJoinIntegrationTests { "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=StreamToTableJoinIntegrationTests-abc", + "--spring.cloud.stream.kafka.streams.bindings.input-x.consumer.topic.properties.cleanup.policy=compact", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString())) { + + embeddedKafka.getZookeeperConnectionString()); + try { + // Testing certain ancillary configuration of GlobalKTable around topics creation. + // See this issue: https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/687 + BinderFactory binderFactory = context.getBeanFactory() + .getBean(BinderFactory.class); + + Binder ktableBinder = binderFactory + .getBinder("ktable", KTable.class); + + KafkaStreamsConsumerProperties inputX = (KafkaStreamsConsumerProperties) ((ExtendedPropertiesBinder) ktableBinder) + .getExtendedConsumerProperties("input-x"); + String cleanupPolicyX = inputX.getTopic().getProperties().get("cleanup.policy"); + + assertThat(cleanupPolicyX).isEqualTo("compact"); // Input 1: Region per user (multiple records allowed per user). List> userRegions = Arrays.asList(new KeyValue<>( @@ -244,21 +252,8 @@ 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.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.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 1f3ea183d..7cba27483 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,6 @@ 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.timeWindow.length=5000", "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId"