diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/Named.java b/streams/src/main/java/org/apache/kafka/streams/kstream/Named.java index 1db031a027955..84bb819700a42 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/Named.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/Named.java @@ -26,6 +26,10 @@ public class Named implements NamedOperation { protected String name; + protected Named(final Named named) { + this(Objects.requireNonNull(named, "named can't be null").name); + } + protected Named(final String name) { this.name = name; if (name != null) { @@ -51,7 +55,7 @@ public Named withName(final String name) { return new Named(name); } - static void validate(final String name) { + protected static void validate(final String name) { if (name.isEmpty()) throw new TopologyException("Name is illegal, it can't be empty"); if (name.equals(".") || name.equals("..")) diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java index e7a7678aa0247..3a90fd222d484 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/InternalStreamsBuilder.java @@ -118,7 +118,7 @@ public KTable table(final String topic, final String sourceName = new NamedInternal(consumed.name()) .orElseGenerateWithPrefix(this, KStreamImpl.SOURCE_NAME); final String tableSourceName = new NamedInternal(consumed.name()) - .suffixWithOrElseGet("-table-source", () -> newProcessorName(KTableImpl.SOURCE_NAME)); + .suffixWithOrElseGet("-table-source", this, KTableImpl.SOURCE_NAME); final KTableSource tableSource = new KTableSource<>(materialized.storeName(), materialized.queryableStoreName()); final ProcessorParameters processorParameters = new ProcessorParameters<>(tableSource, tableSourceName); diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/NamedInternal.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/NamedInternal.java index e83728e92404c..d478e9b6f64b0 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/NamedInternal.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/NamedInternal.java @@ -17,8 +17,6 @@ package org.apache.kafka.streams.kstream.internals; import org.apache.kafka.streams.kstream.Named; -import java.util.Optional; -import java.util.function.Supplier; public class NamedInternal extends Named { @@ -33,14 +31,14 @@ public static NamedInternal with(final String name) { /** * Creates a new {@link NamedInternal} instance. * - * @param internal the internal name. + * @param internal the internal name. */ NamedInternal(final String internal) { super(internal); } /** - * @return a string name. + * @return a string name. */ public String name() { return name; @@ -51,31 +49,31 @@ public NamedInternal withName(final String name) { return new NamedInternal(name); } - /** - * Check whether an internal name is defined. - * @return {@code false} if no name is set. - */ - public boolean isDefined() { - return name != null; - } + String suffixWithOrElseGet(final String suffix, final InternalNameProvider provider, final String prefix) { + // We actually do not need to generate processor names for operation if a name is specified. + // But before returning, we still need to burn index for the operation to keep topology backward compatibility. + if (name != null) { + provider.newProcessorName(prefix); + + final String suffixed = name + suffix; + // Re-validate generated name as suffixed string could be too large. + Named.validate(suffixed); - String suffixWithOrElseGet(final String suffix, final Supplier supplier) { - final Optional suffixed = Optional.ofNullable(this.name).map(s -> s + suffix); - // Creating a new named will re-validate generated name as suffixed string could be too large. - return new NamedInternal(suffixed.orElseGet(supplier)).name(); + return suffixed; + } else { + return provider.newProcessorName(prefix); + } } String orElseGenerateWithPrefix(final InternalNameProvider provider, final String prefix) { - return orElseGet(() -> provider.newProcessorName(prefix)); + // We actually do not need to generate processor names for operation if a name is specified. + // But before returning, we still need to burn index for the operation to keep topology backward compatibility. + if (name != null) { + provider.newProcessorName(prefix); + return name; + } else { + return provider.newProcessorName(prefix); + } } - /** - * Returns the internal name or the value returns from the supplier. - * - * @param supplier the supplier to be used if internal name is empty. - * @return an internal string name. - */ - private String orElseGet(final Supplier supplier) { - return Optional.ofNullable(this.name).orElseGet(supplier); - } } diff --git a/streams/src/test/java/org/apache/kafka/streams/StreamsBuilderTest.java b/streams/src/test/java/org/apache/kafka/streams/StreamsBuilderTest.java index 669cece85425f..93d444bb5a32f 100644 --- a/streams/src/test/java/org/apache/kafka/streams/StreamsBuilderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/StreamsBuilderTest.java @@ -424,7 +424,7 @@ public void shouldUseSpecifiedNameForStreamSourceProcessor() { builder.stream(STREAM_TOPIC_TWO); builder.build(); final ProcessorTopology topology = builder.internalTopologyBuilder.rewriteTopology(new StreamsConfig(props)).build(); - assertSpecifiedNameForOperation(topology, expected, "KSTREAM-SOURCE-0000000000"); + assertSpecifiedNameForOperation(topology, expected, "KSTREAM-SOURCE-0000000001"); } @Test @@ -440,8 +440,8 @@ public void shouldUseSpecifiedNameForTableSourceProcessor() { topology, expected, expected + "-table-source", - "KSTREAM-SOURCE-0000000002", - "KTABLE-SOURCE-0000000003"); + "KSTREAM-SOURCE-0000000004", + "KTABLE-SOURCE-0000000005"); } @Test @@ -467,7 +467,7 @@ public void shouldUseSpecifiedNameForSinkProcessor() { stream.to(STREAM_TOPIC_TWO); builder.build(); final ProcessorTopology topology = builder.internalTopologyBuilder.rewriteTopology(new StreamsConfig(props)).build(); - assertSpecifiedNameForOperation(topology, "KSTREAM-SOURCE-0000000000", expected, "KSTREAM-SINK-0000000001"); + assertSpecifiedNameForOperation(topology, "KSTREAM-SOURCE-0000000000", expected, "KSTREAM-SINK-0000000002"); } private void assertSpecifiedNameForOperation(final ProcessorTopology topology, final String... expected) { diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/NamedInternalTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/NamedInternalTest.java index 98b3a4d776176..1c4c700a2906f 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/NamedInternalTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/NamedInternalTest.java @@ -22,45 +22,55 @@ public class NamedInternalTest { - private static final String TEST_VALUE = "default-value"; + private static final String TEST_PREFIX = "prefix-"; + private static final String TEST_VALUE = "default-value"; private static final String TEST_SUFFIX = "-suffix"; + private static class TestNameProvider implements InternalNameProvider { + int index = 0; + + @Override + public String newProcessorName(final String prefix) { + return prefix + "PROCESSOR-" + index++; + } + + @Override + public String newStoreName(final String prefix) { + return prefix + "STORE-" + index++; + } + + } + @Test public void shouldSuffixNameOrReturnProviderValue() { final String name = "foo"; + final TestNameProvider provider = new TestNameProvider(); + assertEquals( - name + TEST_SUFFIX, - NamedInternal.with(name).suffixWithOrElseGet(TEST_SUFFIX, () -> TEST_VALUE) + name + TEST_SUFFIX, + NamedInternal.with(name).suffixWithOrElseGet(TEST_SUFFIX, provider, TEST_PREFIX) ); + + // 1, not 0, indicates that the named call still burned an index number. assertEquals( - TEST_VALUE, - NamedInternal.with(null).suffixWithOrElseGet(TEST_SUFFIX, () -> TEST_VALUE) + "prefix-PROCESSOR-1", + NamedInternal.with(null).suffixWithOrElseGet(TEST_SUFFIX, provider, TEST_PREFIX) ); } @Test public void shouldGenerateWithPrefixGivenEmptyName() { final String prefix = "KSTREAM-MAP-"; - assertEquals(prefix + "PROCESSOR-NAME", NamedInternal.with(null).orElseGenerateWithPrefix( - new InternalNameProvider() { - @Override - public String newProcessorName(final String prefix) { - return prefix + "PROCESSOR-NAME"; - } - - @Override - public String newStoreName(final String prefix) { - return null; - } - }, - prefix) + assertEquals(prefix + "PROCESSOR-0", NamedInternal.with(null).orElseGenerateWithPrefix( + new TestNameProvider(), + prefix) ); } @Test public void shouldNotGenerateWithPrefixGivenValidName() { final String validName = "validName"; - assertEquals(validName, NamedInternal.with(validName).orElseGenerateWithPrefix(null, "KSTREAM-MAP-") + assertEquals(validName, NamedInternal.with(validName).orElseGenerateWithPrefix(new TestNameProvider(), "KSTREAM-MAP-") ); } } \ No newline at end of file