From d2f735cf7ee98ea255926fbfc32e91d9b39ee165 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Thu, 17 May 2018 18:42:56 -0500 Subject: [PATCH 1/2] return to double-counting for count topo names --- .../kstream/internals/KGroupedStreamImpl.java | 22 +- .../kstream/internals/MaterializedPeek.java | 33 ++ .../internals/SessionWindowedKStreamImpl.java | 26 +- .../internals/TimeWindowedKStreamImpl.java | 26 +- .../apache/kafka/streams/TopologyTest.java | 351 ++++++++++++++++-- 5 files changed, 400 insertions(+), 58 deletions(-) create mode 100644 streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java index 3a9f9197ac7c8..fd68858b468a5 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java @@ -117,14 +117,26 @@ public KTable aggregate(final Initializer initializer, @Override public KTable count() { - return count(Materialized.>with(keySerde, Serdes.Long())); + return doCount(Materialized.with(keySerde, Serdes.Long())); } @Override public KTable count(final Materialized> materialized) { Objects.requireNonNull(materialized, "materialized can't be null"); + + // TODO: remove this when we do a topology-incompatible release + // we used to burn a topology name here, so we have to keep doing it for compatibility + final String givenStoreName = new MaterializedPeek<>(materialized).givenStoreName(); + if (givenStoreName == null) { + builder.newStoreName(AGGREGATE_NAME); + } + + return doCount(materialized); + } + + private KTable doCount(final Materialized> materialized) { final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -133,9 +145,9 @@ public KTable count(final Materialized(materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator), - AGGREGATE_NAME, - materializedInternal); + new KStreamAggregate<>(materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator), + AGGREGATE_NAME, + materializedInternal); } @Override diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java new file mode 100644 index 0000000000000..00784f96888a7 --- /dev/null +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java @@ -0,0 +1,33 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.kstream.internals; + +import org.apache.kafka.streams.kstream.Materialized; +import org.apache.kafka.streams.processor.StateStore; + +class MaterializedPeek extends Materialized { + MaterializedPeek(final Materialized materialized) { + super(materialized); + } + + String givenStoreName() { + if (storeSupplier != null) { + return storeSupplier.name(); + } + return storeName; + } +} diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java index c29c6566c2f4b..abd8d5b20e4a0 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java @@ -68,15 +68,27 @@ public Long apply(final K aggKey, final Long aggOne, final Long aggTwo) { @Override public KTable, Long> count() { - return count(Materialized.>with(keySerde, Serdes.Long())); + return doCount(Materialized.with(keySerde, Serdes.Long())); } - @SuppressWarnings("unchecked") @Override public KTable, Long> count(final Materialized> materialized) { Objects.requireNonNull(materialized, "materialized can't be null"); + + // TODO: remove this when we do a topology-incompatible release + // we used to burn a topology name here, so we have to keep doing it for compatibility + final String givenStoreName = new MaterializedPeek<>(materialized).givenStoreName(); + if (givenStoreName == null) { + builder.newStoreName(AGGREGATE_NAME); + } + + return doCount(materialized); + } + + @SuppressWarnings("unchecked") + private KTable, Long> doCount(final Materialized> materialized) { final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -85,10 +97,10 @@ public KTable, Long> count(final Materialized, Long>) aggregateBuilder.build( - new KStreamSessionWindowAggregate<>(windows, materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator, countMerger), - AGGREGATE_NAME, - materialize(materializedInternal), - materializedInternal.isQueryable()); + new KStreamSessionWindowAggregate<>(windows, materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator, countMerger), + AGGREGATE_NAME, + materialize(materializedInternal), + materializedInternal.isQueryable()); } @Override diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java index d1e5a1758476a..5b268aed805ba 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java @@ -24,9 +24,9 @@ import org.apache.kafka.streams.kstream.KTable; import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.kstream.Reducer; +import org.apache.kafka.streams.kstream.TimeWindowedKStream; import org.apache.kafka.streams.kstream.Window; import org.apache.kafka.streams.kstream.Windowed; -import org.apache.kafka.streams.kstream.TimeWindowedKStream; import org.apache.kafka.streams.kstream.Windows; import org.apache.kafka.streams.state.StoreBuilder; import org.apache.kafka.streams.state.Stores; @@ -63,15 +63,27 @@ public class TimeWindowedKStreamImpl extends AbstractStr @Override public KTable, Long> count() { - return count(Materialized.>with(keySerde, Serdes.Long())); + return doCount(Materialized.with(keySerde, Serdes.Long())); } - @SuppressWarnings("unchecked") @Override public KTable, Long> count(final Materialized> materialized) { Objects.requireNonNull(materialized, "materialized can't be null"); + + // TODO: remove this when we do a topology-incompatible release + // we used to burn a topology name here, so we have to keep doing it for compatibility + final String givenStoreName = new MaterializedPeek<>(materialized).givenStoreName(); + if (givenStoreName == null) { + builder.newStoreName(AGGREGATE_NAME); + } + + return doCount(materialized); + } + + @SuppressWarnings("unchecked") + private KTable, Long> doCount(final Materialized> materialized) { final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -80,9 +92,9 @@ public KTable, Long> count(final Materialized, Long>) aggregateBuilder.build(new KStreamWindowAggregate<>(windows, materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator), - AGGREGATE_NAME, - materialize(materializedInternal), - materializedInternal.isQueryable()); + AGGREGATE_NAME, + materialize(materializedInternal), + materializedInternal.isQueryable()); } diff --git a/streams/src/test/java/org/apache/kafka/streams/TopologyTest.java b/streams/src/test/java/org/apache/kafka/streams/TopologyTest.java index a845be3569f47..f8f9c7d68a538 100644 --- a/streams/src/test/java/org/apache/kafka/streams/TopologyTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/TopologyTest.java @@ -16,8 +16,12 @@ */ package org.apache.kafka.streams; +import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.errors.TopologyException; +import org.apache.kafka.streams.kstream.Materialized; +import org.apache.kafka.streams.kstream.SessionWindows; +import org.apache.kafka.streams.kstream.TimeWindows; import org.apache.kafka.streams.processor.Processor; import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.processor.ProcessorSupplier; @@ -40,6 +44,7 @@ import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; +import static org.junit.Assert.assertEquals; import static org.junit.Assert.fail; public class TopologyTest { @@ -76,12 +81,7 @@ public void shouldNotAllowZeroTopicsWhenAddingSource() { @Test(expected = NullPointerException.class) public void shouldNotAllowNullNameWhenAddingProcessor() { - topology.addProcessor(null, new ProcessorSupplier() { - @Override - public Processor get() { - return new MockProcessorSupplier().get(); - } - }); + topology.addProcessor(null, () -> new MockProcessorSupplier().get()); } @Test(expected = NullPointerException.class) @@ -130,7 +130,7 @@ public void shouldNotAllowToAddSourcesWithSameName() { try { topology.addSource("source", "topic-2"); fail("Should throw TopologyException for duplicate source name"); - } catch (TopologyException expected) { } + } catch (final TopologyException expected) { } } @Test @@ -139,7 +139,7 @@ public void shouldNotAllowToAddTopicTwice() { try { topology.addSource("source-2", "topic-1"); fail("Should throw TopologyException for already used topic"); - } catch (TopologyException expected) { } + } catch (final TopologyException expected) { } } @Test @@ -167,7 +167,7 @@ public void shouldNotAllowToAddProcessorWithSameName() { try { topology.addProcessor("processor", new MockProcessorSupplier(), "source"); fail("Should throw TopologyException for duplicate processor name"); - } catch (TopologyException expected) { } + } catch (final TopologyException expected) { } } @Test(expected = TopologyException.class) @@ -187,7 +187,7 @@ public void shouldNotAllowToAddSinkWithSameName() { try { topology.addSink("sink", "topic-3", "source"); fail("Should throw TopologyException for duplicate sink name"); - } catch (TopologyException expected) { } + } catch (final TopologyException expected) { } } @Test(expected = TopologyException.class) @@ -257,7 +257,7 @@ public void shouldNotAllowToAddStoreWithSameName() { } @Test - public void shouldThrowOnUnassignedStateStoreAccess() throws Exception { + public void shouldThrowOnUnassignedStateStoreAccess() { final String sourceNodeName = "source"; final String goodNodeName = "goodGuy"; final String badNodeName = "badGuy"; @@ -283,7 +283,7 @@ public void shouldThrowOnUnassignedStateStoreAccess() throws Exception { } catch (final StreamsException e) { final String error = e.toString(); final String expectedMessage = "org.apache.kafka.streams.errors.StreamsException: failed to initialize processor " + badNodeName; - + assertThat(error, equalTo(expectedMessage)); } } @@ -295,12 +295,12 @@ private static class LocalMockProcessorSupplier implements ProcessorSupplier { public Processor get() { return new Processor() { @Override - public void init(ProcessorContext context) { + public void init(final ProcessorContext context) { context.getStateStore(STORE_NAME); } @Override - public void process(Object key, Object value) { } + public void process(final Object key, final Object value) { } @Override public void close() { } @@ -324,7 +324,7 @@ public void shouldNotAllowToAddGlobalStoreWithSourceNameEqualsProcessorName() { @Test public void shouldDescribeEmptyTopology() { - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -333,9 +333,9 @@ public void singleSourceShouldHaveSingleSubtopology() { expectedDescription.addSubtopology( new InternalTopologyBuilder.Subtopology(0, - Collections.singleton(expectedSourceNode))); + Collections.singleton(expectedSourceNode))); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -344,9 +344,9 @@ public void singleSourceWithListOfTopicsShouldHaveSingleSubtopology() { expectedDescription.addSubtopology( new InternalTopologyBuilder.Subtopology(0, - Collections.singleton(expectedSourceNode))); + Collections.singleton(expectedSourceNode))); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -355,9 +355,9 @@ public void singleSourcePatternShouldHaveSingleSubtopology() { expectedDescription.addSubtopology( new InternalTopologyBuilder.Subtopology(0, - Collections.singleton(expectedSourceNode))); + Collections.singleton(expectedSourceNode))); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -365,19 +365,19 @@ public void multipleSourcesShouldHaveDistinctSubtopologies() { final TopologyDescription.Source expectedSourceNode1 = addSource("source1", "topic1"); expectedDescription.addSubtopology( new InternalTopologyBuilder.Subtopology(0, - Collections.singleton(expectedSourceNode1))); + Collections.singleton(expectedSourceNode1))); final TopologyDescription.Source expectedSourceNode2 = addSource("source2", "topic2"); expectedDescription.addSubtopology( new InternalTopologyBuilder.Subtopology(1, - Collections.singleton(expectedSourceNode2))); + Collections.singleton(expectedSourceNode2))); final TopologyDescription.Source expectedSourceNode3 = addSource("source3", "topic3"); expectedDescription.addSubtopology( new InternalTopologyBuilder.Subtopology(2, - Collections.singleton(expectedSourceNode3))); + Collections.singleton(expectedSourceNode3))); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -390,7 +390,7 @@ public void sourceAndProcessorShouldHaveSingleSubtopology() { allNodes.add(expectedProcessorNode); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -405,7 +405,7 @@ public void sourceAndProcessorWithStateShouldHaveSingleSubtopology() { allNodes.add(expectedProcessorNode); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @@ -421,7 +421,7 @@ public void sourceAndProcessorWithMultipleStatesShouldHaveSingleSubtopology() { allNodes.add(expectedProcessorNode); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -436,7 +436,7 @@ public void sourceWithMultipleProcessorsShouldHaveSingleSubtopology() { allNodes.add(expectedProcessorNode2); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -451,7 +451,7 @@ public void processorWithMultipleSourcesShouldHaveSingleSubtopology() { allNodes.add(expectedProcessorNode); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -480,7 +480,7 @@ public void multipleSourcesWithProcessorsShouldHaveDistinctSubtopologies() { allNodes3.add(expectedProcessorNode3); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(2, allNodes3)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -509,7 +509,7 @@ public void multipleSourcesWithSinksShouldHaveDistinctSubtopologies() { allNodes3.add(expectedSinkNode3); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(2, allNodes3)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -540,7 +540,7 @@ public void processorsWithSameSinkShouldHaveSameSubtopology() { allNodes.add(expectedSinkNode); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test @@ -570,30 +570,303 @@ public void processorsWithSharedStateShouldHaveSameSubtopology() { allNodes.add(expectedProcessorNode3); expectedDescription.addSubtopology(new InternalTopologyBuilder.Subtopology(0, allNodes)); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test public void shouldDescribeGlobalStoreTopology() { addGlobalStoreToTopologyAndExpectedDescription("globalStore", "source", "globalTopic", "processor", 0); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); } @Test public void shouldDescribeMultipleGlobalStoreTopology() { addGlobalStoreToTopologyAndExpectedDescription("globalStore1", "source1", "globalTopic1", "processor1", 0); addGlobalStoreToTopologyAndExpectedDescription("globalStore2", "source2", "globalTopic2", "processor2", 1); - assertThat(topology.describe(), equalTo((TopologyDescription) expectedDescription)); + assertThat(topology.describe(), equalTo(expectedDescription)); + } + + @Test + public void kGroupedStreamZeroArgCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .count(); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000002\n" + + " Processor: KSTREAM-AGGREGATE-0000000002 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000001])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void kGroupedStreamNamedMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .count(Materialized.as("count-store")); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000001\n" + + " Processor: KSTREAM-AGGREGATE-0000000001 (stores: [count-store])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void kGroupedStreamAnonymousMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .count(Materialized.with(null, Serdes.Long())); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000003\n" + + " Processor: KSTREAM-AGGREGATE-0000000003 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000002])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void timeWindowZeroArgCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .windowedBy(TimeWindows.of(1)) + .count(); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000002\n" + + " Processor: KSTREAM-AGGREGATE-0000000002 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000001])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void timeWindowNamedMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .windowedBy(TimeWindows.of(1)) + .count(Materialized.as("count-store")); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000001\n" + + " Processor: KSTREAM-AGGREGATE-0000000001 (stores: [count-store])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void timeWindowAnonymousMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .windowedBy(TimeWindows.of(1)) + .count(Materialized.with(null, Serdes.Long())); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000003\n" + + " Processor: KSTREAM-AGGREGATE-0000000003 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000002])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void sessionWindowZeroArgCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .windowedBy(SessionWindows.with(1)) + .count(); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000002\n" + + " Processor: KSTREAM-AGGREGATE-0000000002 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000001])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void sessionWindowNamedMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .windowedBy(SessionWindows.with(1)) + .count(Materialized.as("count-store")); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000001\n" + + " Processor: KSTREAM-AGGREGATE-0000000001 (stores: [count-store])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void sessionWindowAnonymousMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("input-topic") + .groupByKey() + .windowedBy(SessionWindows.with(1)) + .count(Materialized.with(null, Serdes.Long())); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000000 (topics: [input-topic])\n" + + " --> KSTREAM-AGGREGATE-0000000003\n" + + " Processor: KSTREAM-AGGREGATE-0000000003 (stores: [KSTREAM-AGGREGATE-STATE-STORE-0000000002])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000000\n\n", + describe.toString() + ); + } + + @Test + public void tableZeroArgCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.table("input-topic") + .groupBy((key, value) -> null) + .count(); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000001 (topics: [input-topic])\n" + + " --> KTABLE-SOURCE-0000000002\n" + + " Processor: KTABLE-SOURCE-0000000002 (stores: [input-topic-STATE-STORE-0000000000])\n" + + " --> KTABLE-SELECT-0000000003\n" + + " <-- KSTREAM-SOURCE-0000000001\n" + + " Processor: KTABLE-SELECT-0000000003 (stores: [])\n" + + " --> KSTREAM-SINK-0000000005\n" + + " <-- KTABLE-SOURCE-0000000002\n" + + " Sink: KSTREAM-SINK-0000000005 (topic: KTABLE-AGGREGATE-STATE-STORE-0000000004-repartition)\n" + + " <-- KTABLE-SELECT-0000000003\n" + + "\n" + + " Sub-topology: 1\n" + + " Source: KSTREAM-SOURCE-0000000006 (topics: [KTABLE-AGGREGATE-STATE-STORE-0000000004-repartition])\n" + + " --> KTABLE-AGGREGATE-0000000007\n" + + " Processor: KTABLE-AGGREGATE-0000000007 (stores: [KTABLE-AGGREGATE-STATE-STORE-0000000004])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000006\n" + + "\n", + describe.toString() + ); + } + + @Test + public void tableNamedMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.table("input-topic") + .groupBy((key, value) -> null) + .count(Materialized.as("count-store")); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000001 (topics: [input-topic])\n" + + " --> KTABLE-SOURCE-0000000002\n" + + " Processor: KTABLE-SOURCE-0000000002 (stores: [input-topic-STATE-STORE-0000000000])\n" + + " --> KTABLE-SELECT-0000000003\n" + + " <-- KSTREAM-SOURCE-0000000001\n" + + " Processor: KTABLE-SELECT-0000000003 (stores: [])\n" + + " --> KSTREAM-SINK-0000000004\n" + + " <-- KTABLE-SOURCE-0000000002\n" + + " Sink: KSTREAM-SINK-0000000004 (topic: count-store-repartition)\n" + + " <-- KTABLE-SELECT-0000000003\n" + + "\n" + + " Sub-topology: 1\n" + + " Source: KSTREAM-SOURCE-0000000005 (topics: [count-store-repartition])\n" + + " --> KTABLE-AGGREGATE-0000000006\n" + + " Processor: KTABLE-AGGREGATE-0000000006 (stores: [count-store])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000005\n" + + "\n", + describe.toString() + ); + } + + @Test + public void tableAnonymousMaterializedCountShouldPreserveTopologyStructure() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.table("input-topic") + .groupBy((key, value) -> null) + .count(Materialized.with(null, Serdes.Long())); + final TopologyDescription describe = builder.build().describe(); + assertEquals( + "Topologies:\n" + + " Sub-topology: 0\n" + + " Source: KSTREAM-SOURCE-0000000001 (topics: [input-topic])\n" + + " --> KTABLE-SOURCE-0000000002\n" + + " Processor: KTABLE-SOURCE-0000000002 (stores: [input-topic-STATE-STORE-0000000000])\n" + + " --> KTABLE-SELECT-0000000003\n" + + " <-- KSTREAM-SOURCE-0000000001\n" + + " Processor: KTABLE-SELECT-0000000003 (stores: [])\n" + + " --> KSTREAM-SINK-0000000005\n" + + " <-- KTABLE-SOURCE-0000000002\n" + + " Sink: KSTREAM-SINK-0000000005 (topic: KTABLE-AGGREGATE-STATE-STORE-0000000004-repartition)\n" + + " <-- KTABLE-SELECT-0000000003\n" + + "\n" + + " Sub-topology: 1\n" + + " Source: KSTREAM-SOURCE-0000000006 (topics: [KTABLE-AGGREGATE-STATE-STORE-0000000004-repartition])\n" + + " --> KTABLE-AGGREGATE-0000000007\n" + + " Processor: KTABLE-AGGREGATE-0000000007 (stores: [KTABLE-AGGREGATE-STATE-STORE-0000000004])\n" + + " --> none\n" + + " <-- KSTREAM-SOURCE-0000000006\n" + + "\n", + describe.toString() + ); } private TopologyDescription.Source addSource(final String sourceName, final String... sourceTopic) { topology.addSource(null, sourceName, null, null, null, sourceTopic); - String allSourceTopics = sourceTopic[0]; + final StringBuilder allSourceTopics = new StringBuilder(sourceTopic[0]); for (int i = 1; i < sourceTopic.length; ++i) { - allSourceTopics += ", " + sourceTopic[i]; + allSourceTopics.append(", ").append(sourceTopic[i]); } - return new InternalTopologyBuilder.Source(sourceName, allSourceTopics); + return new InternalTopologyBuilder.Source(sourceName, allSourceTopics.toString()); } private TopologyDescription.Source addSource(final String sourceName, From 2851a05e213d87be1f908d8fad1949dcd92659f5 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Fri, 1 Jun 2018 16:14:11 -0500 Subject: [PATCH 2/2] move name generation from MI constructor to method * remove MaterializePeek now that MI constructor doesn't mutate storeName --- .../apache/kafka/streams/StreamsBuilder.java | 44 ++++--- .../kstream/internals/KGroupedStreamImpl.java | 19 +-- .../kstream/internals/KGroupedTableImpl.java | 15 ++- .../streams/kstream/internals/KTableImpl.java | 44 ++++--- .../internals/MaterializedInternal.java | 19 ++- .../kstream/internals/MaterializedPeek.java | 33 ------ .../internals/SessionWindowedKStreamImpl.java | 16 +-- .../internals/TimeWindowedKStreamImpl.java | 17 +-- .../internals/InternalStreamsBuilderTest.java | 108 ++++++++---------- .../internals/MaterializedInternalTest.java | 15 ++- .../internals/GlobalStreamThreadTest.java | 10 +- .../KeyValueStoreMaterializerTest.java | 35 +++--- .../processor/internals/StreamThreadTest.java | 4 +- 13 files changed, 186 insertions(+), 193 deletions(-) delete mode 100644 streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsBuilder.java b/streams/src/main/java/org/apache/kafka/streams/StreamsBuilder.java index ead8a7663abf2..517104da323d0 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsBuilder.java @@ -224,9 +224,9 @@ public synchronized KTable table(final String topic, Objects.requireNonNull(materialized, "materialized can't be null"); final ConsumedInternal consumedInternal = new ConsumedInternal<>(consumed); materialized.withKeySerde(consumedInternal.keySerde()).withValueSerde(consumedInternal.valueSerde()); - return internalStreamsBuilder.table(topic, - consumedInternal, - new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-")); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-"); + return internalStreamsBuilder.table(topic, consumedInternal, materializedInternal); } /** @@ -273,12 +273,10 @@ public synchronized KTable table(final String topic, Objects.requireNonNull(topic, "topic can't be null"); Objects.requireNonNull(consumed, "consumed can't be null"); final ConsumedInternal consumedInternal = new ConsumedInternal<>(consumed); - return internalStreamsBuilder.table(topic, - consumedInternal, - new MaterializedInternal<>( - Materialized.>with(consumedInternal.keySerde(), consumedInternal.valueSerde()), - internalStreamsBuilder, - topic + "-")); + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.with(consumedInternal.keySerde(), consumedInternal.valueSerde())); + materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-"); + return internalStreamsBuilder.table(topic, consumedInternal, materializedInternal); } /** @@ -302,8 +300,9 @@ public synchronized KTable table(final String topic, final Materialized> materialized) { Objects.requireNonNull(topic, "topic can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-"); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-"); + return internalStreamsBuilder.table(topic, new ConsumedInternal<>(Consumed.with(materializedInternal.keySerde(), materializedInternal.valueSerde())), @@ -331,14 +330,11 @@ public synchronized GlobalKTable globalTable(final String topic, Objects.requireNonNull(topic, "topic can't be null"); Objects.requireNonNull(consumed, "consumed can't be null"); final ConsumedInternal consumedInternal = new ConsumedInternal<>(consumed); - final MaterializedInternal> materialized = - new MaterializedInternal<>( - Materialized.>with(consumedInternal.keySerde(), consumedInternal.valueSerde()), - internalStreamsBuilder, - topic + "-"); - + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.>with(consumedInternal.keySerde(), consumedInternal.valueSerde())); + materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-"); - return internalStreamsBuilder.globalTable(topic, consumedInternal, materialized); + return internalStreamsBuilder.globalTable(topic, consumedInternal, materializedInternal); } /** @@ -402,9 +398,10 @@ public synchronized GlobalKTable globalTable(final String topic, final ConsumedInternal consumedInternal = new ConsumedInternal<>(consumed); // always use the serdes from consumed materialized.withKeySerde(consumedInternal.keySerde()).withValueSerde(consumedInternal.valueSerde()); - return internalStreamsBuilder.globalTable(topic, - consumedInternal, - new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-")); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-"); + + return internalStreamsBuilder.globalTable(topic, consumedInternal, materializedInternal); } /** @@ -436,8 +433,9 @@ public synchronized GlobalKTable globalTable(final String topic, final Materialized> materialized) { Objects.requireNonNull(topic, "topic can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal = - new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-"); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-"); + return internalStreamsBuilder.globalTable(topic, new ConsumedInternal<>(Consumed.with(materializedInternal.keySerde(), materializedInternal.valueSerde())), diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java index fd68858b468a5..74f930ee357ed 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImpl.java @@ -74,8 +74,10 @@ public KTable reduce(final Reducer reducer, final Materialized> materialized) { Objects.requireNonNull(reducer, "reducer can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, REDUCE_NAME); + + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, REDUCE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -97,8 +99,9 @@ public KTable aggregate(final Initializer initializer, Objects.requireNonNull(aggregator, "aggregator can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -126,8 +129,7 @@ public KTable count(final Materialized(materialized).givenStoreName(); - if (givenStoreName == null) { + if (new MaterializedInternal<>(materialized).storeName() == null) { builder.newStoreName(AGGREGATE_NAME); } @@ -135,8 +137,9 @@ public KTable count(final Materialized doCount(final Materialized> materialized) { - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImpl.java index db119f30fd57f..49f258bf5035b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImpl.java @@ -128,8 +128,9 @@ public KTable reduce(final Reducer adder, Objects.requireNonNull(adder, "adder can't be null"); Objects.requireNonNull(subtractor, "subtractor can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -150,8 +151,9 @@ public KTable reduce(final Reducer adder, @Override public KTable count(final Materialized> materialized) { - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -182,8 +184,9 @@ public KTable aggregate(final Initializer initializer, Objects.requireNonNull(subtractor, "subtractor can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal = - new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableImpl.java index bcd31bba16e94..21c15058bb9d2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableImpl.java @@ -154,7 +154,10 @@ public KTable filter(final Predicate predicate, final Materialized> materialized) { Objects.requireNonNull(predicate, "predicate can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - return doFilter(predicate, new MaterializedInternal<>(materialized, builder, FILTER_NAME), false); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, FILTER_NAME); + + return doFilter(predicate, materializedInternal, false); } @Override @@ -168,7 +171,10 @@ public KTable filterNot(final Predicate predicate, final Materialized> materialized) { Objects.requireNonNull(predicate, "predicate can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - return doFilter(predicate, new MaterializedInternal<>(materialized, builder, FILTER_NAME), true); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, FILTER_NAME); + + return doFilter(predicate, materializedInternal, true); } private KTable doMapValues(final ValueMapperWithKey mapper, @@ -210,7 +216,10 @@ public KTable mapValues(final ValueMapper m Objects.requireNonNull(mapper, "mapper can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - return doMapValues(withKey(mapper), new MaterializedInternal<>(materialized, builder, MAPVALUES_NAME)); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, MAPVALUES_NAME); + + return doMapValues(withKey(mapper), materializedInternal); } @Override @@ -219,7 +228,10 @@ public KTable mapValues(final ValueMapperWithKey(materialized, builder, MAPVALUES_NAME)); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, MAPVALUES_NAME); + + return doMapValues(mapper, materializedInternal); } @Override @@ -233,8 +245,10 @@ public KTable transformValues(final ValueTransformerWithKeySupplier< final Materialized> materialized, final String... stateStoreNames) { Objects.requireNonNull(materialized, "materialized can't be null"); - return doTransformValues(transformerSupplier, - new MaterializedInternal<>(materialized, builder, TRANSFORMVALUES_NAME), stateStoreNames); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, TRANSFORMVALUES_NAME); + + return doTransformValues(transformerSupplier, materializedInternal, stateStoreNames); } private KTable doTransformValues(final ValueTransformerWithKeySupplier transformerSupplier, @@ -304,7 +318,10 @@ public KTable join(final KTable other, Objects.requireNonNull(other, "other can't be null"); Objects.requireNonNull(joiner, "joiner can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - return doJoin(other, joiner, new MaterializedInternal<>(materialized, builder, MERGE_NAME), false, false); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, MERGE_NAME); + + return doJoin(other, joiner, materializedInternal, false, false); } @Override @@ -317,7 +334,10 @@ public KTable outerJoin(final KTable other, public KTable outerJoin(final KTable other, final ValueJoiner joiner, final Materialized> materialized) { - return doJoin(other, joiner, new MaterializedInternal<>(materialized, builder, MERGE_NAME), true, true); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, MERGE_NAME); + + return doJoin(other, joiner, materializedInternal, true, true); } @Override @@ -330,11 +350,9 @@ public KTable leftJoin(final KTable other, public KTable leftJoin(final KTable other, final ValueJoiner joiner, final Materialized> materialized) { - return doJoin(other, - joiner, - new MaterializedInternal<>(materialized, builder, MERGE_NAME), - true, - false); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, MERGE_NAME); + return doJoin(other, joiner, materializedInternal, true, false); } @SuppressWarnings("unchecked") diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedInternal.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedInternal.java index c933b8687f57e..5361e48c1d80e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedInternal.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedInternal.java @@ -25,18 +25,17 @@ public class MaterializedInternal extends Materialized { - private final boolean queryable; + private final boolean queriable; - - public MaterializedInternal(final Materialized materialized, - final InternalNameProvider nameProvider, - final String generatedStorePrefix) { + public MaterializedInternal(final Materialized materialized) { super(materialized); + queriable = storeName() != null; + } + + public void generateStoreNameIfNeeded(final InternalNameProvider nameProvider, + final String generatedStorePrefix) { if (storeName() == null) { - queryable = false; storeName = nameProvider.newStoreName(generatedStorePrefix); - } else { - queryable = true; } } @@ -63,7 +62,7 @@ public boolean loggingEnabled() { return loggingEnabled; } - public Map logConfig() { + Map logConfig() { return topicConfig; } @@ -72,6 +71,6 @@ boolean cachingEnabled() { } boolean isQueryable() { - return queryable; + return queriable; } } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java deleted file mode 100644 index 00784f96888a7..0000000000000 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/MaterializedPeek.java +++ /dev/null @@ -1,33 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.kafka.streams.kstream.internals; - -import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.processor.StateStore; - -class MaterializedPeek extends Materialized { - MaterializedPeek(final Materialized materialized) { - super(materialized); - } - - String givenStoreName() { - if (storeSupplier != null) { - return storeSupplier.name(); - } - return storeName; - } -} diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java index abd8d5b20e4a0..b3cbacd57f623 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/SessionWindowedKStreamImpl.java @@ -77,8 +77,7 @@ public KTable, Long> count(final Materialized(materialized).givenStoreName(); - if (givenStoreName == null) { + if (new MaterializedInternal<>(materialized).storeName() == null) { builder.newStoreName(AGGREGATE_NAME); } @@ -87,8 +86,8 @@ public KTable, Long> count(final Materialized, Long> doCount(final Materialized> materialized) { - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -115,8 +114,8 @@ public KTable, V> reduce(final Reducer reducer, Objects.requireNonNull(reducer, "reducer can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); final Aggregator reduceAggregator = aggregatorForReducer(reducer); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, REDUCE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, REDUCE_NAME); if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -148,8 +147,9 @@ public KTable, VR> aggregate(final Initializer initializer, Objects.requireNonNull(aggregator, "aggregator can't be null"); Objects.requireNonNull(sessionMerger, "sessionMerger can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java index 5b268aed805ba..4f5301b53b249 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TimeWindowedKStreamImpl.java @@ -72,8 +72,7 @@ public KTable, Long> count(final Materialized(materialized).givenStoreName(); - if (givenStoreName == null) { + if (new MaterializedInternal<>(materialized).storeName() == null) { builder.newStoreName(AGGREGATE_NAME); } @@ -82,8 +81,9 @@ public KTable, Long> count(final Materialized, Long> doCount(final Materialized> materialized) { - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -112,8 +112,8 @@ public KTable, VR> aggregate(final Initializer initializer, Objects.requireNonNull(initializer, "initializer can't be null"); Objects.requireNonNull(aggregator, "aggregator can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME); if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } @@ -134,8 +134,9 @@ public KTable, V> reduce(final Reducer reducer, final Materialize Objects.requireNonNull(reducer, "reducer can't be null"); Objects.requireNonNull(materialized, "materialized can't be null"); - final MaterializedInternal> materializedInternal - = new MaterializedInternal<>(materialized, builder, REDUCE_NAME); + final MaterializedInternal> materializedInternal = new MaterializedInternal<>(materialized); + materializedInternal.generateStoreNameIfNeeded(builder, REDUCE_NAME); + if (materializedInternal.keySerde() == null) { materializedInternal.withKeySerde(keySerde); } diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java index 79bf81e6bb2fb..63432ffc43927 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilderTest.java @@ -33,7 +33,6 @@ import org.apache.kafka.test.MockMapper; import org.apache.kafka.test.MockTimestampExtractor; import org.apache.kafka.test.MockValueJoiner; -import org.junit.Before; import org.junit.Test; import java.util.Collections; @@ -58,12 +57,11 @@ public class InternalStreamsBuilderTest { private final InternalStreamsBuilder builder = new InternalStreamsBuilder(new InternalTopologyBuilder()); private final ConsumedInternal consumed = new ConsumedInternal<>(); private final String storePrefix = "prefix-"; - private MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("test-store"), builder, storePrefix); + private final MaterializedInternal> materialized = new MaterializedInternal<>(Materialized.as("test-store")); - @Before - public void setUp() { + { builder.internalTopologyBuilder.setApplicationId(APP_ID); + materialized.generateStoreNameIfNeeded(builder, storePrefix); } @Test @@ -127,12 +125,10 @@ public boolean test(final String key, final String value) { @Test public void shouldStillMaterializeSourceKTableIfMaterializedIsntQueryable() { - KTable table1 = builder.table("topic2", - consumed, - new MaterializedInternal<>( - Materialized.>with(null, null), - builder, - storePrefix)); + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.with(null, null)); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + final KTable table1 = builder.table("topic2", consumed, materializedInternal); final ProcessorTopology topology = builder.internalTopologyBuilder.build(null); @@ -147,38 +143,33 @@ public void shouldStillMaterializeSourceKTableIfMaterializedIsntQueryable() { @Test public void shouldBuildGlobalTableWithNonQueryableStoreName() { - final GlobalKTable table1 = builder.globalTable( - "topic2", - consumed, - new MaterializedInternal<>( - Materialized.>with(null, null), - builder, - storePrefix)); + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.with(null, null)); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + + final GlobalKTable table1 = builder.globalTable("topic2", consumed, materializedInternal); assertNull(table1.queryableStoreName()); } @Test public void shouldBuildGlobalTableWithQueryaIbleStoreName() { - final GlobalKTable table1 = builder.globalTable( - "topic2", - consumed, - new MaterializedInternal<>( - Materialized.>as("globalTable"), - builder, - storePrefix)); + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.as("globalTable")); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + final GlobalKTable table1 = builder.globalTable("topic2", consumed, materializedInternal); assertEquals("globalTable", table1.queryableStoreName()); } @Test public void shouldBuildSimpleGlobalTableTopology() { + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.as("globalTable")); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); builder.globalTable("table", consumed, - new MaterializedInternal<>( - Materialized.>as("globalTable"), - builder, - storePrefix)); + materializedInternal); final ProcessorTopology topology = builder.internalTopologyBuilder.buildGlobalStateTopology(); final List stateStores = topology.globalStateStores(); @@ -199,15 +190,18 @@ private void doBuildGlobalTopologyWithAllGlobalTables() { @Test public void shouldBuildGlobalTopologyWithAllGlobalTables() { - builder.globalTable("table", - consumed, - new MaterializedInternal<>( - Materialized.>as("global1"), builder, storePrefix)); - builder.globalTable("table2", - consumed, - new MaterializedInternal<>( - Materialized.>as("global2"), builder, storePrefix)); - + { + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.as("global1")); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + builder.globalTable("table", consumed, materializedInternal); + } + { + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.as("global2")); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + builder.globalTable("table2", consumed, materializedInternal); + } doBuildGlobalTopologyWithAllGlobalTables(); } @@ -216,25 +210,22 @@ public void shouldAddGlobalTablesToEachGroup() { final String one = "globalTable"; final String two = "globalTable2"; - final GlobalKTable globalTable = builder.globalTable("table", - consumed, - new MaterializedInternal<>( - Materialized.>as(one), builder, storePrefix)); - final GlobalKTable globalTable2 = builder.globalTable("table2", - consumed, - new MaterializedInternal<>( - Materialized.>as(two), builder, storePrefix)); + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.as(one)); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + final GlobalKTable globalTable = builder.globalTable("table", consumed, materializedInternal); - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("not-global"), builder, storePrefix); - builder.table("not-global", consumed, materialized); + final MaterializedInternal> materializedInternal2 = + new MaterializedInternal<>(Materialized.as(two)); + materializedInternal2.generateStoreNameIfNeeded(builder, storePrefix); + final GlobalKTable globalTable2 = builder.globalTable("table2", consumed, materializedInternal2); - final KeyValueMapper kvMapper = new KeyValueMapper() { - @Override - public String apply(final String key, final String value) { - return value; - } - }; + final MaterializedInternal> materializedInternalNotGlobal = + new MaterializedInternal<>(Materialized.as("not-global")); + materializedInternalNotGlobal.generateStoreNameIfNeeded(builder, storePrefix); + builder.table("not-global", consumed, materializedInternalNotGlobal); + + final KeyValueMapper kvMapper = (key, value) -> value; final KStream stream = builder.stream(Collections.singleton("t1"), consumed); stream.leftJoin(globalTable, kvMapper, MockValueJoiner.TOSTRING_JOINER); @@ -260,9 +251,10 @@ public String apply(final String key, final String value) { public void shouldMapStateStoresToCorrectSourceTopics() { final KStream playEvents = builder.stream(Collections.singleton("events"), consumed); - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("table-store"), builder, storePrefix); - final KTable table = builder.table("table-topic", consumed, materialized); + final MaterializedInternal> materializedInternal = + new MaterializedInternal<>(Materialized.as("table-store")); + materializedInternal.generateStoreNameIfNeeded(builder, storePrefix); + final KTable table = builder.table("table-topic", consumed, materializedInternal); assertEquals(Collections.singletonList("table-topic"), builder.internalTopologyBuilder.stateStoreNameToSourceTopics().get("table-store")); final KStream mapped = playEvents.map(MockMapper.selectValueKeyValueMapper()); diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/MaterializedInternalTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/MaterializedInternalTest.java index 5fd76f3acfe47..1a83ef13152ad 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/MaterializedInternalTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/MaterializedInternalTest.java @@ -49,8 +49,9 @@ public void shouldGenerateStoreNameWithPrefixIfProvidedNameIsNull() { EasyMock.replay(nameProvider); - final MaterializedInternal materialized - = new MaterializedInternal<>(Materialized.with(null, null), nameProvider, prefix); + final MaterializedInternal materialized = + new MaterializedInternal<>(Materialized.with(null, null)); + materialized.generateStoreNameIfNeeded(nameProvider, prefix); assertThat(materialized.storeName(), equalTo(generatedName)); EasyMock.verify(nameProvider); @@ -59,8 +60,9 @@ public void shouldGenerateStoreNameWithPrefixIfProvidedNameIsNull() { @Test public void shouldUseProvidedStoreNameWhenSet() { final String storeName = "store-name"; - final MaterializedInternal materialized - = new MaterializedInternal<>(Materialized.as(storeName), nameProvider, prefix); + final MaterializedInternal materialized = + new MaterializedInternal<>(Materialized.as(storeName)); + materialized.generateStoreNameIfNeeded(nameProvider, prefix); assertThat(materialized.storeName(), equalTo(storeName)); } @@ -69,8 +71,9 @@ public void shouldUseStoreNameOfSupplierWhenProvided() { final String storeName = "other-store-name"; EasyMock.expect(supplier.name()).andReturn(storeName).anyTimes(); EasyMock.replay(supplier); - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.as(supplier), nameProvider, prefix); + final MaterializedInternal> materialized = + new MaterializedInternal<>(Materialized.as(supplier)); + materialized.generateStoreNameIfNeeded(nameProvider, prefix); assertThat(materialized.storeName(), equalTo(storeName)); } } \ No newline at end of file diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java index c71f4690c4bfd..f3e9299cb509f 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java @@ -73,19 +73,21 @@ public class GlobalStreamThreadTest { @Before public void before() { final MaterializedInternal> materialized = new MaterializedInternal<>( - Materialized.>with(null, null), + Materialized.with(null, null)); + materialized.generateStoreNameIfNeeded( new InternalNameProvider() { @Override - public String newProcessorName(String prefix) { + public String newProcessorName(final String prefix) { return "processorName"; } @Override - public String newStoreName(String prefix) { + public String newStoreName(final String prefix) { return GLOBAL_STORE_NAME; } }, - "store-"); + "store-" + ); builder.addGlobalStore( (StoreBuilder) new KeyValueStoreMaterializer<>(materialized).materialize().withLoggingDisabled(), diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/KeyValueStoreMaterializerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/KeyValueStoreMaterializerTest.java index 9ba86ac7e141e..fc243c61f8dfd 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/KeyValueStoreMaterializerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/KeyValueStoreMaterializerTest.java @@ -53,10 +53,10 @@ public class KeyValueStoreMaterializerTest { @Test public void shouldCreateBuilderThatBuildsMeteredStoreWithCachingAndLoggingEnabled() { - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("store"), - nameProvider, - storePrefix); + final MaterializedInternal> materialized = + new MaterializedInternal<>(Materialized.as("store")); + materialized.generateStoreNameIfNeeded(nameProvider, storePrefix); + final KeyValueStoreMaterializer materializer = new KeyValueStoreMaterializer<>(materialized); final StoreBuilder> builder = materializer.materialize(); final KeyValueStore store = builder.build(); @@ -69,9 +69,10 @@ public void shouldCreateBuilderThatBuildsMeteredStoreWithCachingAndLoggingEnable @Test public void shouldCreateBuilderThatBuildsStoreWithCachingDisabled() { - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("store") - .withCachingDisabled(), nameProvider, storePrefix); + final MaterializedInternal> materialized = new MaterializedInternal<>( + Materialized.>as("store").withCachingDisabled() + ); + materialized.generateStoreNameIfNeeded(nameProvider, storePrefix); final KeyValueStoreMaterializer materializer = new KeyValueStoreMaterializer<>(materialized); final StoreBuilder> builder = materializer.materialize(); final KeyValueStore store = builder.build(); @@ -81,9 +82,11 @@ public void shouldCreateBuilderThatBuildsStoreWithCachingDisabled() { @Test public void shouldCreateBuilderThatBuildsStoreWithLoggingDisabled() { - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("store") - .withLoggingDisabled(), nameProvider, storePrefix); + final MaterializedInternal> materialized = new MaterializedInternal<>( + Materialized.>as("store") + .withLoggingDisabled() + ); + materialized.generateStoreNameIfNeeded(nameProvider, storePrefix); final KeyValueStoreMaterializer materializer = new KeyValueStoreMaterializer<>(materialized); final StoreBuilder> builder = materializer.materialize(); final KeyValueStore store = builder.build(); @@ -94,10 +97,11 @@ public void shouldCreateBuilderThatBuildsStoreWithLoggingDisabled() { @Test public void shouldCreateBuilderThatBuildsStoreWithCachingAndLoggingDisabled() { - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.>as("store") + final MaterializedInternal> materialized = new MaterializedInternal<>( + Materialized.>as("store") .withCachingDisabled() - .withLoggingDisabled(), nameProvider, storePrefix); + .withLoggingDisabled()); + materialized.generateStoreNameIfNeeded(nameProvider, storePrefix); final KeyValueStoreMaterializer materializer = new KeyValueStoreMaterializer<>(materialized); final StoreBuilder> builder = materializer.materialize(); final KeyValueStore store = builder.build(); @@ -114,8 +118,9 @@ public void shouldCreateKeyValueStoreWithTheProvidedInnerStore() { EasyMock.expect(supplier.get()).andReturn(store); EasyMock.replay(supplier); - final MaterializedInternal> materialized - = new MaterializedInternal<>(Materialized.as(supplier), nameProvider, storePrefix); + final MaterializedInternal> materialized = + new MaterializedInternal<>(Materialized.as(supplier)); + materialized.generateStoreNameIfNeeded(nameProvider, storePrefix); final KeyValueStoreMaterializer materializer = new KeyValueStoreMaterializer<>(materialized); final StoreBuilder> builder = materializer.materialize(); final KeyValueStore built = builder.build(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index f3dce52199836..936c67b8ef838 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -829,7 +829,9 @@ public void shouldUpdateStandbyTask() { final TopicPartition partition2 = t2p1; internalStreamsBuilder.stream(Collections.singleton(topic1), consumed) .groupByKey().count(Materialized.>as(storeName1)); - internalStreamsBuilder.table(topic2, new ConsumedInternal(), new MaterializedInternal(Materialized.as(storeName2), internalStreamsBuilder, "")); + final MaterializedInternal materialized = new MaterializedInternal(Materialized.as(storeName2)); + materialized.generateStoreNameIfNeeded(internalStreamsBuilder, ""); + internalStreamsBuilder.table(topic2, new ConsumedInternal(), materialized); final StreamThread thread = createStreamThread(clientId, config, false); final MockConsumer restoreConsumer = clientSupplier.restoreConsumer;