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 75942ce7608a4..494969460f005 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 @@ -39,13 +39,13 @@ import org.apache.kafka.streams.kstream.ValueTransformerWithKeySupplier; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.CombinedKey; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.CombinedKeySchema; -import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.ForeignJoinSubscriptionProcessorSupplier; -import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.ForeignJoinSubscriptionSendProcessorSupplier; -import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionJoinForeignProcessorSupplier; -import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionResolverJoinProcessorSupplier; +import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.ForeignTableJoinProcessorSupplier; +import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionSendProcessorSupplier; +import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionJoinProcessorSupplier; +import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.ResponseJoinProcessorSupplier; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionResponseWrapper; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionResponseWrapperSerde; -import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionStoreReceiveProcessorSupplier; +import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionReceiveProcessorSupplier; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapper; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapperSerde; import org.apache.kafka.streams.kstream.internals.graph.KTableKTableJoinNode; @@ -1139,9 +1139,9 @@ private KTable doJoinOnForeignKey(final KTable forei ); final KTableValueGetterSupplier primaryKeyValueGetter = valueGetterSupplier(); - final StatefulProcessorNode> subscriptionNode = new StatefulProcessorNode<>( + final StatefulProcessorNode> subscriptionSendNode = new StatefulProcessorNode<>( new ProcessorParameters<>( - new ForeignJoinSubscriptionSendProcessorSupplier<>( + new SubscriptionSendProcessorSupplier<>( foreignKeyExtractor, subscriptionForeignKeySerdePseudoTopic, valueHashSerdePseudoTopic, @@ -1155,7 +1155,7 @@ private KTable doJoinOnForeignKey(final KTable forei Collections.emptySet(), Collections.singleton(primaryKeyValueGetter) ); - builder.addGraphNode(graphNode, subscriptionNode); + builder.addGraphNode(graphNode, subscriptionSendNode); final StreamPartitioner> subscriptionSinkPartitioner = tableJoinedInternal.otherPartitioner() == null @@ -1167,7 +1167,7 @@ private KTable doJoinOnForeignKey(final KTable forei new StaticTopicNameExtractor<>(subscriptionTopicName), new ProducedInternal<>(Produced.with(foreignKeySerde, subscriptionWrapperSerde, subscriptionSinkPartitioner)) ); - builder.addGraphNode(subscriptionNode, subscriptionSink); + builder.addGraphNode(subscriptionSendNode, subscriptionSink); final StreamSourceNode> subscriptionSource = new StreamSourceNode<>( renamed.suffixWithOrElseGet("-subscription-registration-source", builder, SOURCE_NAME), @@ -1197,7 +1197,7 @@ private KTable doJoinOnForeignKey(final KTable forei final StatefulProcessorNode> subscriptionReceiveNode = new StatefulProcessorNode<>( new ProcessorParameters<>( - new SubscriptionStoreReceiveProcessorSupplier<>(subscriptionStore, combinedKeySchema), + new SubscriptionReceiveProcessorSupplier<>(subscriptionStore, combinedKeySchema), renamed.suffixWithOrElseGet("-subscription-receive", builder, SUBSCRIPTION_PROCESSOR) ), Collections.singleton(subscriptionStore), @@ -1206,10 +1206,10 @@ private KTable doJoinOnForeignKey(final KTable forei builder.addGraphNode(subscriptionSource, subscriptionReceiveNode); final KTableValueGetterSupplier foreignKeyValueGetter = ((KTableImpl) foreignKeyTable).valueGetterSupplier(); - final StatefulProcessorNode, Change>>> subscriptionJoinForeignNode = + final StatefulProcessorNode, Change>>> subscriptionJoinNode = new StatefulProcessorNode<>( new ProcessorParameters<>( - new SubscriptionJoinForeignProcessorSupplier<>( + new SubscriptionJoinProcessorSupplier<>( foreignKeyValueGetter ), renamed.suffixWithOrElseGet("-subscription-join-foreign", builder, SUBSCRIPTION_PROCESSOR) @@ -1217,17 +1217,17 @@ private KTable doJoinOnForeignKey(final KTable forei Collections.emptySet(), Collections.singleton(foreignKeyValueGetter) ); - builder.addGraphNode(subscriptionReceiveNode, subscriptionJoinForeignNode); + builder.addGraphNode(subscriptionReceiveNode, subscriptionJoinNode); - final StatefulProcessorNode> foreignJoinSubscriptionNode = new StatefulProcessorNode<>( + final StatefulProcessorNode> foreignTableJoinNode = new StatefulProcessorNode<>( new ProcessorParameters<>( - new ForeignJoinSubscriptionProcessorSupplier<>(subscriptionStore, combinedKeySchema, foreignKeyValueGetter), + new ForeignTableJoinProcessorSupplier<>(subscriptionStore, combinedKeySchema, foreignKeyValueGetter), renamed.suffixWithOrElseGet("-foreign-join-subscription", builder, SUBSCRIPTION_PROCESSOR) ), Collections.singleton(subscriptionStore), Collections.singleton(foreignKeyValueGetter) ); - builder.addGraphNode(((KTableImpl) foreignKeyTable).graphNode, foreignJoinSubscriptionNode); + builder.addGraphNode(((KTableImpl) foreignKeyTable).graphNode, foreignTableJoinNode); final String finalRepartitionTopicName = renamed.suffixWithOrElseGet("-subscription-response", builder, SUBSCRIPTION_RESPONSE) + TOPIC_SUFFIX; @@ -1244,8 +1244,8 @@ private KTable doJoinOnForeignKey(final KTable forei new StaticTopicNameExtractor<>(finalRepartitionTopicName), new ProducedInternal<>(Produced.with(keySerde, responseWrapperSerde, foreignResponseSinkPartitioner)) ); - builder.addGraphNode(subscriptionJoinForeignNode, foreignResponseSink); - builder.addGraphNode(foreignJoinSubscriptionNode, foreignResponseSink); + builder.addGraphNode(subscriptionJoinNode, foreignResponseSink); + builder.addGraphNode(foreignTableJoinNode, foreignResponseSink); final StreamSourceNode> foreignResponseSource = new StreamSourceNode<>( renamed.suffixWithOrElseGet("-subscription-response-source", builder, SOURCE_NAME), @@ -1259,22 +1259,21 @@ private KTable doJoinOnForeignKey(final KTable forei resultSourceNodes.add(foreignResponseSource.nodeName()); builder.internalTopologyBuilder.copartitionSources(resultSourceNodes); - final SubscriptionResolverJoinProcessorSupplier resolverProcessorSupplier = new SubscriptionResolverJoinProcessorSupplier<>( - primaryKeyValueGetter, - valueSerde == null ? null : valueSerde.serializer(), - valueHashSerdePseudoTopic, - joiner, - leftJoin - ); - final StatefulProcessorNode> resolverNode = new StatefulProcessorNode<>( + final StatefulProcessorNode> responseJoinNode = new StatefulProcessorNode<>( new ProcessorParameters<>( - resolverProcessorSupplier, + new ResponseJoinProcessorSupplier<>( + primaryKeyValueGetter, + valueSerde == null ? null : valueSerde.serializer(), + valueHashSerdePseudoTopic, + joiner, + leftJoin + ), renamed.suffixWithOrElseGet("-subscription-response-resolver", builder, SUBSCRIPTION_RESPONSE_RESOLVER_PROCESSOR) ), Collections.emptySet(), Collections.singleton(primaryKeyValueGetter) ); - builder.addGraphNode(foreignResponseSource, resolverNode); + builder.addGraphNode(foreignResponseSource, responseJoinNode); final String resultProcessorName = renamed.suffixWithOrElseGet("-result", builder, FK_JOIN_OUTPUT_NAME); @@ -1308,7 +1307,7 @@ private KTable doJoinOnForeignKey(final KTable forei resultStore ); resultNode.setOutputVersioned(materializedInternal.storeSupplier() instanceof VersionedBytesStoreSupplier); - builder.addGraphNode(resolverNode, resultNode); + builder.addGraphNode(responseJoinNode, resultNode); return new KTableImpl( resultProcessorName, diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionProcessorSupplier.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignTableJoinProcessorSupplier.java similarity index 97% rename from streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionProcessorSupplier.java rename to streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignTableJoinProcessorSupplier.java index 46e2bd24c2577..a9cae20337813 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionProcessorSupplier.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignTableJoinProcessorSupplier.java @@ -40,14 +40,14 @@ import java.nio.ByteBuffer; -public class ForeignJoinSubscriptionProcessorSupplier implements +public class ForeignTableJoinProcessorSupplier implements ProcessorSupplier, K, SubscriptionResponseWrapper> { - private static final Logger LOG = LoggerFactory.getLogger(ForeignJoinSubscriptionProcessorSupplier.class); + private static final Logger LOG = LoggerFactory.getLogger(ForeignTableJoinProcessorSupplier.class); private final StoreBuilder>> storeBuilder; private final CombinedKeySchema keySchema; private final KTableValueGetterSupplier foreignKeyValueGetterSupplier; - public ForeignJoinSubscriptionProcessorSupplier( + public ForeignTableJoinProcessorSupplier( final StoreBuilder>> storeBuilder, final CombinedKeySchema keySchema, final KTableValueGetterSupplier foreignKeyValueGetterSupplier) { diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResolverJoinProcessorSupplier.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java similarity index 89% rename from streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResolverJoinProcessorSupplier.java rename to streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java index c7c0aee7950d4..2fffa89b348e5 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResolverJoinProcessorSupplier.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplier.java @@ -42,18 +42,18 @@ * @param Type of foreign values * @param Type of joined result of primary and foreign values */ -public class SubscriptionResolverJoinProcessorSupplier implements ProcessorSupplier, K, VR> { +public class ResponseJoinProcessorSupplier implements ProcessorSupplier, K, VR> { private final KTableValueGetterSupplier valueGetterSupplier; private final Serializer constructionTimeValueSerializer; private final Supplier valueHashSerdePseudoTopicSupplier; private final ValueJoiner joiner; private final boolean leftJoin; - public SubscriptionResolverJoinProcessorSupplier(final KTableValueGetterSupplier valueGetterSupplier, - final Serializer valueSerializer, - final Supplier valueHashSerdePseudoTopicSupplier, - final ValueJoiner joiner, - final boolean leftJoin) { + public ResponseJoinProcessorSupplier(final KTableValueGetterSupplier valueGetterSupplier, + final Serializer valueSerializer, + final Supplier valueHashSerdePseudoTopicSupplier, + final ValueJoiner joiner, + final boolean leftJoin) { this.valueGetterSupplier = valueGetterSupplier; constructionTimeValueSerializer = valueSerializer; this.valueHashSerdePseudoTopicSupplier = valueHashSerdePseudoTopicSupplier; diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionJoinForeignProcessorSupplier.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionJoinProcessorSupplier.java similarity index 97% rename from streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionJoinForeignProcessorSupplier.java rename to streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionJoinProcessorSupplier.java index 56d6a13321ffb..490f70ffc8d38 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionJoinForeignProcessorSupplier.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionJoinProcessorSupplier.java @@ -39,12 +39,12 @@ * @param Type of foreign key * @param Type of foreign value */ -public class SubscriptionJoinForeignProcessorSupplier +public class SubscriptionJoinProcessorSupplier implements ProcessorSupplier, Change>>, K, SubscriptionResponseWrapper> { private final KTableValueGetterSupplier foreignValueGetterSupplier; - public SubscriptionJoinForeignProcessorSupplier(final KTableValueGetterSupplier foreignValueGetterSupplier) { + public SubscriptionJoinProcessorSupplier(final KTableValueGetterSupplier foreignValueGetterSupplier) { this.foreignValueGetterSupplier = foreignValueGetterSupplier; } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionStoreReceiveProcessorSupplier.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionReceiveProcessorSupplier.java similarity index 97% rename from streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionStoreReceiveProcessorSupplier.java rename to streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionReceiveProcessorSupplier.java index dcbf6c0eaf734..5c386cc735bb5 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionStoreReceiveProcessorSupplier.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionReceiveProcessorSupplier.java @@ -35,14 +35,14 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public class SubscriptionStoreReceiveProcessorSupplier +public class SubscriptionReceiveProcessorSupplier implements ProcessorSupplier, CombinedKey, Change>>> { - private static final Logger LOG = LoggerFactory.getLogger(SubscriptionStoreReceiveProcessorSupplier.class); + private static final Logger LOG = LoggerFactory.getLogger(SubscriptionReceiveProcessorSupplier.class); private final StoreBuilder>> storeBuilder; private final CombinedKeySchema keySchema; - public SubscriptionStoreReceiveProcessorSupplier( + public SubscriptionReceiveProcessorSupplier( final StoreBuilder>> storeBuilder, final CombinedKeySchema keySchema) { diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionSendProcessorSupplier.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionSendProcessorSupplier.java similarity index 92% rename from streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionSendProcessorSupplier.java rename to streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionSendProcessorSupplier.java index 3d8b5dd222ea8..feb06b02e13ea 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionSendProcessorSupplier.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionSendProcessorSupplier.java @@ -45,8 +45,8 @@ import static org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapper.Instruction.PROPAGATE_NULL_IF_NO_FK_VAL_AVAILABLE; import static org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapper.Instruction.PROPAGATE_ONLY_IF_FK_VAL_AVAILABLE; -public class ForeignJoinSubscriptionSendProcessorSupplier implements ProcessorSupplier, KO, SubscriptionWrapper> { - private static final Logger LOG = LoggerFactory.getLogger(ForeignJoinSubscriptionSendProcessorSupplier.class); +public class SubscriptionSendProcessorSupplier implements ProcessorSupplier, KO, SubscriptionWrapper> { + private static final Logger LOG = LoggerFactory.getLogger(SubscriptionSendProcessorSupplier.class); private final Function foreignKeyExtractor; private final Supplier foreignKeySerdeTopicSupplier; @@ -56,13 +56,13 @@ public class ForeignJoinSubscriptionSendProcessorSupplier implements P private Serializer foreignKeySerializer; private Serializer valueSerializer; - public ForeignJoinSubscriptionSendProcessorSupplier(final Function foreignKeyExtractor, - final Supplier foreignKeySerdeTopicSupplier, - final Supplier valueSerdeTopicSupplier, - final Serde foreignKeySerde, - final Serializer valueSerializer, - final boolean leftJoin, - final KTableValueGetterSupplier primaryKeyValueGetterSupplier) { + public SubscriptionSendProcessorSupplier(final Function foreignKeyExtractor, + final Supplier foreignKeySerdeTopicSupplier, + final Supplier valueSerdeTopicSupplier, + final Serde foreignKeySerde, + final Serializer valueSerializer, + final boolean leftJoin, + final KTableValueGetterSupplier primaryKeyValueGetterSupplier) { this.foreignKeyExtractor = foreignKeyExtractor; this.foreignKeySerdeTopicSupplier = foreignKeySerdeTopicSupplier; this.valueSerdeTopicSupplier = valueSerdeTopicSupplier; diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/ForeignKeyJoinSuite.java b/streams/src/test/java/org/apache/kafka/streams/integration/ForeignKeyJoinSuite.java index 5dd8c054eab97..6be9a6808fc5a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/ForeignKeyJoinSuite.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/ForeignKeyJoinSuite.java @@ -19,7 +19,7 @@ import org.apache.kafka.common.utils.BytesTest; import org.apache.kafka.streams.kstream.internals.KTableKTableForeignKeyJoinScenarioTest; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.CombinedKeySchemaTest; -import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionResolverJoinProcessorSupplierTest; +import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.ResponseJoinProcessorSupplierTest; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionResponseWrapperSerdeTest; import org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapperSerdeTest; import org.junit.runner.RunWith; @@ -44,7 +44,7 @@ CombinedKeySchemaTest.class, SubscriptionWrapperSerdeTest.class, SubscriptionResponseWrapperSerdeTest.class, - SubscriptionResolverJoinProcessorSupplierTest.class + ResponseJoinProcessorSupplierTest.class }) public class ForeignKeyJoinSuite { } diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionProcessorSupplierTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignTableJoinProcessorSupplierTest.java similarity index 98% rename from streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionProcessorSupplierTest.java rename to streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignTableJoinProcessorSupplierTest.java index 6757ccca20bb9..f4f35e6ff05e6 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignJoinSubscriptionProcessorSupplierTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ForeignTableJoinProcessorSupplierTest.java @@ -32,7 +32,7 @@ import org.junit.Assert; import org.junit.Test; -public class ForeignJoinSubscriptionProcessorSupplierTest { +public class ForeignTableJoinProcessorSupplierTest { final Map> fks = Collections.singletonMap( "fk1", ValueAndTimestamp.make("foo", 1L) ); @@ -375,8 +375,8 @@ public String[] storeNames() { Change>>, String, SubscriptionResponseWrapper> processor(final KTableValueGetterSupplier valueGetterSupplier) { - final SubscriptionJoinForeignProcessorSupplier supplier = - new SubscriptionJoinForeignProcessorSupplier<>(valueGetterSupplier); + final SubscriptionJoinProcessorSupplier supplier = + new SubscriptionJoinProcessorSupplier<>(valueGetterSupplier); return supplier.get(); } } \ No newline at end of file diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResolverJoinProcessorSupplierTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplierTest.java similarity index 90% rename from streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResolverJoinProcessorSupplierTest.java rename to streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplierTest.java index afd5f490db29b..4c26efe236485 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResolverJoinProcessorSupplierTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/ResponseJoinProcessorSupplierTest.java @@ -37,7 +37,7 @@ import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.collection.IsEmptyCollection.empty; -public class SubscriptionResolverJoinProcessorSupplierTest { +public class ResponseJoinProcessorSupplierTest { private static final StringSerializer STRING_SERIALIZER = new StringSerializer(); private static final ValueJoiner JOINER = (value1, value2) -> "(" + value1 + "," + value2 + ")"; @@ -79,8 +79,8 @@ public void shouldNotForwardWhenHashDoesNotMatch() { final TestKTableValueGetterSupplier valueGetterSupplier = new TestKTableValueGetterSupplier<>(); final boolean leftJoin = false; - final SubscriptionResolverJoinProcessorSupplier processorSupplier = - new SubscriptionResolverJoinProcessorSupplier<>( + final ResponseJoinProcessorSupplier processorSupplier = + new ResponseJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, () -> "value-hash-dummy-topic", @@ -104,8 +104,8 @@ public void shouldIgnoreUpdateWhenLeftHasBecomeNull() { final TestKTableValueGetterSupplier valueGetterSupplier = new TestKTableValueGetterSupplier<>(); final boolean leftJoin = false; - final SubscriptionResolverJoinProcessorSupplier processorSupplier = - new SubscriptionResolverJoinProcessorSupplier<>( + final ResponseJoinProcessorSupplier processorSupplier = + new ResponseJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, () -> "value-hash-dummy-topic", @@ -129,8 +129,8 @@ public void shouldForwardWhenHashMatches() { final TestKTableValueGetterSupplier valueGetterSupplier = new TestKTableValueGetterSupplier<>(); final boolean leftJoin = false; - final SubscriptionResolverJoinProcessorSupplier processorSupplier = - new SubscriptionResolverJoinProcessorSupplier<>( + final ResponseJoinProcessorSupplier processorSupplier = + new ResponseJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, () -> "value-hash-dummy-topic", @@ -155,8 +155,8 @@ public void shouldEmitTombstoneForInnerJoinWhenRightIsNull() { final TestKTableValueGetterSupplier valueGetterSupplier = new TestKTableValueGetterSupplier<>(); final boolean leftJoin = false; - final SubscriptionResolverJoinProcessorSupplier processorSupplier = - new SubscriptionResolverJoinProcessorSupplier<>( + final ResponseJoinProcessorSupplier processorSupplier = + new ResponseJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, () -> "value-hash-dummy-topic", @@ -181,8 +181,8 @@ public void shouldEmitResultForLeftJoinWhenRightIsNull() { final TestKTableValueGetterSupplier valueGetterSupplier = new TestKTableValueGetterSupplier<>(); final boolean leftJoin = true; - final SubscriptionResolverJoinProcessorSupplier processorSupplier = - new SubscriptionResolverJoinProcessorSupplier<>( + final ResponseJoinProcessorSupplier processorSupplier = + new ResponseJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, () -> "value-hash-dummy-topic", @@ -207,8 +207,8 @@ public void shouldEmitTombstoneForLeftJoinWhenRightIsNullAndLeftIsNull() { final TestKTableValueGetterSupplier valueGetterSupplier = new TestKTableValueGetterSupplier<>(); final boolean leftJoin = true; - final SubscriptionResolverJoinProcessorSupplier processorSupplier = - new SubscriptionResolverJoinProcessorSupplier<>( + final ResponseJoinProcessorSupplier processorSupplier = + new ResponseJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, () -> "value-hash-dummy-topic", diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionStoreReceiveProcessorSupplierTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionReceiveProcessorSupplierTest.java similarity index 95% rename from streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionStoreReceiveProcessorSupplierTest.java rename to streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionReceiveProcessorSupplierTest.java index 59636556c402b..dc04d560ed0a8 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionStoreReceiveProcessorSupplierTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionReceiveProcessorSupplierTest.java @@ -45,7 +45,7 @@ import org.junit.Before; import org.junit.Test; -public class SubscriptionStoreReceiveProcessorSupplierTest { +public class SubscriptionReceiveProcessorSupplierTest { private final Properties props = StreamsTestUtils.getStreamsConfig(Serdes.String(), Serdes.String()); private File stateDir; @@ -87,7 +87,7 @@ public void shouldDetectVersionChange() { @Test public void shouldDeleteKeyAndPropagateV0() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -139,7 +139,7 @@ public void shouldDeleteKeyAndPropagateV0() { @Test public void shouldDeleteKeyAndPropagateV1() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -190,7 +190,7 @@ public void shouldDeleteKeyAndPropagateV1() { @Test public void shouldDeleteKeyNoPropagateV0() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -242,7 +242,7 @@ public void shouldDeleteKeyNoPropagateV0() { @Test public void shouldDeleteKeyNoPropagateV1() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -294,7 +294,7 @@ public void shouldDeleteKeyNoPropagateV1() { @Test public void shouldPropagateOnlyIfFKValAvailableV0() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -346,7 +346,7 @@ public void shouldPropagateOnlyIfFKValAvailableV0() { @Test public void shouldPropagateOnlyIfFKValAvailableV1() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -398,7 +398,7 @@ public void shouldPropagateOnlyIfFKValAvailableV1() { @Test public void shouldPropagateNullIfNoFKValAvailableV0() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -450,7 +450,7 @@ public void shouldPropagateNullIfNoFKValAvailableV0() { @Test public void shouldPropagateNullIfNoFKValAvailableV1() { final StoreBuilder>> storeBuilder = storeBuilder(); - final SubscriptionStoreReceiveProcessorSupplier supplier = supplier(storeBuilder); + final SubscriptionReceiveProcessorSupplier supplier = supplier(storeBuilder); final Processor, CombinedKey, @@ -500,10 +500,10 @@ public void shouldPropagateNullIfNoFKValAvailableV1() { } - private SubscriptionStoreReceiveProcessorSupplier supplier( + private SubscriptionReceiveProcessorSupplier supplier( final StoreBuilder>> storeBuilder) { - return new SubscriptionStoreReceiveProcessorSupplier<>(storeBuilder, COMBINED_KEY_SCHEMA); + return new SubscriptionReceiveProcessorSupplier<>(storeBuilder, COMBINED_KEY_SCHEMA); } private StoreBuilder>> storeBuilder() {