[
https://issues.apache.org/jira/browse/KAFKA-20847?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18111749#comment-18111749
]
Matthias J. Sax commented on KAFKA-20847:
-----------------------------------------
Frankly, I am not 100% sure if this one is worth doing – I filed it mainly for
documentation purpose, and to see if anybody is reporting that they would need
a fix for this for their application.
Given that we would use more memory per row in the suppress() buffer, if there
is no clear demand from application developers using KS that they would need
this information added to the changelog, it would actually be a step backward
to do this ticket, as it makes the suppress() operator less effective assuming
the same memory budget.
> Suppress buffer changelog tombstones carry no record metadata
> -------------------------------------------------------------
>
> Key: KAFKA-20847
> URL: https://issues.apache.org/jira/browse/KAFKA-20847
> Project: Kafka
> Issue Type: Improvement
> Components: streams
> Reporter: Matthias J. Sax
> Priority: Minor
>
> When a key is evicted from the suppression buffer,
> InMemoryTimeOrderedKeyValueChangeBuffer writes a delete marker to the
> changelog so the topic can be compacted. That record is written with no
> metadata at all:
> {code:java}
> // lot of nulls here :)
> private void logTombstone(final Bytes key) {
> ((RecordCollector.Supplier) context).recordCollector().send(
> changelogTopic, key, null, null, partition, null, ...);
> } {code}
> Streams itself does not need them. The key is passed as already-serialized
> Bytes with a pass-through BytesSerializer, and on restore it is wrapped
> straight back into Bytes to remove the index entry — no key serde is invoked
> on either side.
> The consequence is for anything reading the changelog outside the restore
> path, such as an internal tool used to inspect buffer state during debugging.
> A key serde that carries part of its contract in the record headers has
> nothing to work with on these records, so the key cannot be decoded at all.
> This is also inconsistent with the headers-aware stores, where the equivalent
> path deliberately keeps the tombstone self-describing:
> {code:java}
> // ChangeLoggingTimestampedKeyValueBytesStoreWithHeaders.delete()
> log(key, null, internalContext.recordContext().timestamp(),
> internalContext.recordContext().headers()); {code}
> A constraint worth recording: the metadata is not available at the point the
> tombstone is written. logTombstone fires from commit() when a dirty key is no
> longer in the index, but the row — and with it the ProcessorRecordContext —
> was already discarded during evictWhile. context.recordContext() at flush
> time belongs to whatever record triggered the commit, not to the evicted key.
> A potential fix would be, to change dirtyKey buffer to also store the rows
> timestamp + headers, to make them available to logTombstone. I would reduce
> the memory savings we get from purging the row pro-actively from the suppress
> buffer, but keeping ts+headers is still smaller than keeping the full row
> with full record context in the suppress buffer.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)