From 8e8467de1e064eef4e06114757e7d3b1d2331940 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 28 Apr 2020 11:25:53 -0500 Subject: [PATCH 1/3] KAFKA-9925: decorate pseudo-topics with app id --- .../streams/kstream/internals/KTableImpl.java | 25 ++++++++++++---- .../foreignkeyjoin/CombinedKeySchema.java | 17 +++++++---- ...JoinSubscriptionSendProcessorSupplier.java | 17 +++++++---- ...criptionResolverJoinProcessorSupplier.java | 10 +++++-- .../SubscriptionWrapperSerde.java | 30 +++++++++++++------ .../internals/InternalTopologyBuilder.java | 4 +++ ...TableKTableForeignKeyJoinScenarioTest.java | 30 ++++++++++--------- .../foreignkeyjoin/CombinedKeySchemaTest.java | 20 ++++++------- ...tionResolverJoinProcessorSupplierTest.java | 12 ++++---- .../SubscriptionWrapperSerdeTest.java | 10 ++++--- 10 files changed, 112 insertions(+), 63 deletions(-) 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 7407246bf2e15..d6feae8ee2a0c 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 @@ -58,6 +58,7 @@ import org.apache.kafka.streams.kstream.internals.suppress.KTableSuppressProcessorSupplier; import org.apache.kafka.streams.kstream.internals.suppress.NamedSuppressed; import org.apache.kafka.streams.kstream.internals.suppress.SuppressedInternal; +import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.processor.ProcessorSupplier; import org.apache.kafka.streams.processor.internals.InternalTopicProperties; import org.apache.kafka.streams.processor.internals.StaticTopicNameExtractor; @@ -77,6 +78,7 @@ import java.util.Objects; import java.util.Set; import java.util.function.Function; +import java.util.function.Supplier; import static org.apache.kafka.streams.kstream.internals.graph.GraphGraceSearchUtil.findAndVerifyWindowGrace; @@ -972,13 +974,26 @@ private KTable doJoinOnForeignKey(final KTable forei //This occurs whenever the extracted foreignKey changes values. enableSendingOldValues(); + final NamedInternal renamed = new NamedInternal(joinName); + + final String subscriptionTopicName = renamed.suffixWithOrElseGet( + "-subscription-registration", + builder, + SUBSCRIPTION_REGISTRATION + ) + TOPIC_SUFFIX; + // the decoration can't be performed until we have the configuration available when the app runs, + // so we pass Suppliers into the components, which they can call at run time + + final Supplier subscriptionPrimaryKeySerdePseudoTopic = + () -> internalTopologyBuilder().decoratePseudoTopic(subscriptionTopicName + "-pk"); + + final Supplier subscriptionForeignKeySerdePseudoTopic = + () -> internalTopologyBuilder().decoratePseudoTopic(subscriptionTopicName + "-fk"); + + final Supplier valueHashSerdePseudoTopic = + () -> internalTopologyBuilder().decoratePseudoTopic(subscriptionTopicName + "-vh"); - final NamedInternal renamed = new NamedInternal(joinName); - final String subscriptionTopicName = renamed.suffixWithOrElseGet("-subscription-registration", builder, SUBSCRIPTION_REGISTRATION) + TOPIC_SUFFIX; - final String subscriptionPrimaryKeySerdePseudoTopic = subscriptionTopicName + "-pk"; - final String subscriptionForeignKeySerdePseudoTopic = subscriptionTopicName + "-fk"; - final String valueHashSerdePseudoTopic = subscriptionTopicName + "-vh"; builder.internalTopologyBuilder.addInternalTopic(subscriptionTopicName, InternalTopicProperties.empty()); final Serde foreignKeySerde = ((KTableImpl) foreignKeyTable).keySerde; diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchema.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchema.java index 92fb72c7ae6b9..57bc646a13e11 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchema.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchema.java @@ -23,24 +23,27 @@ import org.apache.kafka.streams.processor.ProcessorContext; import java.nio.ByteBuffer; +import java.util.function.Supplier; /** * Factory for creating CombinedKey serializers / deserializers. */ public class CombinedKeySchema { - private final String primaryKeySerdeTopic; - private final String foreignKeySerdeTopic; + private final Supplier undecoratedPrimaryKeySerdeTopicSupplier; + private final Supplier undecoratedForeignKeySerdeTopicSupplier; + private String primaryKeySerdeTopic; + private String foreignKeySerdeTopic; private Serializer primaryKeySerializer; private Deserializer primaryKeyDeserializer; private Serializer foreignKeySerializer; private Deserializer foreignKeyDeserializer; - public CombinedKeySchema(final String foreignKeySerdeTopic, + public CombinedKeySchema(final Supplier foreignKeySerdeTopicSupplier, final Serde foreignKeySerde, - final String primaryKeySerdeTopic, + final Supplier primaryKeySerdeTopicSupplier, final Serde primaryKeySerde) { - this.primaryKeySerdeTopic = primaryKeySerdeTopic; - this.foreignKeySerdeTopic = foreignKeySerdeTopic; + undecoratedPrimaryKeySerdeTopicSupplier = primaryKeySerdeTopicSupplier; + undecoratedForeignKeySerdeTopicSupplier = foreignKeySerdeTopicSupplier; primaryKeySerializer = primaryKeySerde == null ? null : primaryKeySerde.serializer(); primaryKeyDeserializer = primaryKeySerde == null ? null : primaryKeySerde.deserializer(); foreignKeyDeserializer = foreignKeySerde == null ? null : foreignKeySerde.deserializer(); @@ -49,6 +52,8 @@ public CombinedKeySchema(final String foreignKeySerdeTopic, @SuppressWarnings("unchecked") public void init(final ProcessorContext context) { + primaryKeySerdeTopic = undecoratedPrimaryKeySerdeTopicSupplier.get(); + foreignKeySerdeTopic = undecoratedForeignKeySerdeTopicSupplier.get(); primaryKeySerializer = primaryKeySerializer == null ? (Serializer) context.keySerde().serializer() : primaryKeySerializer; primaryKeyDeserializer = primaryKeyDeserializer == null ? (Deserializer) context.keySerde().deserializer() : primaryKeyDeserializer; foreignKeySerializer = foreignKeySerializer == null ? (Serializer) context.keySerde().serializer() : foreignKeySerializer; 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/ForeignJoinSubscriptionSendProcessorSupplier.java index ba794f7a972ff..97878750dc1f6 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/ForeignJoinSubscriptionSendProcessorSupplier.java @@ -33,6 +33,7 @@ import java.util.Arrays; import java.util.function.Function; +import java.util.function.Supplier; import static org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapper.Instruction.DELETE_KEY_AND_PROPAGATE; import static org.apache.kafka.streams.kstream.internals.foreignkeyjoin.SubscriptionWrapper.Instruction.DELETE_KEY_NO_PROPAGATE; @@ -43,21 +44,21 @@ public class ForeignJoinSubscriptionSendProcessorSupplier implements P private static final Logger LOG = LoggerFactory.getLogger(ForeignJoinSubscriptionSendProcessorSupplier.class); private final Function foreignKeyExtractor; - private final String foreignKeySerdeTopic; - private final String valueSerdeTopic; + private final Supplier foreignKeySerdeTopicSupplier; + private final Supplier valueSerdeTopicSupplier; private final boolean leftJoin; private Serializer foreignKeySerializer; private Serializer valueSerializer; public ForeignJoinSubscriptionSendProcessorSupplier(final Function foreignKeyExtractor, - final String foreignKeySerdeTopic, - final String valueSerdeTopic, + final Supplier foreignKeySerdeTopicSupplier, + final Supplier valueSerdeTopicSupplier, final Serde foreignKeySerde, final Serializer valueSerializer, final boolean leftJoin) { this.foreignKeyExtractor = foreignKeyExtractor; - this.foreignKeySerdeTopic = foreignKeySerdeTopic; - this.valueSerdeTopic = valueSerdeTopic; + this.foreignKeySerdeTopicSupplier = foreignKeySerdeTopicSupplier; + this.valueSerdeTopicSupplier = valueSerdeTopicSupplier; this.valueSerializer = valueSerializer; this.leftJoin = leftJoin; foreignKeySerializer = foreignKeySerde == null ? null : foreignKeySerde.serializer(); @@ -71,11 +72,15 @@ public Processor> get() { private class UnbindChangeProcessor extends AbstractProcessor> { private Sensor droppedRecordsSensor; + private String foreignKeySerdeTopic; + private String valueSerdeTopic; @SuppressWarnings("unchecked") @Override public void init(final ProcessorContext context) { super.init(context); + foreignKeySerdeTopic = foreignKeySerdeTopicSupplier.get(); + valueSerdeTopic = valueSerdeTopicSupplier.get(); // get default key serde if it wasn't supplied directly at construction if (foreignKeySerializer == null) { foreignKeySerializer = (Serializer) context.keySerde().serializer(); 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/SubscriptionResolverJoinProcessorSupplier.java index 31de0687d6ef7..3cd06368a7131 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/SubscriptionResolverJoinProcessorSupplier.java @@ -29,6 +29,8 @@ import org.apache.kafka.streams.state.ValueAndTimestamp; import org.apache.kafka.streams.state.internals.Murmur3; +import java.util.function.Supplier; + /** * Receives {@code SubscriptionResponseWrapper} events and filters out events which do not match the current hash * of the primary key. This eliminates race-condition results for rapidly-changing foreign-keys for a given primary key. @@ -42,18 +44,18 @@ public class SubscriptionResolverJoinProcessorSupplier implements ProcessorSupplier> { private final KTableValueGetterSupplier valueGetterSupplier; private final Serializer constructionTimeValueSerializer; - private final String valueHashSerdePseudoTopic; + private final Supplier valueHashSerdePseudoTopicSupplier; private final ValueJoiner joiner; private final boolean leftJoin; public SubscriptionResolverJoinProcessorSupplier(final KTableValueGetterSupplier valueGetterSupplier, final Serializer valueSerializer, - final String valueHashSerdePseudoTopic, + final Supplier valueHashSerdePseudoTopicSupplier, final ValueJoiner joiner, final boolean leftJoin) { this.valueGetterSupplier = valueGetterSupplier; constructionTimeValueSerializer = valueSerializer; - this.valueHashSerdePseudoTopic = valueHashSerdePseudoTopic; + this.valueHashSerdePseudoTopicSupplier = valueHashSerdePseudoTopicSupplier; this.joiner = joiner; this.leftJoin = leftJoin; } @@ -61,6 +63,7 @@ public SubscriptionResolverJoinProcessorSupplier(final KTableValueGetterSupplier @Override public Processor> get() { return new AbstractProcessor>() { + private String valueHashSerdePseudoTopic; private Serializer runtimeValueSerializer = constructionTimeValueSerializer; private KTableValueGetter valueGetter; @@ -69,6 +72,7 @@ public Processor> get() { @Override public void init(final ProcessorContext context) { super.init(context); + valueHashSerdePseudoTopic = valueHashSerdePseudoTopicSupplier.get(); valueGetter = valueGetterSupplier.get(); valueGetter.init(context); if (runtimeValueSerializer == null) { diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerde.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerde.java index 42aed940ef4bc..136128c53bbee 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerde.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerde.java @@ -25,16 +25,17 @@ import java.nio.ByteBuffer; import java.util.Objects; +import java.util.function.Supplier; public class SubscriptionWrapperSerde implements Serde> { private final SubscriptionWrapperSerializer serializer; private final SubscriptionWrapperDeserializer deserializer; - public SubscriptionWrapperSerde(final String primaryKeySerializationPseudoTopic, + public SubscriptionWrapperSerde(final Supplier primaryKeySerializationPseudoTopicSupplier, final Serde primaryKeySerde) { - serializer = new SubscriptionWrapperSerializer<>(primaryKeySerializationPseudoTopic, + serializer = new SubscriptionWrapperSerializer<>(primaryKeySerializationPseudoTopicSupplier, primaryKeySerde == null ? null : primaryKeySerde.serializer()); - deserializer = new SubscriptionWrapperDeserializer<>(primaryKeySerializationPseudoTopic, + deserializer = new SubscriptionWrapperDeserializer<>(primaryKeySerializationPseudoTopicSupplier, primaryKeySerde == null ? null : primaryKeySerde.deserializer()); } @@ -51,12 +52,13 @@ public Deserializer> deserializer() { private static class SubscriptionWrapperSerializer implements Serializer>, WrappingNullableSerializer, K> { - private final String primaryKeySerializationPseudoTopic; + private final Supplier primaryKeySerializationPseudoTopicSupplier; + private String primaryKeySerializationPseudoTopic = null; private Serializer primaryKeySerializer; - SubscriptionWrapperSerializer(final String primaryKeySerializationPseudoTopic, + SubscriptionWrapperSerializer(final Supplier primaryKeySerializationPseudoTopicSupplier, final Serializer primaryKeySerializer) { - this.primaryKeySerializationPseudoTopic = primaryKeySerializationPseudoTopic; + this.primaryKeySerializationPseudoTopicSupplier = primaryKeySerializationPseudoTopicSupplier; this.primaryKeySerializer = primaryKeySerializer; } @@ -76,6 +78,10 @@ public byte[] serialize(final String ignored, final SubscriptionWrapper data) throw new UnsupportedVersionException("SubscriptionWrapper version is larger than maximum supported 0x7F"); } + if (primaryKeySerializationPseudoTopic == null) { + primaryKeySerializationPseudoTopic = primaryKeySerializationPseudoTopicSupplier.get(); + } + final byte[] primaryKeySerializedData = primaryKeySerializer.serialize( primaryKeySerializationPseudoTopic, data.getPrimaryKey() @@ -106,12 +112,13 @@ public byte[] serialize(final String ignored, final SubscriptionWrapper data) private static class SubscriptionWrapperDeserializer implements Deserializer>, WrappingNullableDeserializer, K> { - private final String primaryKeySerializationPseudoTopic; + private final Supplier primaryKeySerializationPseudoTopicSupplier; + private String primaryKeySerializationPseudoTopic = null; private Deserializer primaryKeyDeserializer; - SubscriptionWrapperDeserializer(final String primaryKeySerializationPseudoTopic, + SubscriptionWrapperDeserializer(final Supplier primaryKeySerializationPseudoTopicSupplier, final Deserializer primaryKeyDeserializer) { - this.primaryKeySerializationPseudoTopic = primaryKeySerializationPseudoTopic; + this.primaryKeySerializationPseudoTopicSupplier = primaryKeySerializationPseudoTopicSupplier; this.primaryKeyDeserializer = primaryKeyDeserializer; } @@ -144,6 +151,11 @@ public SubscriptionWrapper deserialize(final String ignored, final byte[] dat final byte[] primaryKeyRaw = new byte[data.length - lengthSum]; //The remaining data is the serialized pk buf.get(primaryKeyRaw, 0, primaryKeyRaw.length); + + if (primaryKeySerializationPseudoTopic == null) { + primaryKeySerializationPseudoTopic = primaryKeySerializationPseudoTopicSupplier.get(); + } + final K primaryKey = primaryKeyDeserializer.deserialize(primaryKeySerializationPseudoTopic, primaryKeyRaw); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java index 854ce1cb81714..771ca07dd8bef 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java @@ -1224,6 +1224,10 @@ private List maybeDecorateInternalSourceTopics(final Collection return decoratedTopics; } + public String decoratePseudoTopic(final String topic) { + return decorateTopic(topic); + } + private String decorateTopic(final String topic) { if (applicationId == null) { throw new TopologyException("there are internal topics and " diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java index 84d3552fce648..15513b6f36a4a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java @@ -41,6 +41,7 @@ import java.util.Collections; import java.util.Map; import java.util.Properties; +import java.util.TreeSet; import static org.apache.kafka.common.utils.Utils.mkEntry; import static org.apache.kafka.common.utils.Utils.mkMap; @@ -179,8 +180,9 @@ public void shouldWorkWithDefaultAndProducedSerdes() { @Test public void shouldUseExpectedTopicsWithSerde() { + final String applicationId = "ktable-ktable-joinOnForeignKey"; final Properties streamsConfig = mkProperties(mkMap( - mkEntry(StreamsConfig.APPLICATION_ID_CONFIG, "ktable-ktable-joinOnForeignKey"), + mkEntry(StreamsConfig.APPLICATION_ID_CONFIG, applicationId), mkEntry(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "asdf:0000"), mkEntry(StreamsConfig.STATE_DIR_CONFIG, TestUtils.tempDirectory().getPath()) )); @@ -191,12 +193,12 @@ public void shouldUseExpectedTopicsWithSerde() { final KTable left = builder.table( LEFT_TABLE, Consumed.with(serdeScope.decorateSerde(Serdes.String(), streamsConfig, true), - serdeScope.decorateSerde(Serdes.String(), streamsConfig, false)) + serdeScope.decorateSerde(Serdes.String(), streamsConfig, false)) ); final KTable right = builder.table( RIGHT_TABLE, Consumed.with(serdeScope.decorateSerde(Serdes.String(), streamsConfig, true), - serdeScope.decorateSerde(Serdes.String(), streamsConfig, false)) + serdeScope.decorateSerde(Serdes.String(), streamsConfig, false)) ); left.join( @@ -218,19 +220,19 @@ public void shouldUseExpectedTopicsWithSerde() { } // verifying primarily that no extra pseudo-topics were used, but it's nice to also verify the rest of the // topics our serdes serialize data for - assertThat(serdeScope.registeredTopics(), CoreMatchers.is(mkSet( + assertThat(serdeScope.registeredTopics(), is(mkSet( // expected pseudo-topics - "KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic-fk--key", - "KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic-pk--key", - "KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic-vh--value", + applicationId + "-KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic-fk--key", + applicationId + "-KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic-pk--key", + applicationId + "-KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic-vh--value", // internal topics - "ktable-ktable-joinOnForeignKey-KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic--key", - "ktable-ktable-joinOnForeignKey-KTABLE-FK-JOIN-SUBSCRIPTION-RESPONSE-0000000014-topic--key", - "ktable-ktable-joinOnForeignKey-KTABLE-FK-JOIN-SUBSCRIPTION-RESPONSE-0000000014-topic--value", - "ktable-ktable-joinOnForeignKey-left_table-STATE-STORE-0000000000-changelog--key", - "ktable-ktable-joinOnForeignKey-left_table-STATE-STORE-0000000000-changelog--value", - "ktable-ktable-joinOnForeignKey-right_table-STATE-STORE-0000000003-changelog--key", - "ktable-ktable-joinOnForeignKey-right_table-STATE-STORE-0000000003-changelog--value", + applicationId + "-KTABLE-FK-JOIN-SUBSCRIPTION-REGISTRATION-0000000006-topic--key", + applicationId + "-KTABLE-FK-JOIN-SUBSCRIPTION-RESPONSE-0000000014-topic--key", + applicationId + "-KTABLE-FK-JOIN-SUBSCRIPTION-RESPONSE-0000000014-topic--value", + applicationId + "-left_table-STATE-STORE-0000000000-changelog--key", + applicationId + "-left_table-STATE-STORE-0000000000-changelog--value", + applicationId + "-right_table-STATE-STORE-0000000003-changelog--key", + applicationId + "-right_table-STATE-STORE-0000000003-changelog--value", // output topics "output-topic--key", "output-topic--value" diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java index cb1ef596db3f0..487291431ed9c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java @@ -28,8 +28,8 @@ public class CombinedKeySchemaTest { @Test public void nonNullPrimaryKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>("fkTopic", Serdes.String(), - "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer()); final Integer primary = -999; final Bytes result = cks.toBytes("foreignKey", primary); @@ -40,22 +40,22 @@ public void nonNullPrimaryKeySerdeTest() { @Test(expected = NullPointerException.class) public void nullPrimaryKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>("fkTopic", Serdes.String(), - "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer()); cks.toBytes("foreignKey", null); } @Test(expected = NullPointerException.class) public void nullForeignKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>("fkTopic", Serdes.String(), - "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer()); cks.toBytes(null, 10); } @Test public void prefixKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>("fkTopic", Serdes.String(), - "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer()); final String foreignKey = "someForeignKey"; final byte[] foreignKeySerializedData = Serdes.String().serializer().serialize("fkTopic", foreignKey); @@ -71,8 +71,8 @@ public void prefixKeySerdeTest() { @Test(expected = NullPointerException.class) public void nullPrefixKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>("fkTopic", Serdes.String(), - "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer()); final String foreignKey = null; cks.prefixBytes(foreignKey); } 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/SubscriptionResolverJoinProcessorSupplierTest.java index aae99ec6687d7..b95569f7906c8 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/SubscriptionResolverJoinProcessorSupplierTest.java @@ -82,7 +82,7 @@ public void shouldNotForwardWhenHashDoesNotMatch() { new SubscriptionResolverJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, - "value-hash-dummy-topic", + () -> "value-hash-dummy-topic", JOINER, leftJoin ); @@ -107,7 +107,7 @@ public void shouldIgnoreUpdateWhenLeftHasBecomeNull() { new SubscriptionResolverJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, - "value-hash-dummy-topic", + () -> "value-hash-dummy-topic", JOINER, leftJoin ); @@ -132,7 +132,7 @@ public void shouldForwardWhenHashMatches() { new SubscriptionResolverJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, - "value-hash-dummy-topic", + () -> "value-hash-dummy-topic", JOINER, leftJoin ); @@ -158,7 +158,7 @@ public void shouldEmitTombstoneForInnerJoinWhenRightIsNull() { new SubscriptionResolverJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, - "value-hash-dummy-topic", + () -> "value-hash-dummy-topic", JOINER, leftJoin ); @@ -184,7 +184,7 @@ public void shouldEmitResultForLeftJoinWhenRightIsNull() { new SubscriptionResolverJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, - "value-hash-dummy-topic", + () -> "value-hash-dummy-topic", JOINER, leftJoin ); @@ -210,7 +210,7 @@ public void shouldEmitTombstoneForLeftJoinWhenRightIsNullAndLeftIsNull() { new SubscriptionResolverJoinProcessorSupplier<>( valueGetterSupplier, STRING_SERIALIZER, - "value-hash-dummy-topic", + () -> "value-hash-dummy-topic", JOINER, leftJoin ); diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java index 5c6551c3e7db5..30c80b0693444 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java @@ -21,6 +21,8 @@ import org.apache.kafka.streams.state.internals.Murmur3; import org.junit.Test; +import java.util.Collections; + import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; @@ -31,7 +33,7 @@ public class SubscriptionWrapperSerdeTest { @SuppressWarnings("unchecked") public void shouldSerdeTest() { final String originalKey = "originalKey"; - final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>("pkTopic", Serdes.String()); + final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>(() -> "pkTopic", Serdes.String()); final long[] hashedValue = Murmur3.hash128(new byte[] {(byte) 0xFF, (byte) 0xAA, (byte) 0x00, (byte) 0x19}); final SubscriptionWrapper wrapper = new SubscriptionWrapper<>(hashedValue, SubscriptionWrapper.Instruction.DELETE_KEY_AND_PROPAGATE, originalKey); final byte[] serialized = swSerde.serializer().serialize(null, wrapper); @@ -46,7 +48,7 @@ public void shouldSerdeTest() { @SuppressWarnings("unchecked") public void shouldSerdeNullHashTest() { final String originalKey = "originalKey"; - final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>("pkTopic", Serdes.String()); + final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>(() -> "pkTopic", Serdes.String()); final long[] hashedValue = null; final SubscriptionWrapper wrapper = new SubscriptionWrapper<>(hashedValue, SubscriptionWrapper.Instruction.PROPAGATE_ONLY_IF_FK_VAL_AVAILABLE, originalKey); final byte[] serialized = swSerde.serializer().serialize(null, wrapper); @@ -61,7 +63,7 @@ public void shouldSerdeNullHashTest() { @SuppressWarnings("unchecked") public void shouldThrowExceptionOnNullKeyTest() { final String originalKey = null; - final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>("pkTopic", Serdes.String()); + final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>(() -> "pkTopic", Serdes.String()); final long[] hashedValue = Murmur3.hash128(new byte[] {(byte) 0xFF, (byte) 0xAA, (byte) 0x00, (byte) 0x19}); final SubscriptionWrapper wrapper = new SubscriptionWrapper<>(hashedValue, SubscriptionWrapper.Instruction.PROPAGATE_ONLY_IF_FK_VAL_AVAILABLE, originalKey); swSerde.serializer().serialize(null, wrapper); @@ -71,7 +73,7 @@ public void shouldThrowExceptionOnNullKeyTest() { @SuppressWarnings("unchecked") public void shouldThrowExceptionOnNullInstructionTest() { final String originalKey = "originalKey"; - final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>("pkTopic", Serdes.String()); + final SubscriptionWrapperSerde swSerde = new SubscriptionWrapperSerde<>(() -> "pkTopic", Serdes.String()); final long[] hashedValue = Murmur3.hash128(new byte[] {(byte) 0xFF, (byte) 0xAA, (byte) 0x00, (byte) 0x19}); final SubscriptionWrapper wrapper = new SubscriptionWrapper<>(hashedValue, null, originalKey); swSerde.serializer().serialize(null, wrapper); From a863182d36994a2773be69953159c8a5320b07cc Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 28 Apr 2020 11:51:19 -0500 Subject: [PATCH 2/3] checkstyle --- ...TableKTableForeignKeyJoinScenarioTest.java | 2 -- .../foreignkeyjoin/CombinedKeySchemaTest.java | 30 ++++++++++++------- .../SubscriptionWrapperSerdeTest.java | 2 -- 3 files changed, 20 insertions(+), 14 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java index 15513b6f36a4a..1a43a2f513971 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableForeignKeyJoinScenarioTest.java @@ -33,7 +33,6 @@ import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.utils.UniqueTopicSerdeScope; import org.apache.kafka.test.TestUtils; -import org.hamcrest.CoreMatchers; import org.junit.Rule; import org.junit.Test; import org.junit.rules.TestName; @@ -41,7 +40,6 @@ import java.util.Collections; import java.util.Map; import java.util.Properties; -import java.util.TreeSet; import static org.apache.kafka.common.utils.Utils.mkEntry; import static org.apache.kafka.common.utils.Utils.mkMap; diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java index 487291431ed9c..17f0c7969b4a8 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/CombinedKeySchemaTest.java @@ -28,8 +28,10 @@ public class CombinedKeySchemaTest { @Test public void nonNullPrimaryKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), - () -> "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>( + () -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer() + ); final Integer primary = -999; final Bytes result = cks.toBytes("foreignKey", primary); @@ -40,22 +42,28 @@ public void nonNullPrimaryKeySerdeTest() { @Test(expected = NullPointerException.class) public void nullPrimaryKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), - () -> "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>( + () -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer() + ); cks.toBytes("foreignKey", null); } @Test(expected = NullPointerException.class) public void nullForeignKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), - () -> "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>( + () -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer() + ); cks.toBytes(null, 10); } @Test public void prefixKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), - () -> "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>( + () -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer() + ); final String foreignKey = "someForeignKey"; final byte[] foreignKeySerializedData = Serdes.String().serializer().serialize("fkTopic", foreignKey); @@ -71,8 +79,10 @@ public void prefixKeySerdeTest() { @Test(expected = NullPointerException.class) public void nullPrefixKeySerdeTest() { - final CombinedKeySchema cks = new CombinedKeySchema<>(() -> "fkTopic", Serdes.String(), - () -> "pkTopic", Serdes.Integer()); + final CombinedKeySchema cks = new CombinedKeySchema<>( + () -> "fkTopic", Serdes.String(), + () -> "pkTopic", Serdes.Integer() + ); final String foreignKey = null; cks.prefixBytes(foreignKey); } diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java index 30c80b0693444..dd67b4bd55594 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionWrapperSerdeTest.java @@ -21,8 +21,6 @@ import org.apache.kafka.streams.state.internals.Murmur3; import org.junit.Test; -import java.util.Collections; - import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; From ea21699024ea74196d96e416751ba49b77ca074d Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 28 Apr 2020 11:53:22 -0500 Subject: [PATCH 3/3] checkstyle --- .../org/apache/kafka/streams/kstream/internals/KTableImpl.java | 1 - 1 file changed, 1 deletion(-) 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 d6feae8ee2a0c..360656d9c640e 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 @@ -58,7 +58,6 @@ import org.apache.kafka.streams.kstream.internals.suppress.KTableSuppressProcessorSupplier; import org.apache.kafka.streams.kstream.internals.suppress.NamedSuppressed; import org.apache.kafka.streams.kstream.internals.suppress.SuppressedInternal; -import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.processor.ProcessorSupplier; import org.apache.kafka.streams.processor.internals.InternalTopicProperties; import org.apache.kafka.streams.processor.internals.StaticTopicNameExtractor;