From b4a2950acd88df0f96178ca7ea2ba31e9cff43da Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 28 Feb 2019 16:31:31 -0500 Subject: [PATCH] KafkaStreamsStateStore with multiple input bindings When KafkaStreamsStateStore annotation is used on a method with multiple input bindings, it throws an exception. The reason is that each successive input binding after the first one is trying to recreate the store that is already created. Fixing this issue. Resolves #551 --- ...StreamListenerSetupMethodOrchestrator.java | 155 +++++++++--------- ...afkaStreamsStateStoreIntegrationTests.java | 80 ++++++++- 2 files changed, 161 insertions(+), 74 deletions(-) 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(); + } }