Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -261,6 +261,7 @@
<allow pkg="kafka.log" />
<allow pkg="scala" />
<allow class="kafka.zk.EmbeddedZookeeper"/>
<allow pkg="com.fasterxml.jackson" />
</subpackage>

<subpackage name="test">
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
import java.nio.ByteBuffer;
import java.util.Objects;

public class ChangedDeserializer<T> implements Deserializer<Change<T>>, WrappingNullableDeserializer<Change<T>, T> {
public class ChangedDeserializer<T> implements Deserializer<Change<T>>, WrappingNullableDeserializer<Change<T>, Void, T> {

private static final int NEWFLAG_SIZE = 1;

Expand All @@ -37,9 +37,9 @@ public Deserializer<T> inner() {
}

@Override
public void setIfUnset(final Deserializer<T> defaultDeserializer) {
public void setIfUnset(final Deserializer<Void> defaultKeyDeserializer, final Deserializer<T> defaultValueDeserializer) {
if (inner == null) {
inner = Objects.requireNonNull(defaultDeserializer, "defaultDeserializer cannot be null");
inner = Objects.requireNonNull(defaultValueDeserializer);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
import java.nio.ByteBuffer;
import java.util.Objects;

public class ChangedSerializer<T> implements Serializer<Change<T>>, WrappingNullableSerializer<Change<T>, T> {
public class ChangedSerializer<T> implements Serializer<Change<T>>, WrappingNullableSerializer<Change<T>, Void, T> {

private static final int NEWFLAG_SIZE = 1;

Expand All @@ -38,9 +38,9 @@ public Serializer<T> inner() {
}

@Override
public void setIfUnset(final Serializer<T> defaultSerializer) {
public void setIfUnset(final Serializer<Void> defaultKeySerializer, final Serializer<T> defaultValueSerializer) {
if (inner == null) {
inner = Objects.requireNonNull(defaultSerializer, "defaultSerializer cannot be null");
inner = Objects.requireNonNull(defaultValueSerializer);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,6 @@

import org.apache.kafka.common.serialization.Deserializer;

public interface WrappingNullableDeserializer<Outer, Inner> extends Deserializer<Outer> {
void setIfUnset(final Deserializer<Inner> defaultDeserializer);
}
public interface WrappingNullableDeserializer<Outer, InnerK, InnerV> extends Deserializer<Outer> {
void setIfUnset(final Deserializer<InnerK> defaultKeyDeserializer, final Deserializer<InnerV> defaultValueDeserializer);
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,6 @@

import org.apache.kafka.common.serialization.Serializer;

public interface WrappingNullableSerializer<Outer, Inner> extends Serializer<Outer> {
void setIfUnset(final Serializer<Inner> defaultSerializer);
public interface WrappingNullableSerializer<Outer, InnerK, InnerV> extends Serializer<Outer> {
void setIfUnset(final Serializer<InnerK> defaultKeySerializer, final Serializer<InnerV> defaultValueSerializer);
}
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ public Deserializer<SubscriptionResponseWrapper<V>> deserializer() {
}

private static final class SubscriptionResponseWrapperSerializer<V>
implements Serializer<SubscriptionResponseWrapper<V>>, WrappingNullableSerializer<SubscriptionResponseWrapper<V>, V> {
implements Serializer<SubscriptionResponseWrapper<V>>, WrappingNullableSerializer<SubscriptionResponseWrapper<V>, Void, V> {

private Serializer<V> serializer;

Expand All @@ -55,9 +55,9 @@ private SubscriptionResponseWrapperSerializer(final Serializer<V> serializer) {
}

@Override
public void setIfUnset(final Serializer<V> defaultSerializer) {
public void setIfUnset(final Serializer<Void> defaultKeySerializer, final Serializer<V> defaultValueSerializer) {
if (serializer == null) {
serializer = Objects.requireNonNull(defaultSerializer, "defaultSerializer cannot be null");
serializer = Objects.requireNonNull(defaultValueSerializer);
}
}

Expand Down Expand Up @@ -94,7 +94,7 @@ public byte[] serialize(final String topic, final SubscriptionResponseWrapper<V>
}

private static final class SubscriptionResponseWrapperDeserializer<V>
implements Deserializer<SubscriptionResponseWrapper<V>>, WrappingNullableDeserializer<SubscriptionResponseWrapper<V>, V> {
implements Deserializer<SubscriptionResponseWrapper<V>>, WrappingNullableDeserializer<SubscriptionResponseWrapper<V>, Void, V> {

private Deserializer<V> deserializer;

Expand All @@ -103,9 +103,9 @@ private SubscriptionResponseWrapperDeserializer(final Deserializer<V> deserializ
}

@Override
public void setIfUnset(final Deserializer<V> defaultDeserializer) {
public void setIfUnset(final Deserializer<Void> defaultKeyDeserializer, final Deserializer<V> defaultValueDeserializer) {
if (deserializer == null) {
deserializer = Objects.requireNonNull(defaultDeserializer, "defaultDeserializer cannot be null");
deserializer = Objects.requireNonNull(defaultValueDeserializer);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ public Deserializer<SubscriptionWrapper<K>> deserializer() {
}

private static class SubscriptionWrapperSerializer<K>
implements Serializer<SubscriptionWrapper<K>>, WrappingNullableSerializer<SubscriptionWrapper<K>, K> {
implements Serializer<SubscriptionWrapper<K>>, WrappingNullableSerializer<SubscriptionWrapper<K>, K, Void> {

private final Supplier<String> primaryKeySerializationPseudoTopicSupplier;
private String primaryKeySerializationPseudoTopic = null;
Expand All @@ -63,9 +63,9 @@ private static class SubscriptionWrapperSerializer<K>
}

@Override
public void setIfUnset(final Serializer<K> defaultSerializer) {
public void setIfUnset(final Serializer<K> defaultKeySerializer, final Serializer<Void> defaultValueSerializer) {
if (primaryKeySerializer == null) {
primaryKeySerializer = Objects.requireNonNull(defaultSerializer, "defaultSerializer cannot be null");
primaryKeySerializer = Objects.requireNonNull(defaultKeySerializer);
Comment thread
bellemare marked this conversation as resolved.
Outdated
}
}

Expand Down Expand Up @@ -110,7 +110,7 @@ public byte[] serialize(final String ignored, final SubscriptionWrapper<K> data)
}

private static class SubscriptionWrapperDeserializer<K>
implements Deserializer<SubscriptionWrapper<K>>, WrappingNullableDeserializer<SubscriptionWrapper<K>, K> {
implements Deserializer<SubscriptionWrapper<K>>, WrappingNullableDeserializer<SubscriptionWrapper<K>, K, Void> {

private final Supplier<String> primaryKeySerializationPseudoTopicSupplier;
private String primaryKeySerializationPseudoTopic = null;
Expand All @@ -123,9 +123,9 @@ private static class SubscriptionWrapperDeserializer<K>
}

@Override
public void setIfUnset(final Deserializer<K> defaultDeserializer) {
public void setIfUnset(final Deserializer<K> defaultKeyDeserializer, final Deserializer<Void> defaultValueDeserializer) {
if (primaryKeyDeserializer == null) {
primaryKeyDeserializer = Objects.requireNonNull(defaultDeserializer, "defaultDeserializer cannot be null");
primaryKeyDeserializer = Objects.requireNonNull(defaultKeyDeserializer);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,10 +66,13 @@ public void init(final InternalProcessorContext context) {
valSerializer = (Serializer<V>) 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(

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@guozhangwang This is 1/2 of the areas that needed the fix. The valueSerde was being passed into the underlying Serde, despite it really needing the keySerde.

context.keySerde().serializer(),
context.valueSerde().serializer()
);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,10 +87,13 @@ public void init(final InternalProcessorContext context) {
this.valDeserializer = (Deserializer<V>) 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(
Comment thread
bellemare marked this conversation as resolved.
Outdated
context.keySerde().deserializer(),
context.valueSerde().deserializer()
);
}
}

Expand Down
Loading