Write-Ahead Log, fsync & Group Commit
Part 0. Start here
The problem: 100,000 acknowledged writes a second
We go back to the storage node from the LSM-Trees & Compaction page. It takes 100,000 small writes a second, about 100 bytes each, on an EBS gp3 volume. The key-value store loop promises that "OK" means the write survives a crash. So before we reply, each write's log record must be synced: forced out of memory onto the disk itself.
How long does one sync take? AWS describes gp3 latency only as "single-digit millisecond". We assume about 1 ms per sync of the log. That is an assumption, and real values can be higher; Part 2 shows how to measure yours.
| What we need | What one sync per write gives |
|---|---|
| 100,000 acknowledged writes a second | One sync at a time, 1 ms each: 1 ÷ 1 ms = 1,000 a second, 1% of the need |
| "OK" survives power loss | RocksDB's default (sync = false) answers OK before any sync: fast, but a power loss can lose writes we said OK to |
Syncing every write on its own caps this node at 1,000 writes a second. Not syncing loses acknowledged writes in a power cut. How can the node acknowledge 100,000 writes a second and still never lose one it said OK to?
The big picture
A write-ahead log (WAL) is an append-only file that records every change before the change is applied anywhere else. Here is the whole write path this page takes apart:
Synthesizing vector architecture diagram...
What to notice: many writers, one write, one sync, many acknowledgements. The dotted branch is replication (Part 10): when the promise includes losing a machine or a zone, the OK also waits for other machines.
What you'll be able to do after this page
- Say exactly what an "OK" promises, and against which failure (Part 1).
- Explain what
fsyncand its relatives guarantee, and what a sync costs (Part 2). - Handle a failed
fsynccorrectly (Part 3). - Write a log record's format and the append procedure (Part 4).
- Size group commit: batch size, syncs per second, the wait (Part 5).
- Choose a durability mode for a kind of data (Part 6).
- Recover a log after a power loss that tore a write (Part 7).
- Decide when a log segment may be deleted (Part 8).
- Explain the WAL in a B-tree engine: PostgreSQL, InnoDB, SQLite (Part 9).
- Say what replication adds to durability, and trace one write through every layer (Part 10).
- Map all of it onto AWS (Part 11).
You may have arrived from step 1.1 of the key-value store loop ("WAL, fsync, group commit"), step 1.1 of the message queue loop (the broker's append-only log) or step 1.4 of the stock exchange matching engine loop (journal every input before releasing it). All three rest on this page.
Part 1. What "acknowledged" has to mean
"OK" is a promise, and a promise needs terms: survive what? Before any mechanism, we pin down the rule every engine follows and the failures it has to survive.
The WAL rule
The rule has two halves:
- Log the change before applying it. The record goes into the log before the MemTable, the page or anything else changes.
- Make the log durable before acknowledging. "Durable" means the record survives the failures the promise covers.
The LSM page (Part 2) already walks one write through this order: WAL (durable), then MemTable, then visible, then OK. This page is about the word "durable": what it costs and what it really guarantees.
Four ways to fail
| Failure class | What it destroys | Example |
|---|---|---|
| Process crash | The process's memory: its buffers, its MemTable | A bug, an out-of-memory kill, a deploy |
| OS crash or power loss | The kernel's memory, including the page cache (the kernel's cache of file bytes), and any volatile cache inside the drive | A kernel panic, a power cut to the rack |
| Device or node loss | Everything on that disk or machine | A dead drive, a failed EBS volume, a lost instance store |
| Zone loss | Every machine and volume in one Availability Zone | Power or network loss across an AZ |
A write that survives one class may not survive the next. An "OK" is only meaningful when we say which class it covers.
The durability chain
A record passes through several places on its way to safety. Each one survives a different set of failures:
Synthesizing vector architecture diagram...
What to notice: write() only moves a record from the "Application buffer" box to the "Page cache" box. Getting it to "Stable media" takes a sync, which also has to get through the "Device cache". Only "Other machines" survive losing the disk.
The tiny store, now watching the log
The LSM page followed a tiny store as its tree grew; this page watches its log. The settings are the same, plus the log's own:
| Setting | Tiny store | At real scale |
|---|---|---|
| Records | One record per write, (seq, op, key, value), so a torn tail is visible record by record. seq is the log sequence number (LSN) and also the entry's version in the LSM tree | LevelDB and RocksDB merge a whole commit group into one record, so a torn tail loses the group as a unit. PostgreSQL writes change records plus one commit record per transaction, and numbers positions by byte offset in the log, as InnoDB does |
| Segments (log files) | One per MemTable: wal-1 holds writes 1 to 4, ..., wal-6 writes 21 to 24 | RocksDB: a new WAL file at each MemTable switch. PostgreSQL: fixed 16 MB segments. Kafka: a new segment by size or time |
| Record header | Checksum (CRC32C) 4 bytes, length 2 bytes, type 1 byte | The same in LevelDB and RocksDB; RocksDB's recyclable format adds a 4-byte log number |
| Log blocks | Not modelled: every tiny record fits | 32 KB blocks; long records split into pieces (Part 4) |
| Sync latency L | 1.0 ms, and writing to the page cache takes no time (both are assumptions) | Device-dependent: measure it (Part 2). AWS publishes gp3 as "single-digit millisecond" and io2 Block Express as under 500 µs average for 16 KiB I/O |
| Durability mode | Group commit, sync before OK | Per engine (Part 6) |
| Device atomic unit | 32 bytes, a toy value: the device stores a batch unit by unit, may stop between units in a power loss, and may store the units in any order | 512-byte or 4 KiB sectors. AWS torn-write prevention offers atomic 4, 8 and 16 KiB writes on supported instances (Part 11) |
| Replicas (Part 10 only) | A (leader, AZ a), B (AZ b), C (AZ c); commit when 2 of 3 have the record | The same shape as the key-value and message queue loops |
The same 24 writes as the LSM page; Part 5 gives each an arrival time.
What to remember from Part 1
- "OK" is a promise about a failure class; name it.
- Log first, and make the log durable before saying OK.
- Each place in the chain survives a different failure: the page cache survives a process crash, stable media a power loss, other machines a lost disk or zone.
Part 2. What fsync really guarantees
Everyone says "we fsync the log". But there are several calls, each with a different promise, and a cache inside the drive that can quietly break all of them. This Part is what each call really buys, and what a sync costs.
The calls
| Call | What it does | Process crash | OS crash or power loss |
|---|---|---|---|
write() | Copies the bytes into the page cache and returns | Survives | Lost unless the kernel happened to write them back |
fsync(fd) | Writes the file's dirty data and metadata to the device and waits, "including writing through or flushing a disk cache if present" (Linux fsync(2)) | Survives | Survives |
fdatasync(fd) | Like fsync, but skips metadata not needed to read the data back (timestamps); still flushes a size change | Survives | Survives |
O_DSYNC (open flag) | Each write() returns only once done as if followed by fdatasync | Survives | Survives |
O_SYNC (open flag) | Each write() as if followed by fsync | Survives | Survives |
O_DIRECT (open flag) | Bypasses the page cache | Survives once written | Not guaranteed: it "does not give the guarantees of the O_SYNC flag" (open(2)) |
O_DIRECT is not durability. It skips the kernel's cache, but the drive's own cache can still hold the bytes, and the file's metadata may not be written. An engine using direct I/O still needs fdatasync (or O_DSYNC). PostgreSQL uses fdatasync for its WAL by default on Linux; RocksDB uses fdatasync where it is available.
The device's own cache
A drive with a volatile write cache can report a write as done while the bytes are only in its own RAM. So a sync has two jobs: get the pages to the device, then make the device put them on stable media.
Synthesizing vector architecture diagram...
What to notice: the kernel sends a flush (write everything in your cache to media) or uses FUA (force unit access: this write completes only once it is on non-volatile storage). If the device reports no volatile cache, the Linux block layer simply drops these requests. A power-loss-protected drive treats its cache as stable: it reports no volatile cache, or completes a flush almost at once; the sync still has to reach that cache, and no drive protects against its own loss or the zone's.
A cloud volume adds a network: on EBS, a sync becomes a round trip from the instance to the volume, which AWS replicates within its Availability Zone. AWS does not publish how EBS handles flushes internally, so we claim nothing more than its latency figures.
Creating and growing the log file
A sync of the file is not the whole story:
- A new file's name lives in its directory.
fsyncon the file "does not necessarily ensure that the entry in the directory containing the file has also reached disk" (fsync(2)). So after creating a segment (or renaming one), the engine syncs the directory too. Otherwise a power loss can leave a perfectly synced file that no directory lists. - A growing file changes its size. Each sync of an append that extends the file must also write the new size (metadata). So engines give the log its space up front: PostgreSQL zero-fills each new segment and recycles old ones; RocksDB can recycle log files. Reserving space with
fallocatealone is usually not enough: on common Linux file systems the reserved space starts as unwritten extents, and the first write into each one still has to update file-system metadata. RocksDB's documentation gives the same reason for recycling log files: then "fdatasync does not need to update the inode after each write". - The same is true of the LSM tree's flushed files; the LSM page's Part 3 syncs "the file and its directory" for the same reason.
Trace: the first write
When the tiny store opens, it creates wal-1, fills it with zeros, syncs it and syncs its directory. Then write 1 arrives at 0.0 ms. Only the steps that the LSM page's write trace (its Part 8) does not show are drawn here:
Synthesizing vector architecture diagram...
What to notice: the directory sync happens once per file, at creation; the data sync happens once per batch. Because wal-1 was zero-filled in advance, the sync at 1.0 ms has no file size to change.
What a sync costs
One log syncs one batch at a time, and a sync takes L. So:
| Sync latency L | Ceiling of one log, one write per sync |
|---|---|
| 0.2 ms | 1 ÷ 0.0002 = 5,000 a second |
| 1 ms | 1 ÷ 0.001 = 1,000 a second |
| 5 ms | 1 ÷ 0.005 = 200 a second |
The LSM page counted the log as sequential bytes; with a sync per batch, each sync is also one I/O against the IOPS budget. EBS counts an I/O of up to 256 KiB on SSD volumes as one, and "attempts to merge" small sequential I/Os into one, so one sync of a batch of up to 256 KiB costs about one I/O.
One sync per write, L = 1 ms. What is the ceiling of one log? What if the volume syncs in 0.2 ms, or 5 ms?
Measure, don't assume. Run a test with a sync after every write on the actual volume and instance type, and read the engine's own sync-latency metric in production. Three literal command lines:
textfio --name=walsync --filename=/data/fio.test --size=256m --rw=write --bs=4k --ioengine=sync --fdatasync=1 pg_test_fsync -f /data/pg_test_fsync.out cat /sys/block/nvme1n1/queue/write_cache
fio with --fdatasync=1 syncs after every write and reports the latency distribution; pg_test_fsync times each sync method PostgreSQL can use; the last line shows whether the kernel sees a volatile cache on that device (write back) or not (write through).
Go deeper: layers that ignore flushes. Every guarantee above assumes each layer honours the flush: the file system passes it on, the virtualization layer passes it on, the drive obeys it. A layer that acknowledges a flush without doing it (a misconfigured hypervisor, a RAID card with a failed battery still in write-back mode, a consumer drive that lies) breaks every promise above it silently, and nothing shows until a power loss. Treat it as a risk to test for: pull the power (or the virtual equivalent) under load on the exact stack you run, then check that every acknowledged write is there.
What to remember from Part 2
write()reaches the page cache; only a sync (fsync,fdatasync,O_DSYNC) reaches stable media, through the device cache.O_DIRECTalone is not a sync.- New files need their directory synced; growing files need their size made durable, so logs are zero-filled in advance or recycled.
- One log, one sync at a time: at most 1 ÷ L syncs a second, and each one is an I/O.
Part 3. When fsync fails
A sync can return an error: an I/O error (EIO) from a failing device or a lost network volume, or ENOSPC when the disk is full. The natural reaction, "try again", is exactly the wrong one. This Part shows why, with two branches of the tiny store.
What a failed sync means
When the kernel fails to write dirty pages back, Linux may mark those pages clean anyway and report the error once, to the next sync that asks. The bytes may never have reached the disk, but the kernel no longer considers them waiting to be written. So:
- the failed sync means the data may already be gone, and
- a retried sync finds nothing dirty and returns success.
This behaviour hit PostgreSQL in 2018 (the incident is often called "fsyncgate"). PostgreSQL's fix was to stop retrying: by default it now raises a PANIC on a failed data sync (data_sync_retry = off), crashes and recovers from its WAL. It was not alone: a 2020 USENIX study of five widely used applications (PostgreSQL, LMDB, LevelDB, SQLite and Redis) found that none of their failure-handling strategies was sufficient, and that a failed fsync could still cause data loss or corruption.
Trace: the log's sync fails (branch E1)
Rewind to batch B11 (writes 18, 19 and 20, Part 5). Their records are in wal-5's page cache, and the fdatasync of wal-5 returns EIO. The reference implementation's kernel marks the pages clean without writing them.
| Wrong path: retry | Right path: fatal | |
|---|---|---|
| What the engine does | Retries the fdatasync: nothing is dirty, it returns OK | Treats the error as fatal: fails every writer waiting on this sync, stops taking writes, restarts |
| What clients 18 to 20 get | OK at 11.0 ms | An error |
What wal-5 holds on disk | Only record 17 | Only record 17 |
| Power fails before F11 (the flushed file for writes 17 to 20) is written | Recovery finds 17 only: three acknowledged writes lost | Recovery finds 17; next seq 18; nothing that was acknowledged is missing |
Why "before F11": on the wrong path, the MemTable still holds 18 to 20, so if the flush writing F11 completes first, their data reaches disk that way. The retry turns a crash that loses nothing into a timing lottery.
Trace: the flush's sync fails (branch E2)
This is the classic case. After write 20, the flush writes F11 and calls fsync on it, and the fsync fails:
Synthesizing vector architecture diagram...
What to notice: every step after the retry follows the LSM page's safe flush order (its Part 3), and it still loses data, because the order assumed the sync told the truth. F11's blocks never reached the disk, and wal-5, the only other copy of writes 17 to 20, is gone. Nothing complains until a read or a compaction touches F11, which could be days later.
F11's fsync fails. Should the engine retry it?
The rule
- A failed sync is fatal. Never retry it and carry on.
- Fail every write waiting on that sync. None of them was acknowledged.
- Stop taking writes, restart, and recover from what is on disk.
- A client that got an error does not know the outcome. The bytes may or may not be on the disk, and after recovery the write may or may not exist. So an error means "unknown", not "failed", and the client's retry must be idempotent: safe to apply twice (the Idempotency & Effectively-Once Processing loop primitive covers how).
Go deeper: the page cache can still hold what the disk lacks. After a write-back error the kernel marks the pages clean, but it may keep their contents in memory. A process that restarts without a reboot and reads the log through the page cache can then see records the disk does not have: in branch E1, the reference implementation's restarted process reads records 17 to 20 back from the cache, although only 17 is on the disk. It would replay them as if they were durable, and the next power loss would take them away. The 2020 USENIX study found exactly this on ext4 and XFS: after a failed fsync the newest data stays in the page cache, so applications "seemingly work correctly" until the pages are evicted or the machine restarts. Reading the log with direct I/O, or recovering after a reboot, avoids being fooled.
What to remember from Part 3
- A failed sync means the data may already be gone, and a retry can report success anyway.
- Never retry and carry on: stop, restart and recover from what is on disk.
- An error is an unknown outcome, so client retries must be idempotent.
Part 4. The log record and the append
Recovery reads the log back after a crash, possibly with half a record at the end. So every record has to describe itself: how long it is, and whether its bytes are the ones that were written.
What a record holds
| Field | Size | Why |
|---|---|---|
| Checksum | 4 bytes | CRC32C of the type byte and the payload (LevelDB and RocksDB store it masked, a small fixed transformation; RocksDB's recyclable format also covers the log number). The length field is checked by bounds, not by the checksum |
| Length | 2 bytes | Where the record ends, and so where the next begins |
| Type | 1 byte | FULL for a whole record; FIRST, MIDDLE, LAST for the pieces of a long one (Go deeper) |
Payload: seq | 8 bytes | The sequence number, also the entry's version in the LSM tree |
Payload: op | 1 byte | PUT or DEL |
| Payload: key | 1-byte length + the key | user:1 |
| Payload: value | 1-byte length + the value | alice; empty for a delete |
The header layout is LevelDB's and RocksDB's. The payload layout is the tiny store's own choice (real engines use variable-length integers and put a whole batch in one payload).
That last point matters: LevelDB and RocksDB write a whole commit group as one record. The leader of a group commit (Part 5) merges its group's write batches and writes them as a single record (LevelDB's BuildBatchGroup, RocksDB's MergeBatch), then syncs once. The tiny store writes one record per write, so a torn write shows up record by record; in those engines, a torn group is lost or kept as a unit.
The sequence number has two forms in real engines: RocksDB counts entries (a batch carries the number of its first entry, and the entries in it follow on), while PostgreSQL and InnoDB number log positions by byte offset in the log. Either way, it only goes up, so it orders every change.
Trace: wal-1after write 4
The reference implementation's bytes for the first segment (seq as 8 bytes little-endian, PUT = 1, DEL = 0, FULL = 1; checksums shown unmasked):
| seq | Operation | Offset | Header | Payload bytes | Record bytes | CRC32C |
|---|---|---|---|---|---|---|
| 1 | PUT user:1 alice | 0 | length 22, FULL | 22 | 29 | 0x115b1cef |
| 2 | PUT user:2 bob | 29 | length 20, FULL | 20 | 27 | 0x72ff3437 |
| 3 | PUT order:9 book | 56 | length 22, FULL | 22 | 29 | 0xd3a5555c |
| 4 | PUT user:3 carol | 85 | length 22, FULL | 22 | 29 | 0xc103ba8b |
Each record starts where the previous one ends: 0 + 29 = 29, 29 + 27 = 56, 56 + 29 = 85, and wal-1 holds 85 + 29 = 114 bytes. The header is 7 bytes on every record, so on 100-byte writes the log is about 7% larger than the data.
The append
textAPPEND(op, key, value) -- run by every client write 1. seq = next sequence number 2. payload = (seq, op, key, value) record = checksum(type + payload), length, type, payload 3. join the group-commit queue with record; wait until it is durable -- Part 5 if the sync failed: return an error (outcome unknown) -- Part 3 4. apply: insert into the MemTable (or change the page) 5. make seq visible to readers 6. acknowledge the client
The order of steps 3 and 6 is the whole promise. Compare the right order with the one that produces the "acknowledged but lost" bug:
| Right: log, sync, apply, OK | Wrong: apply, OK, then sync |
|---|---|
| Record written to the log | Change applied in memory |
| Sync returns: the record is on stable media | OK sent to the client |
| Change applied, made visible | Record written to the log |
| OK sent to the client | Sync starts... power fails |
| A crash at any point loses only writes nobody was told succeeded | The client was told OK; the write is gone |
Go deeper: blocks, fragments and recycled files. LevelDB and RocksDB split the log into 32 KB blocks. A record never starts in a block's last 6 bytes (too little room for a 7-byte header), so those bytes are filled with zeros. A record too long for the rest of a block is cut into pieces typed FIRST, MIDDLE and LAST, each with its own header and checksum, and the reader glues them back together. Blocks let a reader resynchronise at the next 32 KB boundary after damage. RocksDB's recyclable record format adds a 4-byte log number to every header, so that the old records still sitting in a reused file are not mistaken for new ones (Part 7).
What to remember from Part 4
- Every record carries its own length and checksum, so recovery can tell a whole record from garbage.
- Append, sync, apply, then acknowledge. Acknowledging before the sync is the "acknowledged but lost" bug.
- The sequence number is both the log position and the version.
Part 5. Group commit
One sync per write caps a log at 1 ÷ L writes a second, however fast the rest of the engine is. Group commit removes that cap without weakening the promise. This Part is the algorithm, the 24-write trace and the arithmetic.
One sync, many writers
A sync of the log file makes every byte written to the file before it durable. So if three writers have each written a record, one sync covers all three. The only question is how to gather them. The naive answer is a timer ("a background thread syncs every millisecond"); real engines such as LevelDB, RocksDB and PostgreSQL do something simpler with no timer at all. PostgreSQL's source describes it: a backend waiting for the WAL write lock rechecks "if the backend that held the lock did it for us already".
The algorithm
textGROUP_COMMIT_WRITE(record) 1. lock the queue; append record 2. while a sync is running and record is not durable: wait 3. if record is durable: unlock; return OK -- a leader carried it 4. become the leader: batch = the whole queue; empty the queue; unlock 5. one write of the whole batch to the log file 6. one fdatasync of the log file 7. lock; mark every record in batch durable (or failed); wake every waiter -- a waiter whose record is still queued becomes the next leader 8. unlock; return OK (or the error)
Batches form by themselves: the next batch is whoever queued during the last sync. At low load a writer finds nothing running and leads a batch of one, with no added wait. At high load the queue fills while each sync runs, and batches grow. Tie rule (the tiny store's choice, needed because its times are exact): a write that arrives at the very instant a sync completes joins the next batch.
Batch B4 of the tiny store, as the reference implementation ran it:
Synthesizing vector architecture diagram...
What to notice: writers 5, 6 and 7 arrived during B3's sync, so they became one batch; writer 5 was first in the queue and led it. One sync, three acknowledgements, all at 4.0 ms.
Trace: 24 writes, 14 syncs
The same 24 writes as the LSM page, now with the time each reaches the engine. Sync latency L = 1.0 ms. Every result below comes from the reference implementation:
| seq | Operation | Arrives (ms) | Batch (sync, ms) | OK at (ms) | Waited (ms) |
|---|---|---|---|---|---|
| 1 | PUT user:1 alice | 0.0 | B1 (0 to 1) | 1.0 | 1.0 |
| 2 | PUT user:2 bob | 0.3 | B2 (1 to 2) | 2.0 | 1.7 |
| 3 | PUT order:9 book | 0.6 | B2 | 2.0 | 1.4 |
| 4 | PUT user:3 carol | 1.1 | B3 (2 to 3) | 3.0 | 1.9 |
| 5 | PUT user:1 alicia | 2.4 | B4 (3 to 4) | 4.0 | 1.6 |
| 6 | DELETE user:2 | 2.5 | B4 | 4.0 | 1.5 |
| 7 | PUT order:7 pen | 2.9 | B4 | 4.0 | 1.1 |
| 8 | PUT user:4 dave | 3.3 | B5 (4 to 5) | 5.0 | 1.7 |
| 9 | PUT user:5 erin | 4.6 | B6 (5 to 6) | 6.0 | 1.4 |
| 10 | DELETE user:1 | 4.8 | B6 | 6.0 | 1.2 |
| 11 | PUT order:9 lamp | 5.0 | B6 (tie rule) | 6.0 | 1.0 |
| 12 | PUT user:2 bea | 5.2 | B7 (6 to 7) | 7.0 | 1.8 |
| 13 | PUT user:6 finn | 6.8 | B8 (7 to 8) | 8.0 | 1.2 |
| 14 | PUT user:3 cara | 7.0 | B8 (tie rule) | 8.0 | 1.0 |
| 15 | DELETE order:7 | 7.4 | B9 (8 to 9) | 9.0 | 1.6 |
| 16 | PUT user:7 gus | 7.9 | B9 | 9.0 | 1.1 |
| 17 | PUT order:9 vase | 9.0 | B10 (9 to 10, tie rule) | 10.0 | 1.0 |
| 18 | PUT user:8 hana | 9.2 | B11 (10 to 11) | 11.0 | 1.8 |
| 19 | PUT user:4 dan | 9.5 | B11 | 11.0 | 1.5 |
| 20 | DELETE user:5 | 9.9 | B11 | 11.0 | 1.1 |
| 21 | PUT user:1 ali | 11.0 | B12 (11 to 12, tie rule) | 12.0 | 1.0 |
| 22 | PUT order:1 mug | 11.2 | B13 (12 to 13) | 13.0 | 1.8 |
| 23 | PUT user:9 ivy | 12.5 | B14 (13 to 14) | 14.0 | 1.5 |
| 24 | PUT user:6 fay | 12.7 | B14 | 14.0 | 1.3 |
14 syncs for 24 writes. The average wait is 33.2 ÷ 24 ≈ 1.4 ms, the worst 1.9 ms (write 4), the best 1.0 ms: always between L and 2L. No batch crosses a segment boundary (B1 to B3 are all in wal-1, and so on); a new segment starts only between batches.
The same arrivals with one sync per write: 24 writes arrive in 12.7 ms, about 24 ÷ 0.0127 ≈ 1,900 a second, above the 1,000-a-second ceiling. So the queue grows: write n is acknowledged at n ms, the last at 24.0 ms. By 12.7 ms, 12 writes are waiting. Write 24 waits 24.0 − 12.7 = 11.3 ms, and the average wait is 6.4 ms.
Synthesizing vector architecture diagram...
What to notice: both rows sync back to back, one sync per millisecond: the disk is equally busy. The "Group commit" row finishes all 24 writes at 14 ms because several of its syncs carry two or three writes; the "One per write" row needs 24 ms, and every write that arrives after 1 ms waits in a queue that keeps growing until the arrivals stop at 12.7 ms.
The math
The rule of thumb for the batch size:
Worked for the hook: 100,000 writes a second × 0.001 s = 100 writes per batch, so 100,000 ÷ 100 = 1,000 syncs a second, each carrying about 10 KB. A writer waits for the rest of the sync in progress (between 0 and L, L ÷ 2 on average) plus its own sync (L): between L and 2L, about 1.5 ms on average. The reference implementation's steady-load run (100,000 writes a second, both evenly spaced and random arrivals) gives exactly that: 1,000 syncs a second, 100 writes per batch, 1.50 ms average wait, 2.0 ms worst.
Three consequences:
- It needs concurrency. Writes in flight = 100,000 × 0.0015 ≈ 150. One client writing one record at a time never gets a batch bigger than one.
- It adapts by itself. At low load, batches of one and a wait of exactly L; at high load, bigger batches and a wait approaching 2L, while syncs a second stay at about 1 ÷ L.
- It stops helping only at the bandwidth limit, when a batch takes longer to write than to sync; above 256 KiB, one sync also costs EBS more than one I/O.
A log takes 50,000 writes a second, and one sync takes 0.5 ms. With leader-based group commit, how big is a batch, how many syncs a second does the log do, and how long does a writer wait?
The window variant
Instead of syncing as soon as the previous sync ends, some designs wait on purpose: collect writes until a window W has passed or the batch reaches a size, then write and sync once. The message queue loop uses "2 ms or 64 KB, whichever comes first". Real examples: Cassandra's commitlog_sync: group (acknowledges after a sync of each group window) and PostgreSQL's commit_delay, which waits before flushing only when at least commit_siblings other transactions are active (defaults: 0 µs and 5, so off by default).
Both on equal terms, on the same inputs, from the reference implementation (the window opens when the first write reaches an empty batch; a write arriving exactly as a window closes joins the next one):
| Leader-based | Window: 2 ms or 64 KB | |
|---|---|---|
| Tiny store, 24 writes: syncs | 14 | 6 (each carries one whole segment: 1 to 4, 5 to 8, ...) |
| Tiny store: average wait | 1.4 ms | 2.5 ms |
| Tiny store: worst wait | 1.9 ms (write 4) | 3.0 ms (writes 1, 5, 9, 13, 17 and 21: each opens its window) |
| 100,000 writes a second: syncs a second | 1,000 | 500 |
| 100,000 writes a second: average (worst) wait | 1.5 ms (2.0) | 2.0 ms (3.0) |
| Best when | Latency matters; the sync rate fits the IOPS budget | Syncs are the bottleneck (a shared journal, a tight IOPS budget), and a few milliseconds of latency are acceptable |
The window halves the syncs at high load and costs up to W of extra latency at every load, even when the disk is idle.
"Group commit" is also used loosely for any batching of durable writes; the distributed email loop "group-commits" messages into an S3 spool. The idea is the same: one expensive durable step shared by many writes.
Go deeper: pipelining and one shared journal. RocksDB's optional pipelined writes (enable_pipelined_write, off by default) let the next group start its WAL write while the previous group is still inserting into the MemTable, so the log and the MemTable work in parallel. At the other extreme, a broker with many partitions can give them one shared journal, so one sync covers every partition's writes instead of one sync per partition log; the message queue loop does this in Round 2 (step R2.5). The whole broker then shares one sync budget.
What to remember from Part 5
- One sync covers every record written before it.
- Batches form by themselves: whoever queued during the last sync. No timer.
- Throughput needs concurrency; latency grows to at most about 2 syncs. A window trades more latency for fewer syncs.
Part 6. Durability modes: what "acknowledged" means
Not all data deserves a sync before OK. Metrics can lose a second; a payment cannot. Every engine lets you choose, and the classic mistake is to choose a mode without stating what it can lose.
The ladder
| Mode | Process crash | OS crash or power loss | Device or node loss | Zone loss |
|---|---|---|---|---|
| Sync per write | Loses nothing acknowledged | Loses nothing acknowledged | Loses everything on that disk | Loses everything in the zone |
| Group commit (sync before OK) | Same as sync per write: the same guarantee, a lower price | Same | Same | Same |
| Periodic sync every T | Loses nothing acknowledged, if each write reaches the OS before the OK | Up to about T of acknowledged writes | Everything on that disk | Everything in the zone |
| No sync (the OS writes back on its own schedule) | Loses nothing acknowledged | Whatever the OS had not written back (on Linux, typically up to about 30 s) | Everything on that disk | Everything in the zone |
| No WAL | Everything since the last flush | Everything since the last flush | Everything on that disk | Everything in the zone |
| Replicated (Part 10) | Depends on each copy's own mode | Survives if a quorum's copies are synced, or their failures are independent | Survives if the quorum spans machines | Survives if the quorum spans zones |
Trace: one power loss, four modes
The same 24 arrivals, and power fails at the same wall-clock instant, 13.5 ms, under four modes. Flushed files are synced in every mode (the LSM page's flush order). Two things the reference implementation derives rather than assumes: under the no-sync modes, the flush of writes 21 to 24 (F12) starts when write 24 arrives at 12.7 ms and needs its own 1 ms file sync, so it cannot be durable before 13.7 ms; and we assume the kernel has not written wal-6's pages back by 13.5 ms.
| Mode | Acknowledged by 13.5 ms | Acknowledged and lost |
|---|---|---|
| (a) No sync: OK as soon as the record is written | 1 to 24 | 21 to 24 (only in wal-6's page cache; 1 to 20 are safe in synced files) |
| (b) Periodic sync every 10 s: the last sync ran before write 1 | 1 to 24 | 21 to 24, the same as (a): no periodic sync happened in the window |
| (c) Sync per write | 1 to 13 (write 13 at 13.0 ms; write 14's sync is running) | None: 11 writes (14 to 24) are still waiting, none of them told OK |
| (d) Group commit | 1 to 22 (B13 ended at 13.0 ms; B14 is running) | None: 23 and 24 were never told OK |
Group commit acknowledged 22 writes by the instant of the power loss and lost none of them; sync per write acknowledged only 13; the fast modes acknowledged all 24 and lost four of them.
How engines name it
| Engine | Setting | Values (what each means) | Default |
|---|---|---|---|
| RocksDB | WriteOptions.sync, disableWAL, manual_wal_flush | sync = true syncs before OK; false writes to the OS only; disableWAL skips the log; manual_wal_flush keeps records in the engine until FlushWAL() | sync = false |
| PostgreSQL | synchronous_commit | on, local (wait for the local WAL flush), remote_write, remote_apply (Part 10), off (don't wait; can lose up to three times wal_writer_delay of commits, never corrupts). With no synchronous standby configured, local, remote_write and remote_apply all give the same local guarantee as on | on |
| MySQL InnoDB | innodb_flush_log_at_trx_commit | 1: write and flush the log at each commit. 2: write at each commit, flush about once a second. 0: write and flush about once a second. "Once a second" is not guaranteed | 1 |
| Cassandra | commitlog_sync | periodic (OK before the sync; sync every commitlog_sync_period), batch (OK after the sync), group (OK after a sync per group window) | periodic, 10,000 ms |
| Redis | appendfsync | always (sync before the replies, with group commit across parallel writes), everysec, no (the OS decides: "normally" every 30 s on Linux) | everysec; the append-only file itself is off by default (appendonly no) |
| OpenSearch | index.translog.durability | request (sync on the primary and every replica before OK), async (sync every index.translog.sync_interval, 5 s by default) | request |
| Kafka | producer acks, topic min.insync.replicas, flush.messages | Durability by replication (Part 10); the flush settings default to never forcing a sync, and the docs recommend leaving them that way | acks=all (current clients), min.insync.replicas=1 |
| SQLite (WAL mode) | PRAGMA synchronous | FULL syncs the WAL at each commit; NORMAL syncs the WAL only before each checkpoint (Part 9) | Depends on how SQLite was built and on the platform: check it, don't assume |
With InnoDB's 2, the log reaches the OS at every commit, so by Part 2's rules a crash of the MySQL process alone should not lose it; an OS crash or power loss can lose up to about the last second.
Losing writes vs corrupting data
PostgreSQL shows the difference in two settings:
synchronous_commit = offacknowledges before the WAL flush. A crash can lose the last few commits (up to three timeswal_writer_delay), but the database comes back consistent: the WAL is still written in order and still protects every page.fsync = offstops syncing altogether, including the data pages and checkpoints. A crash can leave pages and log out of step, and the database may be corrupt, not just behind. PostgreSQL's documentation puts it plainly:fsync = off"can result in unrecoverable data corruption";synchronous_commit = off"does not create any risk of database inconsistency".
Losing the last second of metrics can be a sound business choice. Corrupting the database never is.
Pick a durability mode for (1) a payments ledger, (2) a metrics pipeline taking a million samples a second, (3) a session cache that can be rebuilt from logins.
What to remember from Part 6
- A mode is a loss window per failure class; state both.
- Group commit keeps "sync before OK" and only changes the price.
- Losing recent writes can be a choice; corrupting data never is.
Part 7. Power loss and torn writes
A sync writes many bytes, and a power loss can stop it anywhere. Recovery must find exactly the records that are whole, and nothing else.
What a torn write is
A device persists data in atomic units (sectors): each unit is written completely or not at all. A write larger than one unit is not atomic. If power fails in the middle, some units have the new bytes and some the old, and the units may have reached the media in any order. That is a torn write.
The tiny store uses 32-byte units so the effect shows up in a small example. At real scale the unit is 512 bytes or 4 KiB; two small records then usually share one unit and cannot tear apart, and a torn tail is cut at unit boundaries, not record boundaries.
Trace: power fails during B14 (branch P)
wal-6 holds records 21 to 24. The reference implementation places them at these byte ranges: record 21 at 0 to 27, record 22 at 27 to 55, record 23 at 55 to 82, record 24 at 82 to 109. Units 0 and 1 were already on the media after B13. B14 writes records 23 and 24, which dirties units 1, 2 and 3 (unit 1 is rewritten whole, record 22's tail included), and the fdatasync is running when the power fails:
Synthesizing vector architecture diagram...
Case 1 of branch P: units 1 and 2 reached the media, unit 3 did not. Record 23 is complete on disk; record 24 has its header and its first 7 payload bytes, then zeros. Because each unit is atomic, rewriting unit 1 could not damage record 22's tail: that unit holds either its old bytes or its new ones.
The recovery scan
textRECOVER_SCAN(segments still needed, oldest first) for each segment: offset = 0 loop: 1. read 7 header bytes at offset if fewer than 7, or all zero: stop -- end of what was ever written 2. read length payload bytes if fewer than length: stop -- torn 3. if checksum(type + payload) != the stored checksum: stop -- torn or damaged 4. if seq is already in a flushed file: skip it; else apply it -- safe to repeat 5. offset = offset + 7 + length "stop" ends the whole replay, not just this segment cut the log after the last good record; next seq = last good seq + 1
Synthesizing vector architecture diagram...
What to notice: there is only one way forward and three ways to stop. The scan never skips a bad record to look for good ones after it.
This is the log side of the LSM page's boot sequence (its Part 9), which surrounds this scan with reading the manifest and flushing the rebuilt MemTable; the LSM page's Part 10 lists RocksDB's recovery modes for damage found here.
| Branch P | What reached the disk | The scan | Acknowledged? |
|---|---|---|---|
| Case 1: power fails during B14's sync | Units 1 and 2; not unit 3 | Keeps 21, 22, 23; record 24's checksum fails at offset 82, so the log is cut there; next seq 24 | 21 and 22 yes; 23 and 24 no. Write 23 survived although its client never got an OK |
| Case 2: B14's sync returns at 14.0 ms, power fails before the OKs are sent | Everything | Keeps 21 to 24; stops at offset 109 on an all-zero header (never written); next seq 25 | 21 and 22 yes; 23 and 24 no, yet both survived |
Case 1 depends on the tiny store's one-record-per-write format; in LevelDB or RocksDB, 23 and 24 would be one record and would be lost or kept together. Case 2 happens in every engine: the sync can finish a moment before the power goes, and nobody is told.
Writes that survived without their OK
In a sync-before-OK mode, the torn tail is always safe to cut: nothing in it was acknowledged. The surprise runs the other way: a write can survive without its OK. Clients 23 and 24 saw a timeout or a dropped connection, not an error, and their writes may exist after recovery (both do in case 2). So a timeout is an unknown outcome, like the error in Part 3. The client must retry with an idempotency key, so the write applies once whether or not the first attempt survived (the Idempotency & Effectively-Once Processing loop primitive).
Branch P2: at the same instant, the device had stored units 2 and 3 but not unit 1. So record 24 is complete on the disk with a good checksum, but record 23 is missing its first 9 bytes. Should recovery replay record 24?
Go deeper: recycled segments. An engine that reuses old log files (to avoid growing a file, Part 2) leaves old records in them. After a crash, the bytes just past the new tail may be an old record with a perfectly valid checksum, which the scan would happily apply. The fix is to put the log file's number in each record (RocksDB's recyclable format, Part 4), or a position the old record cannot match (PostgreSQL's WAL pages carry their own log address in a header), so the scan stops at the first record from a previous life.
What to remember from Part 7
- A torn tail fails its length or checksum test and is cut; in sync-before-OK modes it was never acknowledged.
- Replay stops at the first bad record, even if later records look fine.
- A write can survive without its OK: treat timeouts as unknown and make retries idempotent.
Part 8. Segments, rotation and truncation
The log records every write forever unless something deletes it. It cannot simply be deleted by age: a segment may still hold the only durable copy of a write, or data a reader has not consumed yet.
Why segments
One endless file cannot be shortened from the front. So the log is a series of segments, separate files, and old ones are deleted, recycled or archived whole. A new segment starts:
- at a MemTable switch (RocksDB, and the tiny store: one segment per MemTable, so a segment becomes useless exactly when its MemTable is flushed);
- at a size (PostgreSQL's 16 MB segments); or
- at a size or a time (Kafka).
When a segment may go
A segment may be deleted only when both are true:
- Every record in it is durable elsewhere: in a flushed file (LSM), or in data pages written back by a checkpoint (B-tree, Part 9).
- No consumer still needs it: every replica, backup shipper and change-data-capture reader has confirmed a position past its last record.
The LSM page's Part 3 gives the order for the first condition: sync the flushed file and its directory, record it in the manifest, and only then delete the segment. Its crash-at-each-step table shows why that order loses nothing.
Synthesizing vector architecture diagram...
What to notice: a sealed segment has two gates to pass, and the slowest one decides. A flush can release it in milliseconds; a stuck reader can hold it forever.
Trace: after write 24, with and without a stuck reader (branch T)
Fork the main timeline after write 24 is acknowledged and before F12's manifest edit. Manifest edit 9 (F11) has made writes 1 to 20 durable in flushed files, so wal-1 to wal-5 pass the first gate and wal-6 does not. Now add a change-data-capture reader that has confirmed only up to seq 10:
| Segment | Writes | Durable elsewhere? | No reader: | Reader confirmed up to seq 10: |
|---|---|---|---|---|
wal-1 | 1 to 4 | Yes | Deletable | Deletable |
wal-2 | 5 to 8 | Yes | Deletable | Deletable |
wal-3 | 9 to 12 | Yes | Deletable | Kept: the reader still needs 11 and 12 |
wal-4 | 13 to 16 | Yes | Deletable | Kept |
wal-5 | 17 to 20 | Yes | Deletable | Kept |
wal-6 | 21 to 24 | No (F12 not recorded yet) | Kept | Kept |
Without the reader, one segment is kept; with it, four. Every new write adds to the pile until the reader moves or the disk fills.
The log as a change stream
Replicas, backup shippers and change-data-capture tools all read the log, and each one pins it:
- Replication: a follower that falls behind needs every segment since its position (or must be rebuilt from a snapshot).
- Backup shipping: the key-value store loop ships sealed segments to S3 every 60 s; a segment may go only after its upload is confirmed.
- Change data capture: PostgreSQL's replication slots make the server keep WAL until the slot's consumer confirms it. The transactional outbox loop (step R2.7, "a stalled slot retains WAL") and the payments loop (step R3.9) alarm on exactly this. The Change Streams & the Transactional Outbox loop primitive covers the readers themselves.
- RocksDB can archive WAL files for such readers instead of deleting them, bounded by
WAL_ttl_secondsandWAL_size_limit_MB(both 0 by default: no archive).
A PostgreSQL replication slot's consumer has been stuck for 6 hours. What is happening to the writer, and what do you do?
What to remember from Part 8
- A segment goes only when its data is durable elsewhere and no reader needs it.
- The slowest consumer sets the log's size.
- Alarm on consumer lag before disk space.
Part 9. The WAL in B-tree engines
PostgreSQL, InnoDB and SQLite are not LSM trees: they keep data in fixed-size pages (8 KB in PostgreSQL, 16 KB in InnoDB) and change them in place. They log first too, but a page adds two problems: replay must be safe to repeat, and a half-written page cannot be repaired from a small log record.
Log first in a page-based engine
A change modifies a page in memory, logs a record describing it, and the dirty page is written back to its place in the data file later, possibly much later. The WAL rule adds one constraint: a page may be written back only after the log is durable up to that page's last change. The log's commit record is what makes a transaction durable; the page write-back is just catching up.
The words to know: this is ARIES-style recovery, a steal / no-force design. "No-force": commit does not wait for the pages to be written. "Steal": a page may be written back before its transaction commits (so the log also needs undo information; that is beyond this page).
pageLSN makes redo safe to repeat
Each page stores the LSN of the last log record applied to it, its pageLSN. Redo applies a record to a page only if the record's LSN is greater than the pageLSN; otherwise that change is already on the page. That makes replay idempotent: a crash during recovery just runs it again.
textREDO(from the checkpoint's redo point) for each log record r, in LSN order: page = read r's page from the data file if r carries a full-page image: page = the image; page.LSN = r.LSN -- the image already includes r's change else if r.LSN <= page.LSN: skip -- already on the page else: apply r; page.LSN = r.LSN
Checkpoints bound replay
A checkpoint writes dirty pages back and records a redo point: every change before it is on the data pages. Recovery starts there, not at the beginning of the log, and segments before it can go (Part 8's first gate). Checkpoints trade steady background writes for shorter recovery.
Torn pages
A page is 8 or 16 KB, much larger than the device's 512-byte or 4 KiB atomic unit, so a power loss during write-back can leave a page half old and half new. A log record says something like "put order:7 pen in this page"; applied to a page that is half garbage, it cannot rebuild the rest. The page needs a whole good copy from somewhere:
Synthesizing vector architecture diagram...
What to notice: in "Full-page image" the good copy lives in the log (PostgreSQL, full_page_writes = on by default); in "Doublewrite" it lives in a separate area written before the page (InnoDB, on by default); with "Atomic device write" no copy is needed, which is what AWS's torn-write prevention and RDS Optimized Writes provide (Part 11).
Full-page images and doublewrite are the B-tree's extra write cost; the LSM page's Part 7 compares the totals.
Trace: the tiny store on two pages (branch B)
Put writes 1 to 8 on a B-tree with two leaf pages: P1 holds the order: keys, P2 the user: keys. A checkpoint after write 4 writes both pages (P1 pageLSN 3, P2 pageLSN 4); redo will start at LSN 5. Write 5 is P2's first change after the checkpoint, so its record carries P2's full-page image; write 7 is P1's first, so it carries P1's image. P2 is written back after write 6 (pageLSN 6). After write 8, P1's write-back is torn by a power loss.
Synthesizing vector architecture diagram...
What to notice: record 7's image is the only thing that can rebuild P1; record 8 is newer than P2's pageLSN, so it must be applied.
The reference implementation runs recovery both ways. Both end with P1 = order:7 pen, order:9 book (pageLSN 7) and P2 = user:1 alicia, user:3 carol, user:4 dave (pageLSN 8), exactly the state before the power loss:
| Record | With full-page images (PostgreSQL style) | With doublewrite instead (InnoDB style) |
|---|---|---|
| Before redo | P1 torn | P1 torn: its page checksum fails, so it is replaced by its doublewrite copy (pageLSN 7) |
| 5 (P2) | Carries P2's image: restored, P2 pageLSN 5. PostgreSQL restores an image "even if the page in the database appears newer", because it cannot tell a torn page from a good one | Skipped: 5 ≤ pageLSN 6 |
| 6 (P2) | Applied: 6 > 5 | Skipped: 6 ≤ 6 |
| 7 (P1) | Carries P1's image: restored, P1 pageLSN 7 | Skipped: 7 ≤ 7 |
| 8 (P2) | Applied: 8 > 6 | Applied: 8 > 6 |
Without either protection, P1 is lost: record 7 only says "put order:7 pen", and order:9 book was written before the checkpoint, so nothing in the log since the redo point holds it.
A 16 KB page is half written when power fails. Why can't the small redo record for the last change fix it?
SQLite's WAL mode
The mobile news feed and mobile chat loops run SQLite in WAL mode. Its log holds whole pages, not change records:
| Behaviour | What it means for a design |
|---|---|
A commit appends the changed pages to the -wal file; the database file is untouched | Writes are sequential appends; a second file, -shm, indexes the log |
| A reader sees the database plus the log frames committed when its read began | Readers never block the writer, and the writer never blocks readers |
| A checkpoint copies frames back into the database file; by default automatically once the log reaches 1,000 pages | The -wal file's size is bounded by checkpoints |
| A checkpoint cannot copy past a frame an older reader may still need | A long-lived reader lets -wal grow (the news feed loop's step R2.8) |
synchronous = FULL syncs the log at each commit; NORMAL syncs it only at checkpoints | Under NORMAL, a power loss can roll back the latest commits, but the database stays consistent. A design must say which it runs (both mobile loops run FULL on the writer) |
What to remember from Part 9
- B-tree engines log first too; the page's LSN makes replay safe to repeat, and checkpoints bound it.
- A torn page needs a whole-page copy: a full-page image, doublewrite, or an atomic device write.
- In SQLite's WAL mode the log holds pages, and a checkpoint copies them home.
Part 10. Replication as a durability level
Everything so far survives a power loss on one machine. A dead volume, a lost instance or a zone outage still takes the only copy. The next step up is more copies on other machines, and the question becomes: which copies must have the record before we say OK?
What each promise means
| System | Setting | What "acknowledged" means |
|---|---|---|
| Kafka | acks=all with min.insync.replicas | The leader waits for every replica currently in sync. min.insync.replicas (Apache Kafka default 1; MSK's default on a 3-AZ cluster is 2) is the fewest in-sync replicas allowed for the write to succeed. With the default, if the in-sync set shrinks to the leader alone, acks=all is satisfied by the leader alone. And no broker syncs to disk per message by default: "in sync" means in memory or page cache. Pair acks=all with min.insync.replicas=2 on 3 replicas |
| PostgreSQL | synchronous_commit with synchronous standbys | remote_write: the standby has received the WAL and written it to its OS. on: the standby has flushed it to disk. remote_apply: the standby has also applied it, so reads there see it |
| MySQL | Semi-synchronous replication | The source waits until at least one replica acknowledges; a replica acknowledges "only after the events have been written to its relay log and flushed to disk". If no replica answers within the timeout, the source falls back to asynchronous replication |
| Raft (etcd) | Built in | A follower persists the entries before replying; the leader commits once a majority has them. etcd syncs "when etcd persists its log entries to disk before applying them"; watch etcd_disk_wal_fsync_duration_seconds |
| DynamoDB | Built in | "A write is acknowledged to the application once a quorum of peers persists the log record to their local write-ahead logs" (DynamoDB paper, USENIX ATC 2022) |
| Aurora | Built in | The storage fleet keeps six copies across three AZs; a write needs four of them |
The shared mechanism in two sentences: the leader appends the record to its own log and sends it to the followers in parallel; it commits (and acknowledges) once W of the N replicas, itself included, have the record, each by its own sync policy. How a leader is elected and how a lagging follower catches up belong to the Replication, Quorums & Read-Your-Writes loop primitive and Primitive #09: Raft and Paxos.
The latency rule
The leader's local sync and the followers' round trips run at the same time, so commit latency is the slowest leg we must wait for, not the sum: max(local sync, the (W − 1)-th fastest follower's acknowledgement). At W = 2 of 3 that is the faster of the two followers. The slow follower does not count. This is how Raft-style logs work (the Raft paper, etcd, DynamoDB): the leader sends to followers in parallel with its own sync. PostgreSQL is different: it ships only WAL it has already flushed locally, so a synchronous commit costs the local flush plus the standby's round trip (with this branch's numbers, 1.0 + 1.4 = 2.4 ms).
Trace: three replicas take writes 21 to 24 (branch R)
Fork the main timeline after write 20. Replica A leads in AZ a; B (AZ b) and C (AZ c) follow; commit at W = 2 of 3. Assumptions, all labelled: A's local sync takes 1.0 ms, B acknowledges 1.4 ms after A sends (network plus its own sync), and C is slow at 6.0 ms. One batch is in flight at a time, as in Part 5, so each batch now takes max(1.0, 1.4) = 1.4 ms, and the batches regroup. From the reference implementation:
| Batch | Writes | Sent (ms) | A synced | B acks | C acks | Committed and OK |
|---|---|---|---|---|---|---|
| R1 | 21 | 11.0 | 12.0 | 12.4 | 17.0 | 12.4 |
| R2 | 22 | 12.4 | 13.4 | 13.8 | 18.4 | 13.8 |
| R3 | 23, 24 | 13.8 | 14.8 | 15.2 | 19.8 | 15.2 |
Synthesizing vector architecture diagram...
What to notice: the OK waits for the slower of A's own sync (1.0 ms) and B, the fastest follower (1.4 ms), so R3 commits 1.4 ms after it starts. C's 6.0 ms never enters the client's latency.
When the leader is cut off, when everyone loses power
| Scenario | What happens | Lesson |
|---|---|---|
| (ii) A syncs 23 and 24 locally, then is cut off from B and C before sending them | Only 1 of 3 replicas has 23 and 24: never committed, never acknowledged (clients time out). B and C elect B leader; B writes its own first entry (a no-op marking its term) at position 23. When A rejoins, its 23 and 24 conflict with B's log, and A deletes them to match | Durable on one disk does not mean kept. A write on one disk but not on a quorum is rolled back |
| (iii) B and C acknowledge from memory, before syncing; A syncs and acknowledges 23 and 24; then all three lose power together. Assume the followers' last own syncs covered up to 22 | B and C restart first with logs ending at 22. B asks C for a vote; C's log is no longer than B's, so it grants it, and B leads with 2 of 3. B writes at position 23, and A must delete its acknowledged 23 and 24 | One disk holding the data does not save it. Replication replaces a local sync only if the replicas' failures are independent, and a shared power event is not |
Local sync vs quorum, on equal terms
| Local sync only | Quorum of replicas that ack from memory | Quorum of replicas that each sync | |
|---|---|---|---|
| Latency added per write | One local sync, L | One round trip to the fastest follower(s) | max(local sync, the fastest follower's round trip plus its sync) |
| Process crash | Safe | Safe | Safe |
| Power loss on one machine | Safe | Safe (other copies) | Safe |
| Power loss on all of them together | Safe | Unsafe: branch R (iii) | Safe |
| Device or node loss | Unsafe: one copy | Safe | Safe |
| Zone loss | Unsafe | Safe if the quorum spans zones | Safe if the quorum spans zones |
| Cost | One disk's IOPS | N copies of storage and cross-zone traffic | Both |
Neither column alone covers everything: replication covers the lost disk and zone, local syncs cover the correlated power loss. Systems that promise both do both.
A team runs Kafka with acks=all and every other setting at its default, on three brokers. Can a message the producer saw acknowledged be lost?
One acknowledged write through every layer
Synthesizing vector architecture diagram...
What to notice: the local path runs down the middle and the replicas run beside it; the OK waits at the commit check for whichever legs W requires.
| Hop | "Done" at this hop means | Lost by |
|---|---|---|
| Client to engine | The request arrived | A client timeout: outcome unknown, retry idempotently |
| Group-commit queue | The record is in the engine's memory | A process crash |
pwrite into the page cache | The record is in kernel memory | An OS crash or power loss |
| File system, block layer, driver | The I/O is queued to the device, with a flush or FUA | The same |
| Device (or the Nitro card and the network to EBS) | The device has it, in its cache or on media | Power loss, if the cache is volatile and the flush has not completed |
| Stable media (EBS: replicated within its AZ) | Survives power loss | The loss of the device, the instance store or the zone |
| Replicas | W of N have it, each by its own sync policy | A correlated failure of the whole quorum |
| The OK travels back | The client knows | A lost reply: the write happened, the client doesn't know (Part 7) |
Write 23's budget in branch R, with max() for parallel legs and a sum for dependent ones:
| Leg | Time | Why |
|---|---|---|
| Queue wait: batch R2 is still committing | 13.8 − 12.5 = 1.3 ms | Dependent: comes first |
| Commit: max(A's sync 1.0, B's ack 1.4) | 1.4 ms | Parallel legs: the slowest required |
| C's ack, 6.0 ms | not counted | Not required at W = 2 |
| Total, engine arrival to OK | 1.3 + 1.4 = 2.7 ms | The reference implementation agrees: OK at 15.2 ms |
The naive sum of every leg, 1.3 + 1.0 + 1.4 + 6.0 = 9.7 ms, is more than three times too high. And requiring all three replicas (W = 3) would make the commit max(1.0, 1.4, 6.0) = 6.0 ms per batch. The client's own network round trip adds to all of these.
What to remember from Part 10
- Replication replaces a local sync only against independent failures.
- A write on one disk but not on a quorum is rolled back, not kept.
- Commit latency is the slowest leg you must wait for (the fastest follower at W = 2 of 3), not the sum.
Part 11. On AWS
Every AWS database and log service has a write-ahead log or a replicated log inside it; the question for a design is what each one publishes about when it says OK. And if you run your own engine, the log's sync budget comes out of the same volume as everything else.
Managed services that use a log
Stated only from AWS's public documentation and papers:
| Service | What AWS has published |
|---|---|
| Amazon DynamoDB | Each partition is a replication group using Multi-Paxos. Its storage replicas hold a write-ahead log and a B-tree; "log replicas" persist only recent log entries. A write is acknowledged once a quorum of peers has persisted the log record to their local write-ahead logs; the logs are periodically archived to S3 (USENIX ATC 2022 paper) |
| Amazon Aurora | The storage fleet keeps six copies across three AZs, with "a write set of four and a read set of three". Normal reads go to one storage node known to be up to date; a read quorum is needed only to rebuild state after a restart or failover. Data pages are built from the redo log |
| Amazon RDS (PostgreSQL, MySQL) | The engine's own WAL on EBS. A Multi-AZ DB instance replicates synchronously to a standby in another AZ; a Multi-AZ DB cluster uses semi-synchronous replication, which "requires acknowledgment from at least one reader DB instance". RDS Optimized Writes (MySQL) writes 16 KiB pages atomically, so the doublewrite buffer is turned off |
| Amazon MemoryDB | "Successful write operations are durably stored in a distributed Multi-AZ transactional logs before returning to clients" |
| Amazon ElastiCache (Valkey, with durability turned on) | Synchronous mode persists each write in a Multi-AZ transactional log across at least two AZs before replying (latency rises to single-digit milliseconds); asynchronous mode replies first and can lose up to 10 seconds of writes in a failure. Node-based clusters only |
| Amazon MSK | Apache Kafka on brokers with EBS storage: durability by replication. AWS recommends three AZs, a replication factor of at least 3, and min.insync.replicas of at most RF − 1 (so 2 with RF 3). MSK's default configuration sets unclean.leader.election.enable to true for clusters without tiered storage: if every in-sync replica is lost, an out-of-sync one can lead and acknowledged messages are lost, so set it to false for data you promised to keep |
| Amazon Kinesis Data Streams | "Synchronously replicates data across three Availability Zones". The metrics loop uses it as its replicated write-ahead log |
| Amazon OpenSearch Service | OpenSearch's translog, request durability by default (sync on the primary and every replica before OK). On the remote-backed instance families (OR1, OR2, OM2 and OI2), data is "copied synchronously to Amazon S3 as it arrives"; AWS does not spell out the exact acknowledgement point, so do not assume one |
Running it yourself on EC2
| Option | What to know |
|---|---|
| EBS gp3 | Replicated within its AZ. "Single-digit millisecond latency"; 3,000 IOPS and 125 MiB/s baseline, provisionable up to 80,000 IOPS and 2,000 MiB/s. The budget is shared: at the hook's rate, about 1,000 IOPS of log syncs plus about 760 for the LSM tree's flushes and compaction (200 MB/s at 256 KiB per I/O) is about 1,800 of 3,000; the log's 12 to 16 MB/s adds to throughput already above the baseline. Size it as write rate ÷ batch size, plus the engine's other I/O, plus headroom |
| EBS io2 Block Express | "Under 500 microseconds" average latency for 16 KiB I/O on Nitro instances. Use it when the sync latency itself bounds throughput or tail latency: at 0.5 ms, one log's ceiling doubles to 2,000 syncs a second |
| A separate log volume | A design choice: put the WAL on its own volume so its syncs do not queue behind flush, compaction or checkpoint writes |
| Instance-store NVMe | Very low sync latency, but the data survives only a reboot: it is lost when the instance stops, hibernates or terminates, or when the disk fails. Use it only as a cache, or with replicas in other AZs (the matching engine's journal nodes are this pattern) |
| Amazon S3 | Where sealed segments go for backup and point-in-time restore (the key-value store's WAL shipping; the matching engine's sealed journal segments). Stored across at least three AZs. Not the WAL itself: an S3 write per record would be far too slow |
| Torn-write prevention | EBS volumes attached to Nitro-based instances can write aligned 4, 8 and 16 KiB blocks all or nothing, which "eliminates the need for using the doublewrite buffer" (Part 9). On instance store, only some storage-optimized families support it, and not every size. Linux only |
| Measuring | Run the Part 2 commands on the actual volume and instance type, and alarm on the engine's own sync-latency and batch-size metrics |
Not the same thing
| Look-alike | Why a student might map it here | Why it is not |
|---|---|---|
| EBS snapshots | "Backups of the disk" | Point-in-time copies of a volume, stored in S3. Restoring one still depends on the engine replaying its own log |
| Amazon SQS | "A durable queue" | A managed queue with its own durability; putting a message there does not make your database's write durable |
| DynamoDB Streams | "A log of changes" | A change stream for consumers, kept for 24 hours; not what makes the table's writes durable |
| ElastiCache without durability | "It has replicas, so it's durable" | Replication is asynchronous: a failover can lose acknowledged writes. The durable mode in the managed table is the counterpart |
| Amazon EFS | "A file system, so fsync works the same" | EFS documents that, for Regional file systems, a synchronous write (fsync, O_DIRECT) or closing the file makes the data durable across AZs (a One Zone file system keeps it in one AZ). It is a place to put a log, with a network round trip per sync, not a log mechanism of its own |
What to remember from Part 11
- Managed services hide the log but not its acknowledgement rule: know what each one publishes.
- Instance store is fast and lost on stop: replicate.
- A volume's IOPS are shared by the log, flushes and compaction; count all of them.
Part 12. What you've learned
Back to the 100,000 writes a second
Our node had to acknowledge 100,000 writes a second, each surviving a crash, when one sync of its log takes about 1 ms and so caps it at 1,000 syncs a second:
- The WAL rule said what OK must mean: logged before applied, durable before acknowledged, against a named failure class (Part 1).
- A sync, not a write, makes a record durable, through the device's cache, with the directory and file size made durable too (Part 2), and a failed sync stops the engine instead of being retried (Part 3).
- Group commit shares each sync: about 100 writes per batch, about 1,000 syncs a second, waits of 1 to 2 ms, with about 150 writes in flight (Part 5). It keeps the same guarantee as a sync per write (Part 6).
- The shared budget holds: about 1,000 IOPS of syncs plus about 760 for the LSM tree's flushes and compaction is about 1,800 of gp3's 3,000 baseline, and the log's 12 to 16 MB/s adds to throughput the LSM page had already pushed above the baseline, so we provision more gp3 throughput (or move the log to its own volume).
- Checksums and the stop-at-first-bad-record scan make a torn tail harmless (Part 7); truncation keeps the log small unless a reader pins it (Part 8).
- Replicas in other AZs cover what one disk cannot, the lost volume and the lost zone, with commit latency set by the fastest follower (Part 10).
The cheat card
| Topic | Remember |
|---|---|
| WAL rule | Log before applying; durable before acknowledging |
| Failure classes | Process crash; OS crash or power loss; device or node loss; zone loss |
| The calls | write = page cache; fsync = data and metadata to stable media; fdatasync = data and needed metadata; O_DSYNC = each write as if followed by fdatasync; O_DIRECT = no page cache, not durable |
| Files | Sync the directory after creating a file; zero-fill in advance or recycle |
| Sync ceiling | Syncs per second ≤ 1 ÷ L (1 ms → 1,000) |
| Batch size | ≈ arrival rate × L; wait between L and about 2L |
| IOPS | Each sync of a batch up to 256 KiB is about one EBS I/O, shared with flushes and compaction |
| Failed sync | Fatal: stop, restart, recover from disk; an error is an unknown outcome |
| Record | Checksum + length + type + payload |
| Recovery | Cut the torn tail; stop at the first bad record |
| No OK is not "not written" | Durable-but-unacknowledged is normal: idempotent retries |
| Truncation | Durable elsewhere and no reader needs it |
| B-tree | pageLSN makes redo idempotent; full-page image, doublewrite or atomic writes for torn pages |
| Quorum latency | Raft-style: max(local sync, the (W − 1)-th fastest follower); PostgreSQL: local flush + standby round trip |
| Kafka | acks=all plus min.insync.replicas ≥ 2 |
Failure checklist
- Does every "OK" in the design name the failure class it survives?
- Is the sync before the acknowledgement in the code path that actually answers?
- Do sync errors stop the engine, never get retried?
- Are log files zero-filled or recycled, and the directory synced on create?
- Are sync latency and batch size metrics with alarms?
- Is consumer lag on every log reader alarmed before disk space?
- Is torn-page protection on, or are atomic writes proven for this instance and block size?
- Are replicas in separate AZs, with a sync policy that covers a power event they all share?
Think-first drills
Drill 1. Five writers reach one log at 0.0, 0.1, 0.3, 0.5 and 0.7 ms. A sync takes 0.5 ms, and a write that arrives just as a sync ends joins the next batch. With leader-based group commit, what are the batches and when does each writer get its OK? And with one sync per write?
Drill 2. Batch B6 writes records 9, 10 and 11 to the new segment wal-3, at bytes 0 to 28, 28 to 52 and 52 to 81. The device's units are 32 bytes. Power fails during the sync; units 0 (bytes 0 to 31) and 2 (bytes 64 to 95) reached the media, unit 1 (bytes 32 to 63) did not. What does recovery keep, and which of these writes had been acknowledged?
Drill 3. A payments ledger replicates every write to three AZs before acknowledging. An engineer proposes turning off the local fsync to cut latency. Is that safe?
Interview questions
| Question | Model answer |
|---|---|
What does fsync guarantee? And fdatasync, O_DIRECT? | fsync writes the file's data and metadata to stable media, including flushing the device's cache, and returns when done. fdatasync skips metadata not needed to read the data (timestamps), but not the size. Neither makes a new file's directory entry durable: sync the directory. O_DIRECT only bypasses the page cache; it is not a durability guarantee without a sync or O_DSYNC. |
| Why group commit, and what does it cost? | One log syncs one batch at a time, so a sync per write caps it at 1 ÷ L. A sync covers everything written before it, so writers queued during one sync share the next: batch ≈ rate × L. Costs: waits between L and about 2L, it needs concurrent writers, and each sync is still an I/O from a shared budget. |
What happens when fsync fails? | Linux may mark the pages clean, so a retry can succeed with the data gone. Treat it as fatal: fail the waiting writes, stop, restart and recover from disk. Clients that got errors must treat the outcome as unknown and retry idempotently. |
| What is a torn write, and how do the log and the pages survive it? | A write larger than the device's atomic unit, cut by a power loss, possibly with units out of order. The log survives with a length and checksum per record and a scan that stops at the first bad record. Pages survive with a whole-page copy: full-page images in the log, a doublewrite area, or atomic device writes. |
| How durable is Kafka by default? | Durability comes from replication, not syncs. acks=all waits for the in-sync replicas, but min.insync.replicas defaults to 1, so a lone leader can acknowledge. Use 3 replicas in 3 AZs and min.insync.replicas=2, and state the independent-failure assumption. |
| When can the log be deleted? | A segment can go when every record in it is durable elsewhere (flushed file or checkpoint) and every consumer (replica, backup shipper, CDC slot) has confirmed past it. The slowest consumer sets the log's size; alarm on its lag. |
Follow-ups an interviewer may add:
| Follow-up | Model answer |
|---|---|
| "Our drives have power-loss protection. Can we skip the sync?" | No. The sync still has to get the data from the page cache into the protected cache; the protection only makes the flush fast. And it does nothing for a dead device or a lost AZ. |
| "Kafka doesn't fsync. Why is that acceptable?" | Because the acknowledgement waits for replicas in other AZs, with min.insync.replicas=2, and it assumes those replicas do not all lose power at once. If that assumption is not acceptable, force flushes and pay for them. |
| "The client got an OK after a load-balancer retry. Was the write applied once?" | Not necessarily: the first attempt's outcome was unknown, and it may have succeeded too. Use an idempotency key so the second attempt is recognised as a repeat (the Idempotency & Effectively-Once Processing loop primitive). |
Where to go next
- LSM-Trees & Compaction: what the tree does with each record once the log has made it durable.
- Coming as loop primitives: Idempotency & Effectively-Once Processing (retries after unknown outcomes), Change Streams & the Transactional Outbox (the readers that pin the log), Replication, Quorums & Read-Your-Writes (elections and catch-up).
- Primitive #09: Raft and Paxos, Primitive #12: Change Data Capture and the Outbox Pattern, Primitive #21: Isolation Levels.
- Loops: the key-value store (steps 1.1, 1.6, R2.5, R2.10), the message queue (steps 1.1, 1.3, R2.5, R2.7), the matching engine (steps 1.4, 2.3), S3-like object storage (step 1.2), metrics (steps 1.1, 2.2), the mobile news feed (steps 1.3, R2.8), mobile chat (step 1.1), the transactional outbox ledger (steps R2.7, 3.2) and payments (step R3.9).