[
https://issues.apache.org/jira/browse/KAFKA-20847?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Arnab Nandy reassigned KAFKA-20847:
-----------------------------------
Assignee: (was: Arnab Nandy)
> 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)