diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStore.java index 1c104915c13a9..aecf69f495a9b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredTimestampedWindowStore.java @@ -53,6 +53,6 @@ void initStoreSerde(final ProcessorContext context) { serdes = new StateSerdes<>( ProcessorStateManager.storeChangelogTopic(context.applicationId(), name()), keySerde == null ? (Serde) context.keySerde() : keySerde, - valueSerde == null ? new ValueAndTimestampSerde<>((Serde) context.keySerde()) : valueSerde); + valueSerde == null ? new ValueAndTimestampSerde<>((Serde) context.valueSerde()) : valueSerde); } } diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampWindowStoreTest.java index a3522f3114483..28d00ad37aa50 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredTimestampWindowStoreTest.java @@ -24,7 +24,9 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; +import org.apache.kafka.streams.state.ValueAndTimestamp; import org.apache.kafka.streams.state.WindowStore; import org.apache.kafka.test.InternalMockProcessorContext; import org.apache.kafka.test.NoOpRecordCollector; @@ -35,6 +37,7 @@ import org.junit.Test; import static org.junit.Assert.assertNull; +import static org.junit.Assert.fail; public class MeteredTimestampWindowStoreTest { private InternalMockProcessorContext context; @@ -89,4 +92,51 @@ public void shouldNotExceptionIfFetchReturnsNull() { assertNull(store.fetch("a", 0)); } + @Test + public void shouldNotThrowExceptionIfSerdesCorrectlySetFromProcessorContext() { + EasyMock.expect(innerStoreMock.name()).andStubReturn("mocked-store"); + EasyMock.replay(innerStoreMock); + final MeteredTimestampedWindowStore store = new MeteredTimestampedWindowStore<>( + innerStoreMock, + 10L, // any size + "scope", + new MockTime(), + null, + null + ); + store.init(context, innerStoreMock); + + try { + store.put("key", ValueAndTimestamp.make(42L, 60000)); + } catch (final StreamsException exception) { + if (exception.getCause() instanceof ClassCastException) { + fail("Serdes are not correctly set from processor context."); + } + throw exception; + } + } + + @Test + public void shouldNotThrowExceptionIfSerdesCorrectlySetFromConstructorParameters() { + EasyMock.expect(innerStoreMock.name()).andStubReturn("mocked-store"); + EasyMock.replay(innerStoreMock); + final MeteredTimestampedWindowStore store = new MeteredTimestampedWindowStore<>( + innerStoreMock, + 10L, // any size + "scope", + new MockTime(), + Serdes.String(), + new ValueAndTimestampSerde<>(Serdes.Long()) + ); + store.init(context, innerStoreMock); + + try { + store.put("key", ValueAndTimestamp.make(42L, 60000)); + } catch (final StreamsException exception) { + if (exception.getCause() instanceof ClassCastException) { + fail("Serdes are not correctly set from constructor parameters."); + } + throw exception; + } + } }