From a02308a5a356e39666700ad2ddc9f49b5112733d Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 1 Oct 2019 13:58:14 -0400 Subject: [PATCH] Allow binding names to be reused in Kafka Streams. Allow same binding names to be reused from multiple StreamListener methods in Kafka Streams binder. Resolves #760 --- .../streams/GlobalKTableBoundElementFactory.java | 5 +++-- .../kafka/streams/KStreamBoundElementFactory.java | 5 +++-- .../kafka/streams/KTableBoundElementFactory.java | 5 +++-- ...ltiProcessorsWithSameNameAndBindingTests.java} | 15 +++++---------- 4 files changed, 14 insertions(+), 16 deletions(-) rename spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/{MultiProcessorsWithSameNameTests.java => MultiProcessorsWithSameNameAndBindingTests.java} (86%) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java index 6b5e806fd..b99952dc4 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBoundElementFactory.java @@ -90,8 +90,9 @@ public class GlobalKTableBoundElementFactory public void wrap(GlobalKTable delegate) { Assert.notNull(delegate, "delegate cannot be null"); - Assert.isNull(this.delegate, "delegate already set to " + this.delegate); - this.delegate = delegate; + if (this.delegate == null) { + this.delegate = delegate; + } } @Override diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java index 439d935ef..526e55dc7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBoundElementFactory.java @@ -109,8 +109,9 @@ class KStreamBoundElementFactory extends AbstractBindingTargetFactory { public void wrap(KStream delegate) { Assert.notNull(delegate, "delegate cannot be null"); - Assert.isNull(this.delegate, "delegate already set to " + this.delegate); - this.delegate = delegate; + if (this.delegate == null) { + this.delegate = delegate; + } } @Override diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java index f9f3fef02..533801938 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBoundElementFactory.java @@ -86,8 +86,9 @@ class KTableBoundElementFactory extends AbstractBindingTargetFactory { public void wrap(KTable delegate) { Assert.notNull(delegate, "delegate cannot be null"); - Assert.isNull(this.delegate, "delegate already set to " + this.delegate); - this.delegate = delegate; + if (this.delegate == null) { + this.delegate = delegate; + } } @Override 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/MultiProcessorsWithSameNameAndBindingTests.java similarity index 86% rename from spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/MultiProcessorsWithSameNameTests.java rename to spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/MultiProcessorsWithSameNameAndBindingTests.java index b27026ed1..699cee4a7 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/MultiProcessorsWithSameNameAndBindingTests.java @@ -34,7 +34,7 @@ import org.springframework.stereotype.Component; import static org.assertj.core.api.Assertions.assertThat; -public class MultiProcessorsWithSameNameTests { +public class MultiProcessorsWithSameNameAndBindingTests { @ClassRule public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, @@ -44,19 +44,17 @@ public class MultiProcessorsWithSameNameTests { .getEmbeddedKafka(); @Test - public void testBinderStartsSuccessfullyWhenTwoProcessorsWithSameNamesArePresent() { + public void testBinderStartsSuccessfullyWhenTwoProcessorsWithSameNamesAndBindingsPresent() { SpringApplication app = new SpringApplication( - MultiProcessorsWithSameNameTests.WordCountProcessorApplication.class); + MultiProcessorsWithSameNameAndBindingTests.WordCountProcessorApplication.class); app.setWebApplicationType(WebApplicationType.NONE); try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", - "--spring.cloud.stream.bindings.input-2.destination=words", + "--spring.cloud.stream.bindings.input-1.destination=words", "--spring.cloud.stream.bindings.output.destination=counts", "--spring.cloud.stream.bindings.output.contentType=application/json", - "--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())) { StreamsBuilderFactoryBean streamsBuilderFactoryBean1 = context @@ -83,7 +81,7 @@ public class MultiProcessorsWithSameNameTests { @Component static class Bar { @StreamListener - public void process(@Input("input-2") KStream input) { + public void process(@Input("input-1") KStream input) { } } } @@ -93,8 +91,5 @@ public class MultiProcessorsWithSameNameTests { @Input("input-1") KStream input1(); - @Input("input-2") - KStream input2(); - } }