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
Original file line number Diff line number Diff line change
Expand Up @@ -802,8 +802,8 @@ public <K, V> ProducerRecord<K, V> readOutput(final String topic,
if (record == null) {
return null;
}
final K key = keyDeserializer.deserialize(record.topic(), record.key());
final V value = valueDeserializer.deserialize(record.topic(), record.value());
final K key = keyDeserializer.deserialize(record.topic(), record.headers(), record.key());
final V value = valueDeserializer.deserialize(record.topic(), record.headers(), record.value());
return new ProducerRecord<>(record.topic(), record.partition(), record.timestamp(), key, value, record.headers());
}

Expand Down Expand Up @@ -906,8 +906,8 @@ <K, V> TestRecord<K, V> readRecord(final String topic,
if (record == null) {
throw new NoSuchElementException("Empty topic: " + topic);
}
final K key = keyDeserializer.deserialize(record.topic(), record.key());
final V value = valueDeserializer.deserialize(record.topic(), record.value());
final K key = keyDeserializer.deserialize(record.topic(), record.headers(), record.key());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

readOutput is deprecated. Thus not sure if it's worth to fix?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thus not sure if it's worth to fix?

Personally, what we should add to readOutput is "deprecation" rather than "a bug". Hence, it is worthwhile to fix if the fix does not break anything.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Not sure if I can follow? readOutput is marked as deprecated: https://github.com/apache/kafka/blob/trunk/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java#L797

I can still fix it on the side, but nobody should use it any longer and thus the gain seems minimal.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

but nobody should use it any longer and thus the gain seems minimal.

You are right. The benefit is too low.

final V value = valueDeserializer.deserialize(record.topic(), record.headers(), record.value());
return new TestRecord<>(key, value, record.headers(), record.timestamp());
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,13 @@
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.apache.kafka.common.serialization.ByteArraySerializer;
import org.apache.kafka.common.serialization.Deserializer;
import org.apache.kafka.common.serialization.LongDeserializer;
import org.apache.kafka.common.serialization.LongSerializer;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.SystemTime;
Expand Down Expand Up @@ -70,6 +73,7 @@
import java.util.Objects;
import java.util.Properties;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.regex.Pattern;

import static org.apache.kafka.common.utils.Utils.mkEntry;
Expand Down Expand Up @@ -711,6 +715,56 @@ public void shouldUseSourceSpecificDeserializers() {
assertThat(result2.getValue(), equalTo(source2Value));
}

@Test
public void shouldPassRecordHeadersIntoSerializersAndDeserializers() {
testDriver = new TopologyTestDriver(setupSourceSinkTopology(), config);

final AtomicBoolean passedHeadersToKeySerializer = new AtomicBoolean(false);
final AtomicBoolean passedHeadersToValueSerializer = new AtomicBoolean(false);
final AtomicBoolean passedHeadersToKeyDeserializer = new AtomicBoolean(false);
final AtomicBoolean passedHeadersToValueDeserializer = new AtomicBoolean(false);

final Serializer<byte[]> keySerializer = new ByteArraySerializer() {
@Override
public byte[] serialize(final String topic, final Headers headers, final byte[] data) {
passedHeadersToKeySerializer.set(true);
return serialize(topic, data);
}
};
final Serializer<byte[]> valueSerializer = new ByteArraySerializer() {
@Override
public byte[] serialize(final String topic, final Headers headers, final byte[] data) {
passedHeadersToValueSerializer.set(true);
return serialize(topic, data);
}
};

final Deserializer<byte[]> keyDeserializer = new ByteArrayDeserializer() {
@Override
public byte[] deserialize(final String topic, final Headers headers, final byte[] data) {
passedHeadersToKeyDeserializer.set(true);
return deserialize(topic, data);
}
};
final Deserializer<byte[]> valueDeserializer = new ByteArrayDeserializer() {
@Override
public byte[] deserialize(final String topic, final Headers headers, final byte[] data) {
passedHeadersToValueDeserializer.set(true);
return deserialize(topic, data);
}
};

final TestInputTopic<byte[], byte[]> inputTopic = testDriver.createInputTopic(SOURCE_TOPIC_1, keySerializer, valueSerializer);
final TestOutputTopic<byte[], byte[]> outputTopic = testDriver.createOutputTopic(SINK_TOPIC_1, keyDeserializer, valueDeserializer);
inputTopic.pipeInput(testRecord1);
outputTopic.readRecord();

assertThat(passedHeadersToKeySerializer.get(), equalTo(true));
assertThat(passedHeadersToValueSerializer.get(), equalTo(true));
assertThat(passedHeadersToKeyDeserializer.get(), equalTo(true));
assertThat(passedHeadersToValueDeserializer.get(), equalTo(true));
}

@Test
public void shouldUseSinkSpecificSerializers() {
final Topology topology = new Topology();
Expand Down