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 690fb7bb5..23c7be3a1 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 @@ -37,8 +37,6 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -178,7 +176,6 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { } @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class WordCountProcessorApplication { @Bean 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 408331d7e..52674a3ee 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 @@ -33,14 +33,11 @@ import org.apache.kafka.streams.kstream.TimeWindows; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; -import org.junit.Ignore; import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -56,7 +53,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, - "counts"); + "counts", "counts-1"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); @@ -69,7 +66,7 @@ public class KafkaStreamsBinderWordCountFunctionTests { consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "counts"); + embeddedKafka.consumeFromEmbeddedTopics(consumer, "counts", "counts-1"); } @AfterClass @@ -78,7 +75,6 @@ public class KafkaStreamsBinderWordCountFunctionTests { } @Test - @Ignore public void testKstreamWordCountFunction() throws Exception { SpringApplication app = new SpringApplication(WordCountProcessorApplication.class); app.setWebApplicationType(WebApplicationType.NONE); @@ -90,14 +86,14 @@ public class KafkaStreamsBinderWordCountFunctionTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count", + "--spring.cloud.stream.kafka.streams.default.consumer.application-id=testKstreamWordCountFunction", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { - receiveAndValidate(context); + receiveAndValidate("words", "counts"); } } @@ -111,26 +107,26 @@ public class KafkaStreamsBinderWordCountFunctionTests { "--spring.cloud.stream.function.outputBindings.process=output", "--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input.destination=words", - "--spring.cloud.stream.bindings.output.destination=counts", + "--spring.cloud.stream.bindings.input.destination=words-1", + "--spring.cloud.stream.bindings.output.destination=counts-1", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { - receiveAndValidate(context); + receiveAndValidate("words-1", "counts-1"); } } - private void receiveAndValidate(ConfigurableApplicationContext context) throws Exception { + private void receiveAndValidate(String in, String out) { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); try { KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); + template.setDefaultTopic(in); template.sendDefault("foobar"); - ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, "counts"); + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, out); assertThat(cr.value().contains("\"word\":\"foobar\",\"count\":1")).isTrue(); } finally { @@ -189,7 +185,6 @@ public class KafkaStreamsBinderWordCountFunctionTests { } @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class WordCountProcessorApplication { @Bean diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java index 88bb43bc8..76b0f0869 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java @@ -59,7 +59,7 @@ public class KafkaStreamsFunctionStateStoreTests { try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.process_in.destination=words", - "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count-1", + "--spring.cloud.stream.kafka.streams.default.consumer.application-id=testKafkaStreamsFuncionWithMultipleStateStores", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", 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 59b718ce3..2d8216b87 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 @@ -47,8 +47,6 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -437,7 +435,6 @@ public class StreamToTableJoinFunctionTests { } @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class CountClicksPerRegionApplication { @Bean @@ -455,7 +452,6 @@ public class StreamToTableJoinFunctionTests { } @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class BiFunctionCountClicksPerRegionApplication { @Bean 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 18d348fce..d9953c900 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 @@ -70,7 +70,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, - "counts", "error.words.group", "error.word1.groupx", "error.word2.groupx"); + "DeserializationErrorHandlerByKafkaTests-out", "error.DeserializationErrorHandlerByKafkaTests-In.group", "error.word1.groupx", "error.word2.groupx"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule .getEmbeddedKafka(); @@ -84,8 +84,6 @@ public abstract class DeserializationErrorHandlerByKafkaTests { public static void setUp() throws Exception { System.setProperty("spring.cloud.stream.kafka.streams.binder.brokers", embeddedKafka.getBrokersAsString()); - System.setProperty("spring.cloud.stream.kafka.streams.binder.zkNodes", - embeddedKafka.getZookeeperConnectionString()); System.setProperty("server.port", "0"); System.setProperty("spring.jmx.enabled", "false"); @@ -96,22 +94,23 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "counts"); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "DeserializationErrorHandlerByKafkaTests-out"); } @AfterClass public static void tearDown() { consumer.close(); + System.clearProperty("spring.cloud.stream.kafka.streams.binder.brokers"); + System.clearProperty("server.port"); + System.clearProperty("spring.jmx.enabled"); } - // @checkstyle:off @SpringBootTest(properties = { "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.bindings.input.consumer.valueSerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) - // @checkstyle:on public static class DeserializationByKafkaAndDlqTests extends DeserializationErrorHandlerByKafkaTests { @@ -122,7 +121,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); + template.setDefaultTopic("DeserializationErrorHandlerByKafkaTests-In"); template.sendDefault("foobar"); Map consumerProps = KafkaTestUtils.consumerProps("foobar", @@ -131,10 +130,10 @@ public abstract class DeserializationErrorHandlerByKafkaTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); Consumer consumer1 = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.words.group"); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.DeserializationErrorHandlerByKafkaTests-In.group"); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer1, - "error.words.group"); + "error.DeserializationErrorHandlerByKafkaTests-In.group"); assertThat(cr.value().equals("foobar")).isTrue(); // Ensuring that the deserialization was indeed done by Kafka natively @@ -145,11 +144,9 @@ public abstract class DeserializationErrorHandlerByKafkaTests { } - // @checkstyle:off @SpringBootTest(properties = { "spring.cloud.stream.bindings.input.destination=word1,word2", "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.bindings.input.consumer.valueSerde=" 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 bfa4990be..fef08cbfd 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 @@ -79,16 +79,11 @@ public abstract class DeserializtionErrorHandlerByBinderTests { public static void setUp() throws Exception { System.setProperty("spring.cloud.stream.kafka.streams.binder.brokers", embeddedKafka.getBrokersAsString()); - System.setProperty("spring.cloud.stream.kafka.streams.binder.zkNodes", - embeddedKafka.getZookeeperConnectionString()); - System.setProperty("server.port", "0"); System.setProperty("spring.jmx.enabled", "false"); Map consumerProps = KafkaTestUtils.consumerProps("kafka-streams-dlq-tests", "false", embeddedKafka); - // consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, - // Deserializer.class.getName()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); @@ -99,6 +94,9 @@ public abstract class DeserializtionErrorHandlerByBinderTests { @AfterClass public static void tearDown() { consumer.close(); + System.clearProperty("spring.cloud.stream.kafka.streams.binder.brokers"); + System.clearProperty("server.port"); + System.clearProperty("spring.jmx.enabled"); } @SpringBootTest(properties = { 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 2e0fbc4b7..60277c5d5 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 @@ -38,13 +38,11 @@ import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.HealthIndicator; import org.springframework.boot.actuate.health.Status; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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.Output; import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -207,20 +205,10 @@ 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.input.consumer.applicationId=" + "ApplicationHealthTest-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); } private ConfigurableApplicationContext multipleStream() { @@ -237,37 +225,16 @@ 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.input.consumer.applicationId=" + "ApplicationHealthTest-xyz", "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.applicationId=" + "ApplicationHealthTest2-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); } @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class KStreamApplication { @StreamListener("input") @@ -285,7 +252,6 @@ public class KafkaStreamsBinderHealthIndicatorTests { @EnableBinding({ KafkaStreamsProcessor.class, KafkaStreamsProcessorX.class }) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class AnotherKStreamApplication { @StreamListener("input") diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderMultipleInputTopicsTest.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsBinderMultipleInputTopicsTest.java index 9be47a52b..804ab91bc 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 @@ -37,12 +37,10 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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.kafka.streams.annotations.KafkaStreamsProcessor; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -113,18 +111,16 @@ public class KafkaStreamsBinderMultipleInputTopicsTest { "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=WordCountProcessorApplication-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { - receiveAndValidate(context); + receiveAndValidate(); } finally { context.close(); } } - private void receiveAndValidate(ConfigurableApplicationContext context) + private void receiveAndValidate() throws Exception { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( @@ -150,7 +146,6 @@ public class KafkaStreamsBinderMultipleInputTopicsTest { @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) static class WordCountProcessorApplication { @StreamListener 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 62c2c9a69..f372924c9 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 @@ -69,8 +69,6 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { public static void setUp() throws Exception { Map consumerProps = KafkaTestUtils.consumerProps("group-id", "false", embeddedKafka); - // consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, - // Deserializer.class.getName()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put("value.deserializer", LongDeserializer.class); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( @@ -100,19 +98,16 @@ public class KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests { "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=" + "KafkaStreamsBinderPojoInputAndPrimitiveTypeOutputTests-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { - receiveAndValidateFoo(context); + receiveAndValidateFoo(); } finally { context.close(); } } - private void receiveAndValidateFoo(ConfigurableApplicationContext context) - throws Exception { + private void receiveAndValidateFoo() { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); 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 f94f6d7ef..c267750ae 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 @@ -38,19 +38,15 @@ import org.apache.kafka.streams.state.ReadOnlyWindowStore; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; -import org.junit.Ignore; import org.junit.Test; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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.kafka.streams.annotations.KafkaStreamsProcessor; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.test.util.TestUtils; @@ -75,7 +71,7 @@ public class KafkaStreamsBinderWordCountIntegrationTests { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, - "counts"); + "counts", "counts-1"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule .getEmbeddedKafka(); @@ -83,14 +79,14 @@ public class KafkaStreamsBinderWordCountIntegrationTests { private static Consumer consumer; @BeforeClass - public static void setUp() throws Exception { + public static void setUp() { Map consumerProps = KafkaTestUtils.consumerProps("group", "false", embeddedKafka); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "counts"); + embeddedKafka.consumeFromEmbeddedTopics(consumer, "counts", "counts-1"); } @AfterClass @@ -99,7 +95,6 @@ public class KafkaStreamsBinderWordCountIntegrationTests { } @Test - @Ignore public void testKstreamWordCountWithApplicationIdSpecifiedAtDefaultConsumer() throws Exception { SpringApplication app = new SpringApplication( @@ -110,7 +105,7 @@ public class KafkaStreamsBinderWordCountIntegrationTests { "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.kafka.streams.default.consumer.application-id=basic-word-count", + "--spring.cloud.stream.kafka.streams.default.consumer.application-id=testKstreamWordCountWithApplicationIdSpecifiedAtDefaultConsumer", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", @@ -120,7 +115,7 @@ public class KafkaStreamsBinderWordCountIntegrationTests { "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.kafka.binder.brokers=" + embeddedKafka.getBrokersAsString())) { - receiveAndValidate(context); + receiveAndValidate("words", "counts"); } } @@ -133,9 +128,9 @@ public class KafkaStreamsBinderWordCountIntegrationTests { try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input.destination=words", - "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=basic-word-count", + "--spring.cloud.stream.bindings.input.destination=words-1", + "--spring.cloud.stream.bindings.output.destination=counts-1", + "--spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=testKstreamWordCountWithInputBindingLevelApplicationId", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", @@ -146,10 +141,8 @@ public class KafkaStreamsBinderWordCountIntegrationTests { "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", "--spring.cloud.stream.bindings.input.consumer.concurrency=2", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString())) { - receiveAndValidate(context); + + embeddedKafka.getBrokersAsString())) { + receiveAndValidate("words-1", "counts-1"); // Assertions on StreamBuilderFactoryBean StreamsBuilderFactoryBean streamsBuilderFactoryBean = context .getBean("&stream-builder-WordCountProcessorApplication-process", StreamsBuilderFactoryBean.class); @@ -176,17 +169,16 @@ public class KafkaStreamsBinderWordCountIntegrationTests { } } - private void receiveAndValidate(ConfigurableApplicationContext context) - throws Exception { + private void receiveAndValidate(String in, String out) { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); try { KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); + template.setDefaultTopic(in); template.sendDefault("foobar"); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, - "counts"); + out); assertThat(cr.value().contains("\"word\":\"foobar\",\"count\":1")).isTrue(); } finally { @@ -200,7 +192,7 @@ public class KafkaStreamsBinderWordCountIntegrationTests { senderProps); try { KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); + template.setDefaultTopic("words-1"); template.sendDefault(null); ConsumerRecords received = consumer .poll(Duration.ofMillis(5000)); @@ -216,12 +208,8 @@ public class KafkaStreamsBinderWordCountIntegrationTests { @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) static class WordCountProcessorApplication { - @Autowired - private TimeWindows timeWindows; - @StreamListener @SendTo("output") public KStream process( @@ -232,7 +220,7 @@ public class KafkaStreamsBinderWordCountIntegrationTests { value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .map((key, value) -> new KeyValue<>(value, value)) .groupByKey(Serialized.with(Serdes.String(), Serdes.String())) - .windowedBy(timeWindows).count(Materialized.as("foo-WordCounts")) + .windowedBy(TimeWindows.of(Duration.ofSeconds(5))).count(Materialized.as("foo-WordCounts")) .toStream() .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsInteractiveQueryIntegrationTests.java index 786fe974d..171eb7f1c 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 @@ -103,9 +103,7 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { "--spring.cloud.stream.kafka.streams.binder.configuration.application.server" + "=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { receiveAndValidateFoo(context); } @@ -153,7 +151,6 @@ public class KafkaStreamsInteractiveQueryIntegrationTests { @StreamListener("input") @SendTo("output") - @SuppressWarnings("deprecation") public KStream process(KStream input) { return input.filter((key, product) -> product.getId() == 123) diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java index 127f83fd0..fc6510268 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkaStreamsNativeEncodingDecodingTests.java @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.integration; +import java.time.Duration; import java.util.Arrays; import java.util.Map; @@ -34,16 +35,12 @@ import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; -import org.springframework.context.annotation.PropertySource; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -54,6 +51,7 @@ import org.springframework.messaging.handler.annotation.SendTo; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringRunner; +import org.springframework.util.StopWatch; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; @@ -70,7 +68,7 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, - "counts"); + "decode-counts", "decode-counts-1"); private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule .getEmbeddedKafka(); @@ -84,9 +82,6 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { public static void setUp() { System.setProperty("spring.cloud.stream.kafka.streams.binder.brokers", embeddedKafka.getBrokersAsString()); - System.setProperty("spring.cloud.stream.kafka.streams.binder.zkNodes", - embeddedKafka.getZookeeperConnectionString()); - System.setProperty("server.port", "0"); System.setProperty("spring.jmx.enabled", "false"); @@ -96,16 +91,20 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( consumerProps); consumer = cf.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "counts"); + embeddedKafka.consumeFromEmbeddedTopics(consumer, "decode-counts", "decode-counts-1"); } @AfterClass public static void tearDown() { consumer.close(); + System.clearProperty("spring.cloud.stream.kafka.streams.binder.brokers"); + System.clearProperty("server.port"); + System.clearProperty("spring.jmx.enabled"); } @SpringBootTest(properties = { - + "spring.cloud.stream.bindings.input.destination=decode-words-1", + "spring.cloud.stream.bindings.output.destination=decode-counts-1", "spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=NativeEncodingDecodingEnabledTests-abc" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) public static class NativeEncodingDecodingEnabledTests @@ -117,10 +116,10 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); + template.setDefaultTopic("decode-words-1"); template.sendDefault("foobar"); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, - "counts"); + "decode-counts-1"); assertThat(cr.value().equals("Count for foobar : 1")).isTrue(); verify(conversionDelegate, never()).serializeOnOutbound(any(KStream.class)); @@ -130,26 +129,31 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { } - // @checkstyle:off @SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = { + "spring.cloud.stream.bindings.input.destination=decode-words", + "spring.cloud.stream.bindings.output.destination=decode-counts", "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=false", "spring.cloud.stream.bindings.output.producer.useNativeEncoding=false", - "spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" - + "=NativeEncodingDecodingEnabledTests-xyz" }) - // @checkstyle:on + "spring.cloud.stream.kafka.streams.bindings.input3.consumer.applicationId" + + "=hello-NativeEncodingDecodingEnabledTests-xyz" }) public static class NativeEncodingDecodingDisabledTests extends KafkaStreamsNativeEncodingDecodingTests { @Test - public void test() throws Exception { + public void test() { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); + template.setDefaultTopic("decode-words"); template.sendDefault("foobar"); + StopWatch stopWatch = new StopWatch(); + stopWatch.start(); + System.out.println("Starting: "); ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, - "counts"); + "decode-counts"); + stopWatch.stop(); + System.out.println("Total time: " + stopWatch.getTotalTimeSeconds()); assertThat(cr.value().equals("Count for foobar : 1")).isTrue(); verify(conversionDelegate).serializeOnOutbound(any(KStream.class)); @@ -161,13 +165,8 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration - @PropertySource("classpath:/org/springframework/cloud/stream/binder/kstream/integTest-1.properties") - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class WordCountProcessorApplication { - @Autowired - private TimeWindows timeWindows; - @StreamListener("input") @SendTo("output") public KStream process(KStream input) { @@ -177,7 +176,7 @@ public abstract class KafkaStreamsNativeEncodingDecodingTests { value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .map((key, value) -> new KeyValue<>(value, value)) .groupByKey(Serialized.with(Serdes.String(), Serdes.String())) - .windowedBy(timeWindows).count(Materialized.as("foo-WordCounts-x")) + .windowedBy(TimeWindows.of(Duration.ofSeconds(5))).count(Materialized.as("foo-WordCounts-x")) .toStream().map((key, value) -> new KeyValue<>(null, "Count for " + key.key() + " : " + value)); } 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 d7be3249b..431a24637 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 @@ -70,9 +70,7 @@ public class KafkaStreamsStateStoreIntegrationTests { "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=KafkaStreamsStateStoreIntegrationTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { Thread.sleep(2000); receiveAndValidateFoo(context); @@ -97,18 +95,10 @@ public class KafkaStreamsStateStoreIntegrationTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", -// "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", -// "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.configuration.spring.json.value.default.type=" + -// "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsStateStoreIntegrationTests.Product", -// "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.valueSerde=org.springframework.kafka.support.serializer.JsonSerde", -// "--spring.cloud.stream.kafka.streams.bindings.input2.consumer.configuration.spring.json.value.default.type=" + -// "org.springframework.cloud.stream.binder.kafka.streams.integration.KafkaStreamsStateStoreIntegrationTests.Product", "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.applicationId" + "=KafkaStreamsStateStoreIntegrationTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { Thread.sleep(2000); // We are not particularly interested in querying the state store here, as that is verified by the other test 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 b6f117cca..ed113ce56 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkastreamsBinderPojoInputStringOutputIntegrationTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/KafkastreamsBinderPojoInputStringOutputIntegrationTests.java @@ -99,11 +99,9 @@ public class KafkastreamsBinderPojoInputStringOutputIntegrationTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId=ProductCountApplication-xyz", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { - receiveAndValidateFoo(context); + receiveAndValidateFoo(); // Assertions on StreamBuilderFactoryBean StreamsBuilderFactoryBean streamsBuilderFactoryBean = context .getBean("&stream-builder-ProductCountApplication-process", StreamsBuilderFactoryBean.class); @@ -117,8 +115,7 @@ public class KafkastreamsBinderPojoInputStringOutputIntegrationTests { } } - private void receiveAndValidateFoo(ConfigurableApplicationContext context) - throws Exception { + private void receiveAndValidateFoo() { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/MultiProcessorsWithSameNameTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/MultiProcessorsWithSameNameTests.java index 4c23eacaa..b27026ed1 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/MultiProcessorsWithSameNameTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/MultiProcessorsWithSameNameTests.java @@ -23,11 +23,9 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.test.EmbeddedKafkaBroker; @@ -60,9 +58,7 @@ public class MultiProcessorsWithSameNameTests { "--spring.cloud.stream.kafka.streams.bindings.input-1.consumer.application-id=basic-word-count", "--spring.cloud.stream.kafka.streams.bindings.input-2.consumer.application-id=basic-word-count-1", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString())) { + + embeddedKafka.getBrokersAsString())) { StreamsBuilderFactoryBean streamsBuilderFactoryBean1 = context .getBean("&stream-builder-Foo-process", StreamsBuilderFactoryBean.class); assertThat(streamsBuilderFactoryBean1).isNotNull(); @@ -74,7 +70,6 @@ public class MultiProcessorsWithSameNameTests { @EnableBinding(KafkaStreamsProcessorX.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) static class WordCountProcessorApplication { @Component diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java index d95378dc9..1b866e1df 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java @@ -32,7 +32,6 @@ import org.apache.kafka.streams.kstream.KStream; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.ClassRule; -import org.junit.Ignore; import org.junit.Test; import org.springframework.boot.SpringApplication; @@ -97,7 +96,6 @@ public class PerRecordAvroContentTypeTests { } @Test - @Ignore public void testPerRecordAvroConentTypeAndVerifySerialization() throws Exception { SpringApplication app = new SpringApplication(SensorCountAvroApplication.class); app.setWebApplicationType(WebApplicationType.NONE); 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 bf2bcbe20..ff0ac78c0 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 @@ -37,7 +37,6 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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; @@ -47,7 +46,6 @@ 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; @@ -98,9 +96,7 @@ public class StreamToGlobalKTableJoinIntegrationTests { "--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.getBrokersAsString()); 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 @@ -256,7 +252,6 @@ public class StreamToGlobalKTableJoinIntegrationTests { @EnableBinding(CustomGlobalKTableProcessor.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class OrderEnricherApplication { @StreamListener 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 54157cc0c..b86cd567f 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 @@ -44,7 +44,6 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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; @@ -54,7 +53,6 @@ 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.context.annotation.Bean; @@ -114,9 +112,7 @@ public class StreamToTableJoinIntegrationTests { + "=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.getBrokersAsString()); 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 @@ -265,9 +261,7 @@ public class StreamToTableJoinIntegrationTests { "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=10000", "--spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=helloxyz-foobar", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString())) { + + embeddedKafka.getBrokersAsString())) { Thread.sleep(1000L); // Input 2: Region per user (multiple records allowed per user). @@ -368,7 +362,6 @@ public class StreamToTableJoinIntegrationTests { @EnableBinding(KafkaStreamsProcessorX.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class CountClicksPerRegionApplication { @StreamListener @@ -385,7 +378,7 @@ public class StreamToTableJoinIntegrationTests { .map((user, regionWithClicks) -> new KeyValue<>( regionWithClicks.getRegion(), regionWithClicks.getClicks())) .groupByKey(Serialized.with(Serdes.String(), Serdes.Long())) - .reduce((firstClicks, secondClicks) -> firstClicks + secondClicks) + .reduce(Long::sum) .toStream(); } 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 7cba27483..baa38547d 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 @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.integration; +import java.time.Duration; import java.util.Arrays; import java.util.Date; import java.util.Map; @@ -33,16 +34,13 @@ import org.junit.BeforeClass; import org.junit.ClassRule; import org.junit.Test; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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.Output; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStreamsApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; @@ -108,9 +106,7 @@ public class WordCountMultipleBranchesIntegrationTests { "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + "=WordCountMultipleBranchesIntegrationTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString(), - "--spring.cloud.stream.kafka.streams.binder.zkNodes=" - + embeddedKafka.getZookeeperConnectionString()); + + embeddedKafka.getBrokersAsString()); try { receiveAndValidate(context); } @@ -145,12 +141,8 @@ public class WordCountMultipleBranchesIntegrationTests { @EnableBinding(KStreamProcessorX.class) @EnableAutoConfiguration - @EnableConfigurationProperties(KafkaStreamsApplicationSupportProperties.class) public static class WordCountProcessorApplication { - @Autowired - private TimeWindows timeWindows; - @StreamListener("input") @SendTo({ "output1", "output2", "output3" }) @SuppressWarnings("unchecked") @@ -163,7 +155,7 @@ public class WordCountMultipleBranchesIntegrationTests { return input .flatMapValues( value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .groupBy((key, value) -> value).windowedBy(timeWindows) + .groupBy((key, value) -> value).windowedBy(TimeWindows.of(Duration.ofSeconds(5))) .count(Materialized.as("WordCounts-multi")).toStream() .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/resources/org/springframework/cloud/stream/binder/kstream/integTest-1.properties b/spring-cloud-stream-binder-kafka-streams/src/test/resources/org/springframework/cloud/stream/binder/kstream/integTest-1.properties index a2342cc7f..4a0189c19 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/resources/org/springframework/cloud/stream/binder/kstream/integTest-1.properties +++ b/spring-cloud-stream-binder-kafka-streams/src/test/resources/org/springframework/cloud/stream/binder/kstream/integTest-1.properties @@ -1,5 +1,5 @@ -spring.cloud.stream.bindings.input.destination=words -spring.cloud.stream.bindings.output.destination=counts +spring.cloud.stream.bindings.input.destination=DeserializationErrorHandlerByKafkaTests-In +spring.cloud.stream.bindings.output.destination=DeserializationErrorHandlerByKafkaTests-Out spring.cloud.stream.bindings.output.contentType=application/json spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000 spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde