Repository navigation
fix(slatedb): redesign ListOffsets(Timestamp) to match Kafka's time-index semantics - #839
solace-aross wants to merge 12 commits into
Conversation
3b0f25e to
655f363
Compare
sgamelin
left a comment
There was a problem hiding this comment.
Thanks for this. The lookup is right: every batch before the ceiling entry has a max timestamp below that entry's, so starting the forward scan there skips nothing, and inspecting records rather than trusting the header matches Kafka's FileRecords.searchForTimestamp. Moving the index out of the watermark value is the right call, too.
I'm requesting changes for one thing: the compaction rebuild can drop a time index entry written by a concurrent produce (inline on storage.rs:804). After that, ListOffsets(Timestamp) answers "no match" or a later offset for data that exists, until a later compaction removes a record. The other inline comments are non-blocking.
Not inline:
- CHANGELOG.
[Unreleased]has no entry. Clients see three changes: a timestamp lookup now answers the matching record's offset and timestamp instead of the batch's base, the Latest and Earliest timestamps change source, and SlateDB gets a newt/keyspace with a migration on first write. The migration and downgrade behaviour (see thestorage.rs:122comment) are worth a line for operators. - Isolation level. Kafka drops a timestamp match at or above the last fetchable offset (the LSO for
READ_COMMITTED, the high watermark otherwise;Partition.scala#L1589-L1593,#L1622-L1623). TheTimestamparm here still ignoresisolation_level, so aREAD_COMMITTEDlookup can answer an offset inside an open transaction. That was true before this PR, so it doesn't need to be fixed here, but the PR claims Kafka semantics. Could you either filter here or open an issue and note the difference at the code site? - Corrupt batches.
sequential_timestamp_scanskips a batch that fails to decode or inflate and returns a later offset withErrorCode::None, so a consumer resetting by time seeks past the corrupt data without being told. Kafka surfaces an error there. If skipping is deliberate, please say so in a comment at the skip site, and include the decode error in thewarn!so the log shows why it failed. - #838. Both no-match returns in
list_offset_for_timestamp(the shortcut and the scan miss) need #838'sNonewhichever PR merges second. That's mechanical.
Kafka behaviour checked against 3.9.1 by reading code; I didn't build or run this branch.
|
|
||
| match postcard::from_bytes::<Watermark>(&encoded) { | ||
| Ok(watermark) => Ok(watermark), | ||
| Err(_) => { |
There was a problem hiding this comment.
A few things about the migration:
- It runs inside whichever write touches the partition first (produce, EndTxn, DeleteRecords, retention, compaction). It reads every batch of the partition and buffers about one
t/put per batch in that transaction, andtxn_endcan do this for several partitions at once. Right after an upgrade, every busy partition pays this on its first produce, concurrent produces repeat it and then lose the watermark conflict, and nothing logs that a migration is running. Could it run inmaintain(or at startup), committing in chunks? At the least, aninfo!per migrated partition with batch count and elapsed time would help an operator explain slow first produces. - Read-only paths never migrate, so a partition that is never written again (cold, or
retention.ms=-1with no compaction) stays legacy forever, and every timestamp lookup on it scans fromlow. A migration inmaintainwould fix this as well. - The backfill doesn't call
delete_time_indexfirst, although its doc says the caller should. After a downgrade and re-upgrade,t/entries from before the downgrade survive. They can be out of order with the backfilled ones, andprune_time_index_belowstops at the first survivor, so it never removes them. - Treating any decode error as "legacy" works for exactly one format change. When a sixth field is added, a 5-field value fails to decode, then decodes as
WatermarkLegacy(postcard ignores trailing bytes), and the index state is silently dropped and backfilled again. A version marker (an enum wrapper, or a new key) would make the next change explicit.
There was a problem hiding this comment.
Partly done in 4d4bf7a:
- Added an
info!per migrated partition with the batch count andelapsed_ms. Moving the migration intomaintainand chunking it is left for a follow-up. - Same follow-up.
- The backfill now calls
delete_time_indexfirst, and reads batches throughtx. Added a test with stalet/entries. - I didn't add a version marker. Postcard ignores trailing bytes, so a binary built from main can still read the new 5-field value today, which keeps downgrade working. A leading version byte would make that binary fail on every migrated watermark. Instead, a value is treated as legacy only when it decodes as exactly the legacy layout with no bytes left over; anything else is an error, so a future 6-field value can't silently fall back. Happy to revisit if you'd rather trade downgrade for an explicit marker.
There was a problem hiding this comment.
No change on items 1 and 2 this round. Running the migration from maintain, in chunks so cold partitions migrate too, stays a follow-up. Item 3 is in. On item 4, a 3-field value can never decode as the current layout, and postcard ignores trailing bytes, so a future 6-field value decodes as the current layout with the extra field dropped, never as legacy. The exact-length check only guards a value that is neither. The downgrade path stays. Happy to switch to a version byte if you would rather trade that away.
One more fix from this round: prune_time_index_below now reads through the transaction. A migration triggered by DeleteRecords or retention rebuilds the index in the same transaction, and the prune was reading the committed index, so an old entry below the cutoff that shared its key with a rebuilt entry deleted the rebuilt one (legacy_migration_by_delete_records_keeps_a_rebuilt_entry).
There was a problem hiding this comment.
No change on 1 and 2 this round; running the migration from maintain, in chunks so cold partitions migrate too, stays a follow-up. 3 is in. On 4, the exact-length check stands as above, and I am happy to switch to a version byte if you would rather trade the downgrade path for it. One related change: lookups now read the watermark, the index and the batches through one snapshot, so a prune that commits mid-lookup cannot move the index out from under the start point.
| .and_then(|ts| ts.iter().find(|(_, off)| **off >= offset)) | ||
| .and_then(|(ts, _)| to_system_time(*ts).ok()); | ||
| let timestamp = self | ||
| .timestamp_at_offset(metadata.id, topition.partition, offset) |
There was a problem hiding this comment.
Question: Kafka 3.9.1 returns NO_TIMESTAMP for both EARLIEST and LATEST (UnifiedLog.scala#L1325-L1356: "For the earliest and latest, we do not need to return the timestamp"). This PR keeps Nisshi's existing behaviour of returning one, and pays for it: Earliest now does up to three scans (more when low isn't a batch base) plus a full batch inflate on every call, and Latest needs the new persisted last_batch_max_timestamp with its own invalidation in prune, compaction and migration. Lag exporters poll Earliest and Latest on every partition. pg returns record timestamps here too, so this is a cross-backend decision rather than this PR's, but would you consider returning None for both and dropping the field, or doing that in a follow-up?
There was a problem hiding this comment.
I'd do that as a follow-up, since pg returns record timestamps here too and it should change on every backend at once. This PR keeps the current behaviour.
There was a problem hiding this comment.
Still a follow-up. pg returns record timestamps here too, so it should change on every backend at once.
There was a problem hiding this comment.
Still a follow-up for every backend at once; pg returns record timestamps here too. One fix this round in the same code: Earliest reads the next batch when compaction removed the batch at the log start, so the timestamp matches the changelog wording.
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn ceiling_scan_exact_boundary() { |
There was a problem hiding this comment.
The ceiling path is only exercised for speed, not for correctness. Here T=150 finds the last entry, and if the ceiling seek were off by one (scan_from(target_millis + 1)), the fallback to low would still answer 2. As far as I can tell, forcing start_offset = low everywhere also passes every test. A case where T equals a middle entry's timestamp (for example {100→0, 150→2, 200→5} at T=150) would catch both.
There was a problem hiding this comment.
Added the middle-entry case in 4d4bf7a: {100->0, 150->2, 200->5} at T=150 must answer 2. With the seek mutated to target_millis + 1, only this test fails, so it covers the off-by-one. Forcing start_offset = low can't be caught through the answer, though: every batch before the ceiling entry has a max timestamp below T, so a scan from low returns the same offset, only more slowly. That one would need a read counter.
There was a problem hiding this comment.
Nothing new; the middle-entry case is in. A read counter for the start_offset = low case would need a counting Db wrapper, which I would rather not add for one test.
There was a problem hiding this comment.
The middle-entry case is in. A read counter for the start_offset = low case needs a counting Db wrapper, which I would rather not add for one test. Happy to add it if you feel strongly.
7d71950 to
92b86b3
Compare
|
@sgamelin pushed
|
92b86b3 to
16dcf70
Compare
|
Rebased on main and ready for another look. Your seven commits are unchanged. Two commits are new:
I also checked whether the backfill can race a produce here, as it could on #840. It cannot: produce assigns the offset, writes the batch and the watermark, and commits in one transaction, and the backfill runs inside the write's own transaction. One difference between the engines to settle in a single change: on a corrupt batch this PR skips and logs, as you suggested, while #840 fails the request. A storage error fails every partition in the ListOffsets frame, which is harsher than Kafka's per-partition error code, so the right answer for both may be a per-partition error. |
sgamelin
left a comment
There was a problem hiding this comment.
Thanks for the thorough round. The compaction race I blocked on is fixed: policy_compact now begins its transaction before it scans, and reads the batches, the index and the watermark through tx (7936bbd). The isolation filter matches Kafka 3.9.1, which keeps a timestamp match only below the last stable offset under READ_COMMITTED (Partition.scala#L1589-L1593, #L1622-L1623). The straddling-batch fix, the backfill's delete_time_index call, the empty-partition reset, the middle-entry test and the None answers all look right. Reading each ListOffsets request through one snapshot also closes a race I had missed.
Two things I'd like before merge, both inline: a compaction that loses its commit to a produce now ends the whole maintenance pass, and the CHANGELOG's downgrade sentence doesn't say what timestamp lookups do on the older binary. The rest is optional.
Not inline:
- #847.
sequential_timestamp_scanskips a batch whose headermax_timestampis below the target, and the index advances on the same header. Until #847 lands, aCreateTimebatch's header is whatever the producer wrote, so a header that understates its records hides them from the lookup. Kafka'sLogValidatorprevents this by recomputing the header, and #847 does the same in the shared produce path. It would be good to land #847 before this PR or with it. - Follow-ups. Three items are deferred: running the migration from
maintain(a partition that is never written again stays legacy, and every timestamp lookup on it scans from the log start), returningNO_TIMESTAMPfor Earliest and Latest on every backend, and a per-partition error for a corrupt batch on both engines (your note from the last round). Could you open an issue for each, so they don't live only in this thread? - Optional. Four of the five new commits fix the same mistake: a read through
self.dbwhere the operation commits through a transaction or answers from a snapshot. In slatedb 0.14,Db,DbSnapshotandDbTransactionall implementDbReadOps, so the read helpers could takeR: DbReadOps. Thenpolicy_delete,delete_recordsanddelete_topic, which still scanself.dbinside a transaction, could move totx.scan_prefix. Those three are correct today only because their watermark read conflicts. Fine as a follow-up.
Earlier points:
- Compaction race (
storage.rs:804): fixed. - Straddling batch (
:402): fixed, with a test. - Migration (
:122): 3 is fixed. On 4, I agree with keeping the downgrade path, and the exact-length check covers the case I was worried about. 1 is partly done with theinfo!line; the rest of 1, and 2, are the follow-up above. latest_indexed_timestampreset (:277): fixed.- Earliest and Latest timestamps (
:1799): fine as a follow-up. - Comment wording (
:197): fixed. - Middle-entry test (
tests.rs:1513): added. I agree a read counter isn't worth aDbwrapper. - CHANGELOG, isolation level, corrupt batches and #838: done, with one CHANGELOG fix inline.
I compared against Kafka 3.9.1 by reading code; I didn't build or run this branch.
Replace the per-partition `timestamps` map (keyed by base_timestamp instead of max_timestamp, and silently dropping batches whenever a new one's timestamp regressed) with a proper `t/` time index maintained by Kafka's own TimeIndex.maybeAppend rule: only a strictly-increasing max_timestamp is appended, floored at NO_TIMESTAMP (-1) so a negative timestamp is never indexed. ListOffsets(Timestamp) now does a ceiling range-scan of the index to find a starting offset, then a sequential scan of `b/` batches forward from there, since a batch's header max_timestamp can be stale after compaction (overstated, never understated) and the index itself can have a coverage gap below its first surviving entry once pruning has run (delete_records/policy_delete only prune `t/` entries; Kafka never hits this because its index is per-segment). policy_compact instead rebuilds the index from scratch from the surviving batches, replaying the monotonic rule, which is what Kafka's own LogCleaner does on segment compaction. A legacy (pre-time-index) watermark is detected by any postcard decode failure against the new (5-field) shape, and backfilled in place from `b/` batch headers the first time a write transaction sees it; until then, reads fall back to a slower but correct index-empty scan. ListOffsets(Latest) answers its timestamp from a new watermark field (the most recent batch's own max_timestamp) set unconditionally at every write site, rather than a batch_base_at_or_before binary search per call, since it's a consumer hot path. delete_topic gains the `t/` keyspace cleanup step this change needs (nothing wildcard-deletes by topic uuid, so a new keyspace leaks forever without one), built from a properly postcard-encoded prefix key rather than the raw-byte-vec pattern its existing (broken, filed separately) consumer-offset cleanup step uses. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
Required changes from review round 1: - Strengthen compaction_rebuild_finds_correct_entry_after_removal with a new discriminating scenario (removed entry sits in the MIDDLE of the index, not the head), since the original test passed under both the real full-rebuild-from-survivors logic and a naive per-entry prune. - Strengthen legacy_watermark_backfills_and_answers_correctly so it actually depends on backfill_time_index running: the new scenario puts the legacy data's max timestamp above the post-migration produce's, so only the O(1) no-match shortcut (gated on latest_indexed_timestamp) discriminates backfilled vs. not. Also assert time_index_count == 2. - Fix a real regression: last_batch_max_timestamp (Latest's O(1) fast path) lingered after delete_records/policy_delete fully emptied a partition, so Latest kept answering with a pruned batch's stale timestamp instead of None. Added Engine::clear_latest_if_partition_empty, wired into both prune call sites, plus regression tests for each. - Add warn! logging at the undecodable/un-inflatable batch skip points in sequential_timestamp_scan, timestamp_at_offset, and backfill_time_index, so corrupt data is observable instead of silently invisible. - Add the ticket's own worked example as a real test. - Document the reviewed-and-confirmed-harmless cold-start migration quirk (backfill reads via self.db.scan, not tx.scan, so it can transiently index an offset the same transaction is about to delete). All scenarios verified to actually discriminate: each new/strengthened test was run against a temporarily reverted "naive" version of the corresponding fix and confirmed to fail, then restored. Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com> Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs
Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com> Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs
Drops ticket citations used as history narration (the rule's "no ticket numbers as scar tissue"), states a couple of requirements directly instead of framing them as "on main, before this change", and renames a test from ticket_worked_example to describe what it actually exercises. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
- Compaction begins its transaction first and reads the batches and the time index through it, so a produce that commits during compaction fails the commit instead of losing its index entry. - The legacy-watermark migration deletes existing `t/` entries before it rebuilds them, reads batches through the transaction, and logs the batch count and elapsed time per partition. A stored watermark is read as legacy only when it decodes as exactly the legacy layout. - A timestamp lookup with no usable index entry starts at the batch that holds the log start offset and skips records below it. - A prune that empties a partition also resets `latest_indexed_timestamp`, so a refill with older timestamps is indexed again. - The timestamp scan logs the decode error of a skipped batch, and comments state the deliberate differences from Kafka (corrupt batch, isolation level). - Comments describe the end state; CHANGELOG gains a Fixed entry. - New tests: ceiling lookup at a middle index entry, a log start inside a batch, refill after emptying, partial prune, stale-entry migration. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011P97dHLhTRdJ35fMJFpYPg Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…mestamp lookups A timestamp lookup with no record at or after the target now answers no offset, on both the shortcut and the scan miss, which the broker maps to -1 as every other engine does. Under READ_COMMITTED, a match at or past the last stable offset is no match, as in Apache Kafka. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A migration triggered by DeleteRecords or retention rebuilds the index in the transaction before the prune runs. The prune read the committed index, so an old entry below the new log start that shared its key with a rebuilt entry deleted the rebuilt entry from the transaction, and a later lookup skipped records. The prune now reads through the transaction, as the rebuild does. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
… compaction Compaction removes a batch whose records are all superseded and leaves the log start where it is, so no batch starts at the log start. The Earliest lookup then fell back to the batch before the log start, which holds no record at or after it, and answered no timestamp. It now reads the next batch when the batch before the log start does not reach it, so the timestamp is the first surviving record's, as the changelog states. The two warn lines on this path now carry the decode error. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The lookup read the watermark, the time index and the batches with separate reads of the live database. The index is read twice, for its first entry and for the ceiling entry, and the choice to trust the ceiling depends on the first entry. A DeleteRecords or retention prune that committed between the two reads removed the entries below the new log start from the second read only, so the ceiling skipped the batches between the new log start and the first surviving entry, and the lookup answered a later offset than the first match. ListOffsets now takes one snapshot per request and reads everything through it. The batch probes take the snapshot as a parameter, so Fetch takes one too, and its probe and its scan see the same batches. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The Earliest lookup read the batch that holds the log start and, when that batch did not reach it, the next batch. A batch at the log start that failed to decode or inflate answered no timestamp without reading the next batch, while the timestamp lookup scan skips such a batch and keeps going. The Earliest timestamp is now the lookup scan with a target that every timestamp reaches, started at the batch that holds the log start. The scan steps over a batch that does not reach the log start, which is how compaction leaves the batch before it, and over any run of corrupt batches, so both paths answer the first readable record at or after the log start. The helper that read one batch is gone. Three tests: the log start inside a five-record batch with a second batch after it, answered from inside the straddling batch; the same layout with the straddling batch un-inflatable, answered from the next batch; and an un-inflatable batch that starts at the log start, answered from the next batch. The log start is written to the watermark directly, because DeleteRecords removes a batch whose base offset is below the new log start. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…conflict A produce that commits to a partition while compaction or retention reads it makes the maintenance transaction fail to commit, as intended. The error then left `maintain` before the partitions after that one and before lake maintenance, and the broker logs it only at debug. Topics iterate in name order, so a partition under steady writes with records to compact held back the same partitions on every tick. A commit that fails with a transaction conflict is now logged at warn with the topic and partition, and the tick moves on to the next partition; the next tick retries. Any other error is still returned. The per-partition body of each policy is now a function that takes the transaction, so a test can begin a transaction, run the body, produce to the partition, and assert that the commit fails and that the produce's batch and time index entry survive. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The changelog is generated from the squash commit of each pull request when a release is cut, so a pull request does not edit CHANGELOG.md. The entry's content moves to the pull request description. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
16dcf70 to
e718378
Compare
What changes for a user or operator
ListOffsetswith a timestamp answers the offset and timestamp of the first record at or after that timestamp, instead of the base offset of its batch. A lookup with no such record answers no match (offset -1), as on every other engine.READ_COMMITTED, a match at or past the last stable offset answers no match, as in Apache Kafka.Earliesttimestamp comes from the first record at or after the log start offset, read from the batch that holds the log start or from the next batch. TheLatesttimestamp comes from the newest batch's maximum timestamp, and is available on each partition after the partition's first write on this version; before that write it is missing.t/time-index keyspace, and builds it for a partition on the first write to that partition after the upgrade. Each build logs aninfoline with the batch count and elapsed time.warnand retried on the next tick; the partitions after it, and lake maintenance, run in the same tick.lowandhigh, and ignores thet/keyspace, but it answers a timestamp lookup and theEarliestandLatesttimestamps only from the watermark'stimestampsmap, which this version always writes asNone. After a downgrade, a timestamp lookup therefore answers no match on every partition that this version wrote to, until the older binary appends a batch to it; from then on a lookup for any earlier target answers that batch's offset, so a consumer that resets by time skips the data before it. A topic that the older binary deletes leaves itst/entries behind. On the next upgrade, the first write to each partition that the older binary wrote to rebuilds that partition's index from its batches.Problem
SlateDB's
ListOffsets(Timestamp)did not answer Kafka's question. It kept, inside each partition's watermark, abase_timestamp -> base_offsetmap and answered with the first entry at or after T. That differs from Kafka in five ways:base_timestampshare a map key, so the later batch overwrites the earlier one and the lookup skips the earlier batch.Design
t/keys, keyed bymax_timestamp(big-endian, so offset-ordered alongside it), instead of an in-watermark map. A batch is indexed only when itsmax_timestampraises the partition's running maximum (append_time_index, Kafka's ownTimeIndex.maybeAppend), so entries rise in both time and offset together. This fixes bugs 1 to 3 and bounds index growth to new maxima rather than one entry per produce. A timestamp of -1 or below (NO_TIMESTAMP) is never indexed.list_offset_for_timestamp): an O(1) no-match answer when T is beyond everything ever indexed; otherwise find the index's ceiling entry and scanb/batches forward from its offset, inspecting each batch's own records for the first one at or after the log start whose timestamp is>= T(fixes bug 4). A batch header can overstate but never understate the true max after compaction, so a qualifying header is a reason to look inside the batch, not a match by itself. When the index cannot rule out a gap below its first surviving entry (pruning deletest/entries without uncovering new ones, which Kafka's per-segment index never has to handle), the scan starts at the batch that holds the log start, and records below the log start are skipped, so a batch that straddles the log start is handled.list_offsetsreads the watermark, the index and the batches of every partition in the request through one SlateDB snapshot. The index is read twice (its first entry decides whether the ceiling entry can be trusted), and a prune that committed between two reads of the live database removed entries from the second read only, so the lookup skipped the batches between the new log start and the first surviving entry. Fetch takes a snapshot for its batch probe and scan too, since it shares the probe helpers.READ_COMMITTED, a match at or past the last stable offset answers no match, as in Kafka.delete_recordsandpolicy_deleteprunet/entries below the new log start through the write transaction (a prefix delete that stops at the first survivor, since entries are both timestamp- and offset-ordered).policy_compactdeletes and replays the index from the surviving batches in offset order, because a partially compacted batch keeps its original, now possibly stale, header, and pruning individual entries can leave the index naming an offset that no longer carries the batch that justified it. The compaction reads the batches, the watermark and the index through its transaction, so a produce that commits during the compaction makes the compaction fail to commit instead of losing its index entry. When a prune empties the partition, the watermark's time-index state is reset, soListOffsets(Latest)does not answer a deleted batch's timestamp and a refill with older timestamps is indexed again.compact_partition,delete_expired_prefix) takes the transaction, andcommit_maintenancecommits it. A commit that fails with a transaction conflict is logged atwarnwith the topic and partition and skipped until the next tick, instead of returned, because a returned error ended the tick before the partitions after that one and before lake maintenance, and the broker logs it only atdebug; topics iterate in name order, so a partition under steady writes held back the same partitions on every tick. Any other error is still returned.Earliesttimestamp is the first readable record at or after the log start. It is the lookup scan with a target that every timestamp reaches, started at the batch that holds the log start, so it steps over the batch before the log start when compaction removed the batch at it (compaction leaves the log start in place) and over a corrupt batch. TheLatesttimestamp is the newest batch'smax_timestamp, kept on the watermark so the answer is O(1).Earliesttimestamp, and the backfill skip a batch that fails to decode or inflate, and log the decode error atwarnwith topic, partition and offset. Kafka returns an error for a corrupt batch here; this is a deliberate difference, noted at the skip site. A failure in the backfill would abort the write that triggered the migration, and every later write to the partition would fail the same way.Migration
A legacy 3-field watermark is distinguished from the new 5-field shape by its decoding: a value is treated as legacy only when it fails to decode as the current layout and decodes as exactly the legacy layout with no bytes left over. On the first write that touches the partition, it is migrated in place: the time index is rebuilt by replaying the monotonic rule over every batch in the partition's
b/keyspace (backfill_time_index), reading through the write's transaction, rather than trusting the old, lossy, wrongly keyed map. The backfill deletes anyt/entries already present first, so entries left behind by a downgrade and re-upgrade cannot survive out of order. Each migration logs aninfoline with the batch count and elapsed time. Read-only paths (list_offsets,offset_stage) tolerate the legacy shape without migrating, and skip the O(1) shortcut in that case, since an index that is not built yet understates reality.A binary without the time index still reads
lowandhighfrom the 5-field value, because postcard ignores trailing bytes. What its lookups answer after a downgrade is described above.Conventions
Noneis the no-match answer on both the O(1) shortcut and the scan miss, as on every engine since #838 merged. The broker maps it to -1.Testing
cargo nextest run -p nisshi-storage-slatedb --all-features: 63 passed, including:delete_records).target + 1).Earliesttimestamp comes from inside it (fails when the batch before the log start is dropped, or when the next batch is tried first).Earliesttimestamp when the batch that straddles the log start, or the batch that starts at it, cannot be inflated: it comes from the next batch (both fail when the scan stops at a corrupt batch).t/entry survive, nothing of the pass applies, and the next tick applies it (both fail when the conflict is returned as an error).delete_recordsthat must keep a rebuilt entry, and a migration that must delete stalet/entries.READ_COMMITTEDstops at the last stable offset; a future timestamp and a scan miss both answer no match.delete_recordsthat empties a partition clears theLatesttimestamp, a refill with older timestamps is indexed again, and a partial prune keeps the running maximum.Earliestafter compaction removes the batch at the log start answers the first surviving record's timestamp.delete_topicleaves not/keys behind.delete_recordsanswers from the state before it (fails when the index or batch reads go to the live database).cargo nextest run -p nisshi-broker -p nisshi-storage -p nisshi-storage-slatedb --all-features --no-fail-fast -E '!test(/::pg::/)': 527 run, 527 passed.cargo fmt --all --check,cargo clippy --workspace --all-features --all-targets -- -D warnings, andRUSTDOCFLAGS="-D warnings" cargo doc --workspace --all-features --no-deps --document-private-itemsare clean.Related
READ_COMMITTEDbound, answerNonefor no match, and never indexNO_TIMESTAMP. They differ on a batch that cannot be read: fix(dynostore): ListOffsets by timestamp now answers from record time, not object write time #840 answersKAFKA_STORAGE_ERRORin that partition'serror_codewhile the other partitions of the request answer normally; this PR skips the batch with awarnin the lookup scan, the backfill and theEarliesttimestamp. A follow-up ticket has been raised to align them.🤖 Generated with Claude Code
https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD