aliehsaeedii commented on code in PR #22165:
URL: https://github.com/apache/kafka/pull/22165#discussion_r3558729017
##########
streams/src/test/java/org/apache/kafka/streams/state/internals/TimeOrderedKeyValueBufferTest.java:
##########
@@ -306,10 +306,74 @@ public void shouldReturnPriorValueForBufferedKey(final
String testName, final Fu
context.setRecordContext(recordContext);
buffer.put(1L, new Record<>("A", new Change<>("new-value",
"old-value"), 0L), recordContext);
buffer.put(1L, new Record<>("B", new Change<>("new-value", null), 0L),
recordContext);
- assertThat(buffer.priorValueForBuffered("A"),
is(Maybe.defined(ValueTimestampHeaders.make("old-value", -1, new
RecordHeaders()))));
+ assertThat(buffer.priorValueForBuffered("A"),
is(Maybe.defined(ValueTimestampHeaders.make("old-value", 0L, new
RecordHeaders()))));
Review Comment:
These only cover a single put per key. The value-vs-metadata mismatch shows
up when the same key is put twice before eviction (prior value from put #1,
timestamp/headers from put #2) — worth a test for that case.
##########
streams/src/test/java/org/apache/kafka/streams/state/internals/TimeOrderedKeyValueBufferTest.java:
##########
@@ -306,10 +306,74 @@ public void shouldReturnPriorValueForBufferedKey(final
String testName, final Fu
context.setRecordContext(recordContext);
buffer.put(1L, new Record<>("A", new Change<>("new-value",
"old-value"), 0L), recordContext);
buffer.put(1L, new Record<>("B", new Change<>("new-value", null), 0L),
recordContext);
- assertThat(buffer.priorValueForBuffered("A"),
is(Maybe.defined(ValueTimestampHeaders.make("old-value", -1, new
RecordHeaders()))));
+ assertThat(buffer.priorValueForBuffered("A"),
is(Maybe.defined(ValueTimestampHeaders.make("old-value", 0L, new
RecordHeaders()))));
assertThat(buffer.priorValueForBuffered("B"), is(Maybe.defined(null)));
}
+ @ParameterizedTest
+ @MethodSource("parameters")
+ public void shouldPropagateHeadersThroughEviction(final String testName,
final Function<String, B> bufferSupplier) {
+ setup(testName, bufferSupplier);
+ final TimeOrderedKeyValueBuffer<String, String, Change<String>> buffer
= bufferSupplier.apply(testName);
+ final MockInternalProcessorContext<?, ?> context = makeContext();
+ buffer.init(context, buffer);
+
+ final RecordHeaders headers = new RecordHeaders(new Header[]{new
RecordHeader("h1", "v1".getBytes(UTF_8))});
Review Comment:
Every getBytes(...) call in this file specifies UTF_8 and UTF_8 is already
imported. Keeping getBytes(UTF_8) matches that convention and avoids relying on
the platform default charset.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]