From 27eef43c706b387ccb0e307f0f8fae74e7d55022 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Wed, 3 Oct 2018 16:48:01 -0400 Subject: [PATCH] MINOR: fix bug where StoreBuilder added after connecting store with processor --- .../graph/StatefulProcessorNode.java | 8 ++-- .../internals/graph/StreamsGraphTest.java | 46 +++++++++++++++++++ 2 files changed, 50 insertions(+), 4 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StatefulProcessorNode.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StatefulProcessorNode.java index c2b445e83a469..a6a15ad871143 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StatefulProcessorNode.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/graph/StatefulProcessorNode.java @@ -60,13 +60,13 @@ public void writeToTopology(final InternalTopologyBuilder topologyBuilder) { topologyBuilder.addProcessor(processorName, processorSupplier, parentNodeNames()); - if (storeNames != null && storeNames.length > 0) { - topologyBuilder.connectProcessorAndStateStores(processorName, storeNames); - } - if (storeBuilder != null) { topologyBuilder.addStateStore(storeBuilder, processorName); } + + if (storeNames != null && storeNames.length > 0) { + topologyBuilder.connectProcessorAndStateStores(processorName, storeNames); + } } public static StatefulProcessorNodeBuilder statefulProcessorNodeBuilder() { diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/graph/StreamsGraphTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/graph/StreamsGraphTest.java index 75e9f5120b033..f4211bb6ca4a4 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/graph/StreamsGraphTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/graph/StreamsGraphTest.java @@ -17,17 +17,27 @@ package org.apache.kafka.streams.kstream.internals.graph; +import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; +import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.JoinWindows; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.TimeWindows; import org.apache.kafka.streams.kstream.ValueJoiner; +import org.apache.kafka.streams.kstream.internals.ConsumedInternal; +import org.apache.kafka.streams.processor.Processor; +import org.apache.kafka.streams.processor.ProcessorContext; +import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder; +import org.apache.kafka.streams.state.StoreBuilder; +import org.apache.kafka.streams.state.Stores; import org.junit.Test; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.Locale; import java.util.Properties; @@ -39,6 +49,8 @@ public class StreamsGraphTest { private final Pattern repartitionTopicPattern = Pattern.compile("Sink: .*-repartition"); + private final InternalTopologyBuilder builder = new InternalTopologyBuilder(); + private final Serde stringSerde = Serdes.String(); // Test builds topology in succesive manner but only graph node not yet processed written to topology @@ -66,6 +78,40 @@ public void shouldBeAbleToBuildTopologyIncrementally() { } + @Test + public void shouldBeAbleToAddStoresWithStatefulProcessorNodes() { + + final StreamSourceNode sourceNode = new StreamSourceNode<>("source", Collections.singletonList("topic"), new ConsumedInternal<>(Consumed.with(Serdes.String(), Serdes.String()))); + + final StoreBuilder storeBuilder = Stores.keyValueStoreBuilder( + Stores.inMemoryKeyValueStore("store"), + Serdes.String(), + Serdes.String()); + + final ProcessorParameters processorParameters = new ProcessorParameters<>(() -> new Processor() { + @Override + public void init(final ProcessorContext context) { + + } + + @Override + public void process(final String key, final String value) { + + } + + @Override + public void close() { + + } + }, "processor"); + final StatefulProcessorNode processorNode = new StatefulProcessorNode<>("node", processorParameters, new String[]{"store"}, storeBuilder, false); + sourceNode.addChild(processorNode); + builder.addSource(null, "source", null, stringSerde.deserializer(), stringSerde.deserializer(), "topic"); + + processorNode.writeToTopology(builder); + } + + @Test public void shouldNotOptimizeWithValueOrKeyChangingOperatorsAfterInitialKeyChange() {