From 46d16fa6f92a4537aba07439cad361acb66fcd61 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Thu, 23 May 2019 14:42:07 -0700 Subject: [PATCH 1/2] MINOR: improve error message for Serde type miss match --- .../kafka/streams/state/StateSerdes.java | 13 +++++- .../internals/ValueAndTimestampSerde.java | 2 +- .../ValueAndTimestampSerializer.java | 2 +- .../kafka/streams/state/StateSerdesTest.java | 40 +++++++++++++++++-- 4 files changed, 49 insertions(+), 8 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java b/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java index 55e9fded99c38..6fcf691fa4c3b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java @@ -21,6 +21,7 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.serialization.Serializer; import org.apache.kafka.streams.errors.StreamsException; +import org.apache.kafka.streams.state.internals.ValueAndTimestampSerializer; import java.util.Objects; @@ -190,12 +191,20 @@ public byte[] rawValue(final V value) { try { return valueSerde.serializer().serialize(topic, value); } catch (final ClassCastException e) { - final String valueClass = value == null ? "unknown because value is null" : value.getClass().getName(); + final String valueClass; + final Class serializerClass; + if (valueSerializer() instanceof ValueAndTimestampSerializer) { + serializerClass = ((ValueAndTimestampSerializer) valueSerializer()).valueSerializer.getClass(); + valueClass = value == null ? "unknown because value is null" : ((ValueAndTimestamp) value).value().getClass().getName(); + } else { + serializerClass = valueSerializer().getClass(); + valueClass = value == null ? "unknown because value is null" : value.getClass().getName(); + } throw new StreamsException( String.format("A serializer (value: %s) is not compatible to the actual value type " + "(value type: %s). Change the default Serdes in StreamConfig or " + "provide correct Serdes via method parameters.", - valueSerializer().getClass().getName(), + serializerClass.getName(), valueClass), e); } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerde.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerde.java index 8be11f32897ea..c02992f2f7ab7 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerde.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerde.java @@ -28,7 +28,7 @@ public class ValueAndTimestampSerde implements Serde> { private final ValueAndTimestampSerializer valueAndTimestampSerializer; private final ValueAndTimestampDeserializer valueAndTimestampDeserializer; - ValueAndTimestampSerde(final Serde valueSerde) { + public ValueAndTimestampSerde(final Serde valueSerde) { Objects.requireNonNull(valueSerde); valueAndTimestampSerializer = new ValueAndTimestampSerializer<>(valueSerde.serializer()); valueAndTimestampDeserializer = new ValueAndTimestampDeserializer<>(valueSerde.deserializer()); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerializer.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerializer.java index 17903952dd715..6db8ccd0b8289 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerializer.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/ValueAndTimestampSerializer.java @@ -24,7 +24,7 @@ import java.util.Map; import java.util.Objects; -class ValueAndTimestampSerializer implements Serializer> { +public class ValueAndTimestampSerializer implements Serializer> { public final Serializer valueSerializer; private final Serializer timestampSerializer; diff --git a/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java b/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java index 56ff71d74dfde..e231d5bd304e5 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java @@ -19,11 +19,16 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.errors.StreamsException; +import org.apache.kafka.streams.state.internals.ValueAndTimestampSerde; import org.junit.Assert; import org.junit.Test; import java.nio.ByteBuffer; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.equalTo; +import static org.junit.Assert.assertThrows; + @SuppressWarnings("unchecked") public class StateSerdesTest { @@ -88,20 +93,47 @@ public void shouldThrowIfValueClassIsNull() { new StateSerdes<>("anyName", Serdes.ByteArray(), null); } - @Test(expected = StreamsException.class) + @Test public void shouldThrowIfIncompatibleSerdeForValue() throws ClassNotFoundException { final Class myClass = Class.forName("java.lang.String"); final StateSerdes stateSerdes = new StateSerdes("anyName", Serdes.serdeFrom(myClass), Serdes.serdeFrom(myClass)); final Integer myInt = 123; - stateSerdes.rawValue(myInt); + final Exception e = assertThrows(StreamsException.class, () -> stateSerdes.rawValue(myInt)); + assertThat( + e.getMessage(), + equalTo( + "A serializer (value: org.apache.kafka.common.serialization.StringSerializer) " + + "is not compatible to the actual value type (value type: java.lang.Integer). " + + "Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.")); + } + + @Test + public void shouldSkipValueAndTimestampeInformationForErrorOnTimestampAndValueSerialization() throws ClassNotFoundException { + final Class myClass = Class.forName("java.lang.String"); + final StateSerdes stateSerdes = + new StateSerdes("anyName", Serdes.serdeFrom(myClass), new ValueAndTimestampSerde(Serdes.serdeFrom(myClass))); + final Integer myInt = 123; + final Exception e = assertThrows(StreamsException.class, () -> stateSerdes.rawValue(ValueAndTimestamp.make(myInt, 0L))); + assertThat( + e.getMessage(), + equalTo( + "A serializer (value: org.apache.kafka.common.serialization.StringSerializer) " + + "is not compatible to the actual value type (value type: java.lang.Integer). " + + "Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.")); } - @Test(expected = StreamsException.class) + @Test public void shouldThrowIfIncompatibleSerdeForKey() throws ClassNotFoundException { final Class myClass = Class.forName("java.lang.String"); final StateSerdes stateSerdes = new StateSerdes("anyName", Serdes.serdeFrom(myClass), Serdes.serdeFrom(myClass)); final Integer myInt = 123; - stateSerdes.rawKey(myInt); + final Exception e = assertThrows(StreamsException.class, () -> stateSerdes.rawKey(myInt)); + assertThat( + e.getMessage(), + equalTo( + "A serializer (key: org.apache.kafka.common.serialization.StringSerializer) " + + "is not compatible to the actual key type (key type: java.lang.Integer). " + + "Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.")); } } From 8c2a1d623efa772be9d4df576cc7b7eed6541ad6 Mon Sep 17 00:00:00 2001 From: "Matthias J. Sax" Date: Thu, 23 May 2019 23:12:45 -0700 Subject: [PATCH 2/2] Github comment --- .../java/org/apache/kafka/streams/state/StateSerdes.java | 4 ++-- .../org/apache/kafka/streams/state/StateSerdesTest.java | 6 +++--- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java b/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java index 6fcf691fa4c3b..1182e50a8868f 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/StateSerdes.java @@ -172,7 +172,7 @@ public byte[] rawKey(final K key) { } catch (final ClassCastException e) { final String keyClass = key == null ? "unknown because key is null" : key.getClass().getName(); throw new StreamsException( - String.format("A serializer (key: %s) is not compatible to the actual key type " + + String.format("A serializer (%s) is not compatible to the actual key type " + "(key type: %s). Change the default Serdes in StreamConfig or " + "provide correct Serdes via method parameters.", keySerializer().getClass().getName(), @@ -201,7 +201,7 @@ public byte[] rawValue(final V value) { valueClass = value == null ? "unknown because value is null" : value.getClass().getName(); } throw new StreamsException( - String.format("A serializer (value: %s) is not compatible to the actual value type " + + String.format("A serializer (%s) is not compatible to the actual value type " + "(value type: %s). Change the default Serdes in StreamConfig or " + "provide correct Serdes via method parameters.", serializerClass.getName(), diff --git a/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java b/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java index e231d5bd304e5..8cb1c7b747ea7 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/StateSerdesTest.java @@ -102,7 +102,7 @@ public void shouldThrowIfIncompatibleSerdeForValue() throws ClassNotFoundExcepti assertThat( e.getMessage(), equalTo( - "A serializer (value: org.apache.kafka.common.serialization.StringSerializer) " + + "A serializer (org.apache.kafka.common.serialization.StringSerializer) " + "is not compatible to the actual value type (value type: java.lang.Integer). " + "Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.")); } @@ -117,7 +117,7 @@ public void shouldSkipValueAndTimestampeInformationForErrorOnTimestampAndValueSe assertThat( e.getMessage(), equalTo( - "A serializer (value: org.apache.kafka.common.serialization.StringSerializer) " + + "A serializer (org.apache.kafka.common.serialization.StringSerializer) " + "is not compatible to the actual value type (value type: java.lang.Integer). " + "Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.")); } @@ -131,7 +131,7 @@ public void shouldThrowIfIncompatibleSerdeForKey() throws ClassNotFoundException assertThat( e.getMessage(), equalTo( - "A serializer (key: org.apache.kafka.common.serialization.StringSerializer) " + + "A serializer (org.apache.kafka.common.serialization.StringSerializer) " + "is not compatible to the actual key type (key type: java.lang.Integer). " + "Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.")); }