diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index b451b3f9a..ddda09e96 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -17,9 +17,11 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.lang.reflect.Method; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Properties; @@ -78,9 +80,9 @@ import org.springframework.util.StringUtils; /** * Kafka Streams specific implementation for {@link StreamListenerSetupMethodOrchestrator} * that overrides the default mechanisms for invoking StreamListener adapters. - * + *

* The orchestration primarily focus on the following areas: - * + *

* 1. Allow multiple KStream output bindings (KStream branching) by allowing more than one * output values on {@link SendTo} 2. Allow multiple inbound bindings for multiple KStream * and or KTable/GlobalKTable types. 3. Each StreamListener method that it orchestrates @@ -110,6 +112,8 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator private final Map methodStreamsBuilderFactoryBeanMap = new HashMap<>(); + private final Map> registeredStoresPerMethod = new HashMap<>(); + private final CleanupConfig cleanupConfig; private ConfigurableApplicationContext applicationContext; @@ -160,9 +164,9 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } @Override - @SuppressWarnings({ "rawtypes", "unchecked" }) + @SuppressWarnings({"rawtypes", "unchecked"}) public void orchestrateStreamListenerSetupMethod(StreamListener streamListener, - Method method, Object bean) { + Method method, Object bean) { String[] methodAnnotatedOutboundNames = getOutboundBindingTargetNames(method); validateStreamListenerMethod(streamListener, method, methodAnnotatedOutboundNames); @@ -223,10 +227,10 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } @Override - @SuppressWarnings({ "unchecked" }) + @SuppressWarnings({"unchecked"}) public Object[] adaptAndRetrieveInboundArguments(Method method, String inboundName, - ApplicationContext applicationContext, - StreamListenerParameterAdapter... adapters) { + ApplicationContext applicationContext, + StreamListenerParameterAdapter... adapters) { Object[] arguments = new Object[method.getParameterTypes().length]; for (int parameterIndex = 0; parameterIndex < arguments.length; parameterIndex++) { MethodParameter methodParameter = MethodParameter.forExecutable(method, @@ -276,14 +280,14 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator Topology.AutoOffsetReset autoOffsetReset = null; if (startOffset != null) { switch (startOffset) { - case earliest: - autoOffsetReset = Topology.AutoOffsetReset.EARLIEST; - break; - case latest: - autoOffsetReset = Topology.AutoOffsetReset.LATEST; - break; - default: - break; + case earliest: + autoOffsetReset = Topology.AutoOffsetReset.EARLIEST; + break; + case latest: + autoOffsetReset = Topology.AutoOffsetReset.LATEST; + break; + default: + break; } } if (extendedConsumerProperties.isResetOffsets()) { @@ -367,30 +371,30 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } private GlobalKTable getGlobalKTable(StreamsBuilder streamsBuilder, - Serde keySerde, Serde valueSerde, String materializedAs, - String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { + Serde keySerde, Serde valueSerde, String materializedAs, + String bindingDestination, Topology.AutoOffsetReset autoOffsetReset) { return materializedAs != null ? materializedAsGlobalKTable(streamsBuilder, bindingDestination, - materializedAs, keySerde, valueSerde, autoOffsetReset) + materializedAs, keySerde, valueSerde, autoOffsetReset) : streamsBuilder.globalTable(bindingDestination, - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); + Consumed.with(keySerde, valueSerde) + .withOffsetResetPolicy(autoOffsetReset)); } private KTable getKTable(StreamsBuilder streamsBuilder, Serde keySerde, - Serde valueSerde, String materializedAs, String bindingDestination, - Topology.AutoOffsetReset autoOffsetReset) { + Serde valueSerde, String materializedAs, String bindingDestination, + Topology.AutoOffsetReset autoOffsetReset) { return materializedAs != null ? materializedAs(streamsBuilder, bindingDestination, materializedAs, - keySerde, valueSerde, autoOffsetReset) + keySerde, valueSerde, autoOffsetReset) : streamsBuilder.table(bindingDestination, - Consumed.with(keySerde, valueSerde) - .withOffsetResetPolicy(autoOffsetReset)); + Consumed.with(keySerde, valueSerde) + .withOffsetResetPolicy(autoOffsetReset)); } private KTable materializedAs(StreamsBuilder streamsBuilder, - String destination, String storeName, Serde k, Serde v, - Topology.AutoOffsetReset autoOffsetReset) { + String destination, String storeName, Serde k, Serde v, + Topology.AutoOffsetReset autoOffsetReset) { return streamsBuilder.table( this.bindingServiceProperties.getBindingDestination(destination), Consumed.with(k, v).withOffsetResetPolicy(autoOffsetReset), @@ -414,31 +418,32 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator private StoreBuilder buildStateStore(KafkaStreamsStateStoreProperties spec) { try { + Serde keySerde = this.keyValueSerdeResolver .getStateStoreKeySerde(spec.getKeySerdeString()); Serde valueSerde = this.keyValueSerdeResolver .getStateStoreValueSerde(spec.getValueSerdeString()); StoreBuilder builder; switch (spec.getType()) { - case KEYVALUE: - builder = Stores.keyValueStoreBuilder( - Stores.persistentKeyValueStore(spec.getName()), keySerde, - valueSerde); - break; - case WINDOW: - builder = Stores - .windowStoreBuilder( - Stores.persistentWindowStore(spec.getName(), - spec.getRetention(), 3, spec.getLength(), false), - keySerde, valueSerde); - break; - case SESSION: - builder = Stores.sessionStoreBuilder(Stores.persistentSessionStore( - spec.getName(), spec.getRetention()), keySerde, valueSerde); - break; - default: - throw new UnsupportedOperationException( - "state store type (" + spec.getType() + ") is not supported!"); + case KEYVALUE: + builder = Stores.keyValueStoreBuilder( + Stores.persistentKeyValueStore(spec.getName()), keySerde, + valueSerde); + break; + case WINDOW: + builder = Stores + .windowStoreBuilder( + Stores.persistentWindowStore(spec.getName(), + spec.getRetention(), 3, spec.getLength(), false), + keySerde, valueSerde); + break; + case SESSION: + builder = Stores.sessionStoreBuilder(Stores.persistentSessionStore( + spec.getName(), spec.getRetention()), keySerde, valueSerde); + break; + default: + throw new UnsupportedOperationException( + "state store type (" + spec.getType() + ") is not supported!"); } if (spec.isCacheEnabled()) { builder = builder.withCachingEnabled(); @@ -455,10 +460,10 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } private KStream getkStream(String inboundName, - KafkaStreamsStateStoreProperties storeSpec, - BindingProperties bindingProperties, StreamsBuilder streamsBuilder, - Serde keySerde, Serde valueSerde, - Topology.AutoOffsetReset autoOffsetReset) { + KafkaStreamsStateStoreProperties storeSpec, + BindingProperties bindingProperties, StreamsBuilder streamsBuilder, + Serde keySerde, Serde valueSerde, + Topology.AutoOffsetReset autoOffsetReset) { if (storeSpec != null) { StoreBuilder storeBuilder = buildStateStore(storeSpec); streamsBuilder.addStateStore(storeBuilder); @@ -499,7 +504,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } private void enableNativeDecodingForKTableAlways(Class parameterType, - BindingProperties bindingProperties) { + BindingProperties bindingProperties) { if (parameterType.isAssignableFrom(KTable.class) || parameterType.isAssignableFrom(GlobalKTable.class)) { if (bindingProperties.getConsumer() == null) { @@ -511,9 +516,9 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } } - @SuppressWarnings({ "unchecked" }) + @SuppressWarnings({"unchecked"}) private void buildStreamsBuilderAndRetrieveConfig(Method method, - ApplicationContext applicationContext, String inboundName) { + ApplicationContext applicationContext, String inboundName) { ConfigurableListableBeanFactory beanFactory = this.applicationContext .getBeanFactory(); @@ -557,7 +562,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator StreamsBuilderFactoryBean streamsBuilder = this.cleanupConfig == null ? new StreamsBuilderFactoryBean(kafkaStreamsConfiguration) : new StreamsBuilderFactoryBean(kafkaStreamsConfiguration, - this.cleanupConfig); + this.cleanupConfig); streamsBuilder.setAutoStartup(false); BeanDefinition streamsBuilderBeanDefinition = BeanDefinitionBuilder .genericBeanDefinition( @@ -578,7 +583,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator } private void validateStreamListenerMethod(StreamListener streamListener, - Method method, String[] methodAnnotatedOutboundNames) { + Method method, String[] methodAnnotatedOutboundNames) { String methodAnnotatedInboundName = streamListener.value(); if (methodAnnotatedOutboundNames != null) { for (String s : methodAnnotatedOutboundNames) { @@ -620,7 +625,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator @SuppressWarnings("unchecked") private boolean isDeclarativeInput(String targetBeanName, - MethodParameter methodParameter) { + MethodParameter methodParameter) { if (!methodParameter.getParameterType().isAssignableFrom(Object.class) && this.applicationContext.containsBean(targetBeanName)) { Class targetBeanClass = this.applicationContext.getType(targetBeanName); @@ -642,23 +647,27 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator return null; } - @SuppressWarnings({ "unchecked" }) - private static KafkaStreamsStateStoreProperties buildStateStoreSpec(Method method) { - KafkaStreamsStateStore spec = AnnotationUtils.findAnnotation(method, - KafkaStreamsStateStore.class); - if (spec != null) { - Assert.isTrue(!ObjectUtils.isEmpty(spec.name()), "name cannot be empty"); - Assert.isTrue(spec.name().length() >= 1, "name cannot be empty."); - KafkaStreamsStateStoreProperties props = new KafkaStreamsStateStoreProperties(); - props.setName(spec.name()); - props.setType(spec.type()); - props.setLength(spec.lengthMs()); - props.setKeySerdeString(spec.keySerde()); - props.setRetention(spec.retentionMs()); - props.setValueSerdeString(spec.valueSerde()); - props.setCacheEnabled(spec.cache()); - props.setLoggingDisabled(!spec.logging()); - return props; + @SuppressWarnings({"unchecked"}) + private KafkaStreamsStateStoreProperties buildStateStoreSpec(Method method) { + if (!this.registeredStoresPerMethod.containsKey(method)) { + KafkaStreamsStateStore spec = AnnotationUtils.findAnnotation(method, + KafkaStreamsStateStore.class); + if (spec != null) { + Assert.isTrue(!ObjectUtils.isEmpty(spec.name()), "name cannot be empty"); + Assert.isTrue(spec.name().length() >= 1, "name cannot be empty."); + this.registeredStoresPerMethod.put(method, new ArrayList<>()); + this.registeredStoresPerMethod.get(method).add(spec.name()); + KafkaStreamsStateStoreProperties props = new KafkaStreamsStateStoreProperties(); + props.setName(spec.name()); + props.setType(spec.type()); + props.setLength(spec.lengthMs()); + props.setKeySerdeString(spec.keySerde()); + props.setRetention(spec.retentionMs()); + props.setValueSerdeString(spec.valueSerde()); + props.setCacheEnabled(spec.cache()); + props.setLoggingDisabled(!spec.logging()); + return props; + } } return null; } 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 8d7cb662d..c304b4fa1 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 @@ -85,6 +85,40 @@ public class KafkaStreamsStateStoreIntegrationTests { } } + @Test + public void testSameStateStoreIsCreatedOnlyOnceWhenMultipleInputBindingsArePresent() throws Exception { + SpringApplication app = new SpringApplication(ProductCountApplicationWithMultipleInputBindings.class); + app.setWebApplicationType(WebApplicationType.NONE); + ConfigurableApplicationContext context = app.run("--server.port=0", + "--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.input.destination=foobar", + "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.bindings.input1.consumer.applicationId" + + "=KafkaStreamsStateStoreIntegrationTests-abc", + "--spring.cloud.stream.kafka.streams.binder.brokers=" + + embeddedKafka.getBrokersAsString(), + "--spring.cloud.stream.kafka.streams.binder.zkNodes=" + + embeddedKafka.getZookeeperConnectionString()); + try { + Thread.sleep(2000); + // We are not particularly interested in querying the state store here, as that is verified by the other test + // in this class. This test verifies that the same store is not attempted to be created by multiple input bindings. + // Normally, that will cause an exception to be thrown. However by not getting any exceptions, we are verifying + // that the binder is handling it appropriately. + //For more info, see this issue: https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/551 + } + catch (Exception e) { + throw e; + } + finally { + context.close(); + } + } + private void receiveAndValidateFoo(ConfigurableApplicationContext context) throws Exception { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); @@ -138,7 +172,44 @@ public class KafkaStreamsStateStoreIntegrationTests { } }, "mystate"); } + } + @EnableBinding(KafkaStreamsProcessorY.class) + @EnableAutoConfiguration + public static class ProductCountApplicationWithMultipleInputBindings { + + WindowStore state; + + boolean processed; + + @StreamListener + @KafkaStreamsStateStore(name = "mystate", type = KafkaStreamsStateStoreProperties.StoreType.WINDOW, lengthMs = 300000) + @SuppressWarnings({ "deprecation", "unchecked" }) + public void process(@Input("input1")KStream input, @Input("input2")KStream input2) { + + input.process(() -> new Processor() { + + @Override + public void init(ProcessorContext processorContext) { + state = (WindowStore) processorContext.getStateStore("mystate"); + } + + @Override + public void process(Object s, Product product) { + processed = true; + } + + @Override + public void close() { + if (state != null) { + state.close(); + } + } + }, "mystate"); + + //simple use of input2, we are not using input2 for anything other than triggering some test behavior. + input2.foreach((key, value) -> { }); + } } public static class Product { @@ -159,7 +230,14 @@ public class KafkaStreamsStateStoreIntegrationTests { @Input("input") KStream input(); - } + interface KafkaStreamsProcessorY { + + @Input("input1") + KStream input1(); + + @Input("input2") + KStream input2(); + } }