diff --git a/checkstyle/import-control.xml b/checkstyle/import-control.xml
index 71192655f4f17..8c3cb2865c63d 100644
--- a/checkstyle/import-control.xml
+++ b/checkstyle/import-control.xml
@@ -261,6 +261,7 @@
+
diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedDeserializer.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedDeserializer.java
index 90d5882887f3f..433a18ddc365b 100644
--- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedDeserializer.java
+++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedDeserializer.java
@@ -22,7 +22,7 @@
import java.nio.ByteBuffer;
import java.util.Objects;
-public class ChangedDeserializer implements Deserializer>, WrappingNullableDeserializer, T> {
+public class ChangedDeserializer implements Deserializer>, WrappingNullableDeserializer, Void, T> {
private static final int NEWFLAG_SIZE = 1;
@@ -37,9 +37,9 @@ public Deserializer inner() {
}
@Override
- public void setIfUnset(final Deserializer defaultDeserializer) {
+ public void setIfUnset(final Deserializer defaultKeyDeserializer, final Deserializer defaultValueDeserializer) {
if (inner == null) {
- inner = Objects.requireNonNull(defaultDeserializer, "defaultDeserializer cannot be null");
+ inner = Objects.requireNonNull(defaultValueDeserializer);
}
}
diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedSerializer.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedSerializer.java
index 551d948392266..f5d63cdaf7ba8 100644
--- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedSerializer.java
+++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/ChangedSerializer.java
@@ -23,7 +23,7 @@
import java.nio.ByteBuffer;
import java.util.Objects;
-public class ChangedSerializer implements Serializer>, WrappingNullableSerializer, T> {
+public class ChangedSerializer implements Serializer>, WrappingNullableSerializer, Void, T> {
private static final int NEWFLAG_SIZE = 1;
@@ -38,9 +38,9 @@ public Serializer inner() {
}
@Override
- public void setIfUnset(final Serializer defaultSerializer) {
+ public void setIfUnset(final Serializer defaultKeySerializer, final Serializer defaultValueSerializer) {
if (inner == null) {
- inner = Objects.requireNonNull(defaultSerializer, "defaultSerializer cannot be null");
+ inner = Objects.requireNonNull(defaultValueSerializer);
}
}
diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableDeserializer.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableDeserializer.java
index a57e9a15315c4..d0c0b14c2883a 100644
--- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableDeserializer.java
+++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableDeserializer.java
@@ -18,6 +18,6 @@
import org.apache.kafka.common.serialization.Deserializer;
-public interface WrappingNullableDeserializer extends Deserializer {
- void setIfUnset(final Deserializer defaultDeserializer);
-}
+public interface WrappingNullableDeserializer extends Deserializer {
+ void setIfUnset(final Deserializer defaultKeyDeserializer, final Deserializer defaultValueDeserializer);
+}
\ No newline at end of file
diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableSerializer.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableSerializer.java
index 2d28e52db2bbd..8854a8d9009c5 100644
--- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableSerializer.java
+++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/WrappingNullableSerializer.java
@@ -18,6 +18,6 @@
import org.apache.kafka.common.serialization.Serializer;
-public interface WrappingNullableSerializer extends Serializer {
- void setIfUnset(final Serializer defaultSerializer);
+public interface WrappingNullableSerializer extends Serializer {
+ void setIfUnset(final Serializer defaultKeySerializer, final Serializer defaultValueSerializer);
}
diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerde.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerde.java
index 31317c500df6c..86191115f1928 100644
--- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerde.java
+++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/foreignkeyjoin/SubscriptionResponseWrapperSerde.java
@@ -46,7 +46,7 @@ public Deserializer> deserializer() {
}
private static final class SubscriptionResponseWrapperSerializer
- implements Serializer>, WrappingNullableSerializer, V> {
+ implements Serializer>, WrappingNullableSerializer, Void, V> {
private Serializer serializer;
@@ -55,9 +55,9 @@ private SubscriptionResponseWrapperSerializer(final Serializer serializer) {
}
@Override
- public void setIfUnset(final Serializer defaultSerializer) {
+ public void setIfUnset(final Serializer defaultKeySerializer, final Serializer defaultValueSerializer) {
if (serializer == null) {
- serializer = Objects.requireNonNull(defaultSerializer, "defaultSerializer cannot be null");
+ serializer = Objects.requireNonNull(defaultValueSerializer);
}
}
@@ -94,7 +94,7 @@ public byte[] serialize(final String topic, final SubscriptionResponseWrapper
}
private static final class SubscriptionResponseWrapperDeserializer
- implements Deserializer>, WrappingNullableDeserializer, V> {
+ implements Deserializer>, WrappingNullableDeserializer, Void, V> {
private Deserializer deserializer;
@@ -103,9 +103,9 @@ private SubscriptionResponseWrapperDeserializer(final Deserializer deserializ
}
@Override
- public void setIfUnset(final Deserializer defaultDeserializer) {
+ public void setIfUnset(final Deserializer defaultKeyDeserializer, final Deserializer defaultValueDeserializer) {
if (deserializer == null) {
- deserializer = Objects.requireNonNull(defaultDeserializer, "defaultDeserializer cannot be null");
+ deserializer = Objects.requireNonNull(defaultValueDeserializer);
}
}
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 136128c53bbee..d2cc989b99e79 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
@@ -50,7 +50,7 @@ public Deserializer> deserializer() {
}
private static class SubscriptionWrapperSerializer
- implements Serializer>, WrappingNullableSerializer, K> {
+ implements Serializer>, WrappingNullableSerializer, K, Void> {
private final Supplier primaryKeySerializationPseudoTopicSupplier;
private String primaryKeySerializationPseudoTopic = null;
@@ -63,9 +63,9 @@ private static class SubscriptionWrapperSerializer
}
@Override
- public void setIfUnset(final Serializer defaultSerializer) {
+ public void setIfUnset(final Serializer defaultKeySerializer, final Serializer defaultValueSerializer) {
if (primaryKeySerializer == null) {
- primaryKeySerializer = Objects.requireNonNull(defaultSerializer, "defaultSerializer cannot be null");
+ primaryKeySerializer = Objects.requireNonNull(defaultKeySerializer);
}
}
@@ -110,7 +110,7 @@ public byte[] serialize(final String ignored, final SubscriptionWrapper data)
}
private static class SubscriptionWrapperDeserializer
- implements Deserializer>, WrappingNullableDeserializer, K> {
+ implements Deserializer>, WrappingNullableDeserializer, K, Void> {
private final Supplier primaryKeySerializationPseudoTopicSupplier;
private String primaryKeySerializationPseudoTopic = null;
@@ -123,9 +123,9 @@ private static class SubscriptionWrapperDeserializer
}
@Override
- public void setIfUnset(final Deserializer defaultDeserializer) {
+ public void setIfUnset(final Deserializer defaultKeyDeserializer, final Deserializer defaultValueDeserializer) {
if (primaryKeyDeserializer == null) {
- primaryKeyDeserializer = Objects.requireNonNull(defaultDeserializer, "defaultDeserializer cannot be null");
+ primaryKeyDeserializer = Objects.requireNonNull(defaultKeyDeserializer);
}
}
diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/SinkNode.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/SinkNode.java
index e0f2510de9364..9b0a254b7ea8a 100644
--- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/SinkNode.java
+++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/SinkNode.java
@@ -66,10 +66,13 @@ public void init(final InternalProcessorContext context) {
valSerializer = (Serializer) context.valueSerde().serializer();
}
- // if value serializers are internal wrapping serializers that may need to be given the default serializer
+ // if serializers are internal wrapping serializers that may need to be given the default serializer
// then pass it the default one from the context
if (valSerializer instanceof WrappingNullableSerializer) {
- ((WrappingNullableSerializer) valSerializer).setIfUnset(context.valueSerde().serializer());
+ ((WrappingNullableSerializer) valSerializer).setIfUnset(
+ context.keySerde().serializer(),
+ context.valueSerde().serializer()
+ );
}
}
diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/SourceNode.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/SourceNode.java
index 717495ec640a0..caf30ce21b3ad 100644
--- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/SourceNode.java
+++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/SourceNode.java
@@ -87,10 +87,13 @@ public void init(final InternalProcessorContext context) {
this.valDeserializer = (Deserializer) context.valueSerde().deserializer();
}
- // if value deserializers are internal wrapping deserializers that may need to be given the default
+ // if deserializers are internal wrapping deserializers that may need to be given the default
// then pass it the default one from the context
if (valDeserializer instanceof WrappingNullableDeserializer) {
- ((WrappingNullableDeserializer) valDeserializer).setIfUnset(context.valueSerde().deserializer());
+ ((WrappingNullableDeserializer) valDeserializer).setIfUnset(
+ context.keySerde().deserializer(),
+ context.valueSerde().deserializer()
+ );
}
}
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 ab84e053cb8fb..eb5a4cd935685 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
@@ -16,6 +16,8 @@
*/
package org.apache.kafka.streams.kstream.internals;
+import org.apache.kafka.common.serialization.IntegerDeserializer;
+import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
@@ -61,17 +63,17 @@ public class KTableKTableForeignKeyJoinScenarioTest {
@Test
public void shouldWorkWithDefaultSerdes() {
final StreamsBuilder builder = new StreamsBuilder();
- final KTable aTable = builder.table("A");
- final KTable bTable = builder.table("B");
+ final KTable aTable = builder.table("A");
+ final KTable bTable = builder.table("B");
- final KTable fkJoinResult = aTable.join(
+ final KTable fkJoinResult = aTable.join(
bTable,
- value -> value.split("-")[0],
+ value -> Integer.parseInt(value.split("-")[0]),
(aVal, bVal) -> "(" + aVal + "," + bVal + ")",
Materialized.as("asdf")
);
- final KTable finalJoinResult = aTable.join(
+ final KTable finalJoinResult = aTable.join(
fkJoinResult,
(aVal, fkJoinVal) -> "(" + aVal + "," + fkJoinVal + ")"
);
@@ -84,17 +86,17 @@ public void shouldWorkWithDefaultSerdes() {
@Test
public void shouldWorkWithDefaultAndConsumedSerdes() {
final StreamsBuilder builder = new StreamsBuilder();
- final KTable aTable = builder.table("A", Consumed.with(Serdes.String(), Serdes.String()));
- final KTable bTable = builder.table("B");
+ final KTable aTable = builder.table("A", Consumed.with(Serdes.Integer(), Serdes.String()));
+ final KTable bTable = builder.table("B");
- final KTable fkJoinResult = aTable.join(
+ final KTable fkJoinResult = aTable.join(
bTable,
- value -> value.split("-")[0],
+ value -> Integer.parseInt(value.split("-")[0]),
(aVal, bVal) -> "(" + aVal + "," + bVal + ")",
Materialized.as("asdf")
);
- final KTable finalJoinResult = aTable.join(
+ final KTable finalJoinResult = aTable.join(
fkJoinResult,
(aVal, fkJoinVal) -> "(" + aVal + "," + fkJoinVal + ")"
);
@@ -107,20 +109,19 @@ public void shouldWorkWithDefaultAndConsumedSerdes() {
@Test
public void shouldWorkWithDefaultAndJoinResultSerdes() {
final StreamsBuilder builder = new StreamsBuilder();
- final KTable aTable = builder.table("A");
- final KTable bTable = builder.table("B");
+ final KTable aTable = builder.table("A");
+ final KTable bTable = builder.table("B");
- final KTable fkJoinResult = aTable.join(
+ final KTable fkJoinResult = aTable.join(
bTable,
- value -> value.split("-")[0],
+ value -> Integer.parseInt(value.split("-")[0]),
(aVal, bVal) -> "(" + aVal + "," + bVal + ")",
- Materialized
- .>as("asdf")
- .withKeySerde(Serdes.String())
- .withValueSerde(Serdes.String())
+ Materialized.>as("asdf")
+ .withKeySerde(Serdes.Integer())
+ .withValueSerde(Serdes.String())
);
- final KTable finalJoinResult = aTable.join(
+ final KTable finalJoinResult = aTable.join(
fkJoinResult,
(aVal, fkJoinVal) -> "(" + aVal + "," + fkJoinVal + ")"
);
@@ -133,20 +134,20 @@ public void shouldWorkWithDefaultAndJoinResultSerdes() {
@Test
public void shouldWorkWithDefaultAndEquiJoinResultSerdes() {
final StreamsBuilder builder = new StreamsBuilder();
- final KTable aTable = builder.table("A");
- final KTable bTable = builder.table("B");
+ final KTable aTable = builder.table("A");
+ final KTable bTable = builder.table("B");
- final KTable fkJoinResult = aTable.join(
+ final KTable fkJoinResult = aTable.join(
bTable,
- value -> value.split("-")[0],
+ value -> Integer.parseInt(value.split("-")[0]),
(aVal, bVal) -> "(" + aVal + "," + bVal + ")",
Materialized.as("asdf")
);
- final KTable finalJoinResult = aTable.join(
+ final KTable finalJoinResult = aTable.join(
fkJoinResult,
(aVal, fkJoinVal) -> "(" + aVal + "," + fkJoinVal + ")",
- Materialized.with(Serdes.String(), Serdes.String())
+ Materialized.with(Serdes.Integer(), Serdes.String())
);
finalJoinResult.toStream().to("output");
@@ -157,22 +158,22 @@ public void shouldWorkWithDefaultAndEquiJoinResultSerdes() {
@Test
public void shouldWorkWithDefaultAndProducedSerdes() {
final StreamsBuilder builder = new StreamsBuilder();
- final KTable aTable = builder.table("A");
- final KTable bTable = builder.table("B");
+ final KTable aTable = builder.table("A");
+ final KTable bTable = builder.table("B");
- final KTable fkJoinResult = aTable.join(
+ final KTable fkJoinResult = aTable.join(
bTable,
- value -> value.split("-")[0],
+ value -> Integer.parseInt(value.split("-")[0]),
(aVal, bVal) -> "(" + aVal + "," + bVal + ")",
Materialized.as("asdf")
);
- final KTable finalJoinResult = aTable.join(
+ final KTable finalJoinResult = aTable.join(
fkJoinResult,
(aVal, fkJoinVal) -> "(" + aVal + "," + fkJoinVal + ")"
);
- finalJoinResult.toStream().to("output", Produced.with(Serdes.String(), Serdes.String()));
+ finalJoinResult.toStream().to("output", Produced.with(Serdes.Integer(), Serdes.String()));
validateTopologyCanProcessData(builder);
}
@@ -189,20 +190,20 @@ public void shouldUseExpectedTopicsWithSerde() {
final UniqueTopicSerdeScope serdeScope = new UniqueTopicSerdeScope();
final StreamsBuilder builder = new StreamsBuilder();
- final KTable left = builder.table(
+ final KTable left = builder.table(
LEFT_TABLE,
- Consumed.with(serdeScope.decorateSerde(Serdes.String(), streamsConfig, true),
- serdeScope.decorateSerde(Serdes.String(), streamsConfig, false))
+ Consumed.with(serdeScope.decorateSerde(Serdes.Integer(), streamsConfig, true),
+ 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))
+ final KTable right = builder.table(
+ RIGHT_TABLE,
+ Consumed.with(serdeScope.decorateSerde(Serdes.Integer(), streamsConfig, true),
+ serdeScope.decorateSerde(Serdes.String(), streamsConfig, false))
);
left.join(
right,
- value -> value.split("\\|")[1],
+ value -> Integer.parseInt(value.split("\\|")[1]),
(value1, value2) -> "(" + value1 + "," + value2 + ")",
Materialized.with(null, serdeScope.decorateSerde(Serdes.String(), streamsConfig, false)
))
@@ -212,10 +213,10 @@ public void shouldUseExpectedTopicsWithSerde() {
final Topology topology = builder.build(streamsConfig);
try (final TopologyTestDriver driver = new TopologyTestDriver(topology, streamsConfig)) {
- final TestInputTopic leftInput = driver.createInputTopic(LEFT_TABLE, new StringSerializer(), new StringSerializer());
- final TestInputTopic rightInput = driver.createInputTopic(RIGHT_TABLE, new StringSerializer(), new StringSerializer());
- leftInput.pipeInput("lhs1", "lhsValue1|rhs1");
- rightInput.pipeInput("rhs1", "rhsValue1");
+ final TestInputTopic leftInput = driver.createInputTopic(LEFT_TABLE, new IntegerSerializer(), new StringSerializer());
+ final TestInputTopic rightInput = driver.createInputTopic(RIGHT_TABLE, new IntegerSerializer(), new StringSerializer());
+ leftInput.pipeInput(2, "lhsValue1|1");
+ rightInput.pipeInput(1, "rhsValue1");
}
// 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
@@ -243,17 +244,17 @@ private void validateTopologyCanProcessData(final StreamsBuilder builder) {
final String safeTestName = safeUniqueTestName(getClass(), testName);
config.setProperty(StreamsConfig.APPLICATION_ID_CONFIG, "dummy-" + safeTestName);
config.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy");
- config.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class.getName());
+ config.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.IntegerSerde.class.getName());
config.setProperty(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class.getName());
config.setProperty(StreamsConfig.STATE_DIR_CONFIG, TestUtils.tempDirectory().getAbsolutePath());
try (final TopologyTestDriver topologyTestDriver = new TopologyTestDriver(builder.build(), config)) {
- final TestInputTopic aTopic = topologyTestDriver.createInputTopic("A", new StringSerializer(), new StringSerializer());
- final TestInputTopic bTopic = topologyTestDriver.createInputTopic("B", new StringSerializer(), new StringSerializer());
- final TestOutputTopic output = topologyTestDriver.createOutputTopic("output", new StringDeserializer(), new StringDeserializer());
- aTopic.pipeInput("a1", "b1-alpha");
- bTopic.pipeInput("b1", "beta");
- final Map x = output.readKeyValuesToMap();
- assertThat(x, is(Collections.singletonMap("a1", "(b1-alpha,(b1-alpha,beta))")));
+ final TestInputTopic aTopic = topologyTestDriver.createInputTopic("A", new IntegerSerializer(), new StringSerializer());
+ final TestInputTopic bTopic = topologyTestDriver.createInputTopic("B", new IntegerSerializer(), new StringSerializer());
+ final TestOutputTopic output = topologyTestDriver.createOutputTopic("output", new IntegerDeserializer(), new StringDeserializer());
+ aTopic.pipeInput(1, "999-alpha");
+ bTopic.pipeInput(999, "beta");
+ final Map x = output.readKeyValuesToMap();
+ assertThat(x, is(Collections.singletonMap(1, "(999-alpha,(999-alpha,beta))")));
}
}
}