Tiered Storage relies on an external object store and a metadata store. As a principal engineer, what consistency and durability invariants must hold, and where can things go wrong?
answer
- readable iff durable bytes AND COPY_FINISHED
- RLMM is source of truth, not bucket listing
- STARTED/FINISHED = crash-safe, idempotent
- orphans acceptable; missing-data not
- DELETE_STARTED before byte delete
basics
~20 sA segment must be durably in the object store before its metadata is marked finished/readable; metadata must be the single source of truth for what's readable. Risks include orphaned objects, eventual-consistency reads, metadata/data divergence, and double-offload across leader changes.
solid answer
~50 sThe core invariant: **a remote segment becomes readable only after its bytes are durably stored AND its metadata is in COPY_SEGMENT_FINISHED**. The RLMM (default: the replicated __remote_log_metadata topic) is the authoritative index of what exists remotely — never the object-store listing. This ordering gives crash safety: a crash mid-upload leaves the segment in COPY_SEGMENT_STARTED, treated as absent, and retried, possibly producing **orphaned objects** the RSM must tolerate/garbage-collect. Failure modes to reason about: (1) **object-store eventual consistency** — read-after-write or list-after-delete lag can make a just-uploaded segment briefly unreadable; (2) **metadata/data divergence** — FINISHED metadata but missing/corrupt bytes breaks reads, so durability must precede FINISHED; (3) **leader changes mid-offload** causing duplicate uploads, reconciled via RLMM and idempotent object keys; (4) **deletion ordering** — DELETE_STARTED before deleting bytes, FINISHED after, so a half-deleted segment is never served. Metadata is partitioned per user-partition and itself replicated, so metadata durability rides on Kafka's own guarantees.
go deeper
Understand that the bytes must really be saved remotely before Kafka says the data is available there.
Explain the STARTED/FINISHED states give crash safety and that metadata, not the bucket listing, defines what's readable.
Discuss eventual consistency, metadata/data divergence, and idempotent re-uploads across leader changes.
Reason end-to-end about ordering invariants, orphan-vs-missing-data tradeoffs, metadata-store availability/DR, and how to design RSM/RLMM implementations that stay correct under partial failure.
## Why invariants matter here Tiered storage puts an **external object store** and a **separate metadata store** into Kafka's read/durability path. Kafka's local-only design previously made 'is this data readable?' a local-disk question. Now correctness depends on two distributed systems agreeing. The design uses ordering rules and a two-phase state machine to keep them consistent. ## The central invariant **Readable iff (bytes durable in object store) AND (RLMM state == COPY_SEGMENT_FINISHED).** Consequences: - The **RLMM is the source of truth** for what is remotely readable — the broker never trusts an object-store bucket listing to decide readability. (Listings are eventually consistent and may include orphans.) - The RSM must guarantee the upload is **durable** (e.g., S3 returns success only after the object is committed) **before** the broker writes FINISHED to the RLMM. Writing FINISHED first would risk serving a read for bytes that aren't actually there yet. ## The state machine as a correctness tool COPY_SEGMENT_STARTED → COPY_SEGMENT_FINISHED, and DELETE_SEGMENT_STARTED → DELETE_SEGMENT_FINISHED. This two-phase pattern provides: - **Crash safety on copy**: a broker dying after STARTED but before FINISHED leaves a segment that is *not readable* and will be retried; the previously uploaded object becomes an **orphan** that GC/idempotent re-upload must handle. - **Crash safety on delete**: marking DELETE_STARTED first ensures the segment stops being served before its bytes vanish; FINISHED after deletion confirms cleanup. A crash mid-delete leaves a known-in-progress deletion to resume. ## Failure modes a principal must anticipate 1. **Eventual consistency of the object store**: read-after-write, list-after-write, and delete visibility can lag. Modern S3 is strongly read-after-write consistent for new objects, but list operations and some stores are not. Design must not depend on list consistency for correctness; rely on RLMM. 2. **Metadata/data divergence**: FINISHED but missing bytes → read failures; orphan bytes but no FINISHED → wasted storage but no correctness break. The safe direction is to risk orphans (cost) over missing-data (correctness). 3. **Leader failover mid-offload**: a new leader, seeing a STARTED-but-not-FINISHED segment via RLMM, re-uploads. Object keys derived from RemoteLogSegmentId keep this **idempotent** so re-upload overwrites rather than corrupts; the old partial object may be orphaned. 4. **Metadata store availability**: with the default TopicBasedRLMM, if the __remote_log_metadata topic is unavailable, offload/read of remote data stalls — the metadata topic's own replication factor and ISR settings become part of your durability story. 5. **Retention races**: remote retention (delete) must coordinate with reads; a segment being deleted should already be past retention, but in-flight reads against a just-deleted segment must fail cleanly (OffsetOutOfRange) rather than corrupt. 6. **Disaster recovery**: because metadata and data are decoupled, restoring a cluster requires both the object store and the metadata to be consistent; losing the metadata topic while keeping objects leaves you needing reconstruction. ## Design takeaways - Order operations so the **authoritative metadata transition trails durable side effects** (FINISHED after durable upload; DELETE_STARTED before byte deletion). - Treat the **RLMM as the single source of truth**; never the bucket listing. - Make object operations **idempotent** via stable IDs so retries and failovers are safe. - Budget for **orphan cleanup** as the acceptable cost of preferring correctness over storage efficiency.
- Why must the broker write COPY_SEGMENT_FINISHED only after the object-store upload is durable, and not before?Because FINISHED makes the segment readable. If FINISHED were written first and the upload then failed or wasn't yet durable, a read could be routed to bytes that don't exist, breaking correctness. The safe ordering risks at worst an orphan object, never missing data.
- Why should the broker trust the RLMM rather than listing the object-store bucket to decide what is readable?Object-store listings are often eventually consistent and may show orphaned partial uploads or lag behind deletes. The RLMM (replicated metadata) is the authoritative, ordered record of which segments are truly FINISHED and readable.
saying these in an interview costs you the question
- Saying the bucket listing is the source of truth for remote segments.
- Claiming you can mark a segment readable before its bytes are durable.
- Assuming object stores are always strongly consistent for all operations (lists/deletes often aren't).
- Ignoring orphaned objects from interrupted uploads / failovers.
- Treating metadata-topic availability as irrelevant to durability.