From v0.7 to v0.8: How group commit changed Turso's tail latency

Cover image for From v0.7 to v0.8: How group commit changed Turso's tail latency

About a year after I joined, Pekka Enberg, Turso's CTO, handed me a barebones pull request for MVCC (multi-version concurrency control). It built on an open source MVCC experiment he had worked on with Piotr Sarna and Avinash Sajjanshetty, based on the research behind Hekaton, Microsoft's in-memory database engine.

I suspect Pekka was also trying to cure me of a few stereotypically Spanish habits, like long naps and late breakfasts. I neither confirm nor deny indulging in either. Fittingly, the project turned out to be about keeping writers from sleeping on the job.

Thankfully, I wasn't alone for long. Everyone on the Turso team and a lot of open source contributors helped bring MVCC to life, so much so that we've been able to shift our focus from reaching feature parity with SQLite to improving Turso's performance. With 0.8, Turso has much lower tail latency and higher throughput than SQLite under concurrent writes. You can see the full results in Turso 0.8: Concurrent writes without SQLite's single-writer bottleneck.

This post covers one piece of that work. In v0.7, BEGIN CONCURRENT already let Turso move past SQLite's single-writer limit, but concurrent writes didn't scale the way they should have. Group commit fixed that in v0.8.

#What a commit means in Turso's MVCC

To see where commit time goes, it helps to know what a commit does. Here's the map.

TransactionMEMORY
Active→Preparing→CommittedSees a snapshot. Its changes stay private until commit.
adds row versions→
MVStoreMEMORY
SkipMap<Rowid, Mutex<RowVersions>>Lock-free map of row versions. Old ones are garbage collected.
commit: write rows, then fsync↓
Database fileDISK
Regular SQLite format, in pages.
checkpoint←
Logical logDISK
Only the rows a transaction changed, not whole pages.Replayed on restart.
The pieces a commit touches. Row versions live in memory. A commit appends the changed rows to the logical log and fsyncs it. A checkpoint later moves the data into the database file.

Turso's MVCC keeps recent row versions in memory, in a structure we call the MVStore. Rows live in a lock-free skip map that points each row ID to a list of that row's versions:

SkipMap<Rowid, Mutex<RowVersions>>

The skip map is lock-free, so transactions reading different rows never block each other. A naive design, such as a hash map behind one mutex, would make every reader wait for every other. Each row's version list does have its own lock, which matters later.

Every transaction moves through a few states. While a transaction is Active, its changes are visible only to itself. Other transactions see a consistent snapshot of committed data, which is how we implemented snapshot isolation. An Active transaction can also abort, on ROLLBACK or an error. When it commits, the transaction moves to Preparing. Turso checks whether another transaction already committed a change to the same rows. If so, it aborts with a write-write conflict. If not, its changes are written to disk.

Run a transaction
Active
Reads a snapshot. Its changes are private.
COMMIT→
Preparing
Checks for conflicts on the same rows.
no conflict: write log, fsync→
Committed
Durable. Visible to new snapshots.
ROLLBACK or error↓
write-write conflict↓
Aborted
Changes are thrown away. The transaction can retry.
Hover a state or an arrow, or run a transaction.
The states of a transaction. It becomes Committed only after its changes are durable on disk.

Those changes go to a file called the logical log. It's similar to SQLite's write-ahead log (WAL), with one important difference. The WAL records whole pages, typically 4 KB each, so changing one row still writes an entire page. The logical log records only the rows a transaction changed: new rows, updates, and deletes (written as tombstones), with schema changes tracked separately. On restart, Turso replays the logical log to rebuild the in-memory state.

The logical log can't grow forever, and memory is limited. So past a threshold, a checkpoint moves data into the regular SQLite-format database file, and garbage collection frees row versions nobody can see anymore.

In short, committing a transaction takes two steps:

  1. Write the transaction's rows to the logical log.
  2. Fsync the logical log, so the data is actually on disk.

The second step is where I went down a performance rabbit hole. Unlike most rabbit holes, this one led somewhere good.

Figure 1 steps through that path in v0.7. A transaction does all its work in memory. The commit is the only place it touches the disk, and the disk is where the milliseconds are. Press play, or click details on any box.

step 1/10 · t = 1 µs of 914 µs
in memory, µson disk, msheld under the commit lock
step 1Connection takes begin_ts = 41. Its snapshot is fixed. Other connections keep running.
1 · RUN, ALL IN MEMORY
1BEGIN CONCURRENT1 µs
→
2Write rows150 µs
→
3MvStore: version chains
2 · COMMIT · DECIDE
4Validate8 µs
→
5Wait for dependencies1 µs
→
6Build the log record15 µs
3 · MAKE DURABLE, PUBLISH · lock held
7Take the commit lock2 µs
→
8pwrite25 µs
→
9fsync700 µs
→
10Publish, release the lock12 µs
MVSTORE · MEMORY
no versions from this transaction yet
LOG RECORD · MEMORY BUFFER
not built yet
.DB-LOG · DISKcommit lock: free
page cache
nothing pending
durable
frames 1…57
Figure 1: A BEGIN CONCURRENT transaction never writes a page. It adds row versions in memory, and at commit it validates, waits for dependencies, serializes one log frame, writes it with pwrite, fsyncs, and only then flips its state to Committed. Readers and writers never block each other during the run. In v0.7 the commit lock, whose job is to serialize writes to the logical log, was taken before the pwrite and released after the publish, so every other committer waited for a flush that was not theirs.

#The v0.7 baseline: MVCC that didn't scale

The whole point of MVCC is that write throughput should grow as you add concurrent writers. SQLite allows one writer at a time, so its throughput stays flat no matter how many connections you open. With BEGIN CONCURRENT, Turso lets transactions work in parallel, so more connections should mean more commits per second.

In v0.7, that wasn't really happening. I noticed it while preparing for my upcoming P99 CONF talk, when I ran the write benchmarks Pekka had built. Adding connections didn't buy us nearly as much throughput as it should have.

Those benchmarks are the same ones in the 0.8 release post. One is a closed-loop throughput run: each connection inserts 100 rows on disjoint keys, then immediately starts the next transaction. The other is an open-loop latency run: 1,000 transactions per second with Poisson arrivals, spread across 1, 8, 16, and 32 connections, measured from the moment a transaction is scheduled until it commits. SQLite stays flat on the first and grows a long tail on the second, because it has one writer. Turso in v0.7 should have climbed. It didn't.

So I profiled it. The profiles showed where the time went: a large share of every commit was spent in fsync. You can see that in Figure 1 already. The fsync box is the only step measured in milliseconds. Everything else is microseconds. Figure 2 shows what that does when more than one connection reaches COMMIT at the same time.

Connections
Speed
running, in memorywaiting for the lockpwritefsync
v0.7one fsync per commit
commits fsyncs 0
Figure 2: The lock lane on top shows who holds the commit lock, which serializes writes to the logical log. A connection that reaches COMMIT while another holds the lock sits in the red band until that fsync finishes, then does its own pwrite and its own fsync. With k connections committing at once the last one waits for k flushes. At 1× one millisecond of engine time takes one second on screen, so a 700 µs fsync takes 0.7 s.

#The bottleneck: every transaction paid for its own fsync

Writing data to a file doesn't mean the data is on disk. When a write call returns, the data may still be sitting in the operating system's page cache or in the disk's own cache. If the machine loses power at that moment, the data is gone. A database that promises durability has to call fsync, which forces everything written so far onto stable storage, and wait for it to finish.

That wait is expensive. Fsync is one of the slowest operations in the commit path, and it blocks the disk while it runs.

In v0.7, every transaction did both steps on its own. With N transactions committing, that meant:

  1. N writes to the logical log
  2. N fsyncs

The writes are genuinely per transaction: each one has its own rows to append. The fsyncs are different. Every fsync does the same thing, making the logical log durable up to that point. When many transactions commit at nearly the same moment, most of those fsyncs are redundant. One fsync after all their writes would make every one of them durable.

So with more concurrent writers, we were queueing more and more identical, slow disk flushes. That's why adding connections didn't scale the way MVCC should.

#The fix: group commit

The fix is a technique almost every serious database uses: group commit (or batch commit). Instead of each transaction flushing on its own, transactions that are committing at about the same time share one fsync.

Group commit wasn't a new idea for us. I had tried implementing it a long time ago, but at the time fsync wasn't what limited us, so it added complexity without a clear payoff. It only became the right thing to build once the profiles showed fsync blocking the database from scaling.

Here's how it works now in v0.8:

  1. Transactions gather into a group. When a transaction reaches the Preparing state, it enqueues its log record and gets a ticket. There is no timer. Transactions that arrive while a flush is already running join the next group.
  2. One transaction becomes the leader. The leader is not chosen. The first transaction to take the commit lock while the queue is non-empty becomes the leader and drives all writes and the fsync.
  3. The leader does the disk work for everyone. It writes every queued transaction's rows to the logical log, then calls fsync once.
  4. The others wait for the leader. The other transactions in the group park until the leader's fsync covers their ticket. Once it does, their data is durable, and each continues independently to finish its commit.

Figure 3 is the same run as Figure 1, with the durable phase rewritten. One leader writes every queued record and flushes once. The others sleep until that flush covers them.

step 1/12 · t = 1 µs of 917 µs
in memoryon diskleader holds the lockparked, no CPU
step 1Connection A takes begin_ts = 41. Connections B, C and D are running their own transactions.
1 · RUN, ALL IN MEMORY
1BEGIN CONCURRENT1 µs
→
2Write rows150 µs
→
3MvStore: version chains
2 · COMMIT · DECIDE
4Validate8 µs
→
5Wait for dependencies1 µs
→
6Build the log record15 µs
3 · GROUP, FLUSH ONCE, PUBLISH
7Enqueue → ticket1 µs
→
8Lock free? lead. Busy? park.2 µs
→
9pwrite × batch25 µs
→
10fsync, once700 µs
→
11durable_through → wake2 µs
→
12Publish, release the lock12 µs
MVSTORE · MEMORY
no versions from this transaction yet
LOG RECORD · MEMORY
not built yet
COMMIT COORDINATORlock: free
pending queue · tickets
empty
durable_through
9
.DB-LOG · DISK
page cache
nothing pending
durable
frames …9
Figure 3: In v0.8 a committer enqueues its record and receives a ticket, its position in the log order. Whoever takes the lock first while the queue is non-empty leads: it writes every queued record in ticket order, fsyncs once, and marks the log durable through the last ticket. Waiters park on a completion the coordinator owns and run again only when their ticket is durable. Each transaction still flips to Committed only after the fsync that covers it, so durability and visibility are unchanged.

The arithmetic is simple. For N transactions:

WritesFsyncs
v0.7NN
v0.8N1 per group

The writes don't go away, because every transaction still has its own rows to append. But the slowest, most redundant step now happens once per group instead of once per transaction.

Figure 4 puts the two commit paths on the same connections, with the same transaction durations. Only the commit path differs. Count the fsync blocks.

Connections
Speed
runningwaiting for the lockin a batch, waiting for its fsyncpwritefsync, v0.7fsync, v0.8 batch
v0.7one fsync per commit
commits fsyncs 0
v0.8one fsync per batch
commits fsyncs 0
Figure 4: In v0.8 connections still wait for the lock (red): a connection that reaches COMMIT while a batch is flushing parks with a ticket until that flush is done. But it waits for at most one flush. When the lock frees, the next leader writes every queued record and fsyncs them together, and the others sleep in that batch until it is durable. The lock lane still shows one holder at a time, but each hold now covers a batch. Count the fsync blocks in each lane to see the difference: v0.7 pays one per commit, v0.8 one per batch.

#The results: v0.7 versus v0.8

To measure the effect of group commit, I reran the transaction latency benchmark with three setups: SQLite, Turso v0.7 (without group commit), and Turso v0.8 (with group commit). The benchmark schedules 1,000 transactions per second with Poisson arrivals and measures each transaction from when it's scheduled until it commits, at 1, 8, 16 and 32 connections.

Connections
Y axis
Turso v0.8Turso v0.7SQLite
0.1 ms1 ms10 ms100 ms1 s10 s100 s0%25%50%75%100%0%90%99%99.9%99.99%99.999%Transaction latency →p99.9 730 msp99.9 13.1 sp99.9 5.85 ms
Enginep50p99p99.9max
Turso v0.80.87 ms1.67 ms5.85 ms46.3 ms
Turso v0.76.67 s12.8 s13.1 s13.1 s
SQLite1.35 ms129 ms730 ms2.53 s
Figure 5: Transaction latency at 1,000 transactions per second with 1, 8, 16 and 32 connections, measured from when a transaction is scheduled until it commits, pooled over three runs. Markers show the median and 99th percentile, and dashed lines the 99.9th percentile. Latency is on a log scale; switch the y axis to Tail to zoom into the slowest transactions, and hover to read any percentile.

The honest headline is that without group commit, Turso's MVCC had a far worse tail than SQLite. At every connection count, the slowest transactions took seconds, and at 32 connections the p99.9 reached 157 seconds. Allowing concurrent writes didn't help, because every transaction still waited on its own fsync.

The numbers are this extreme because the benchmark offers a fixed load of 1,000 transactions per second. Paying one fsync per transaction, v0.7 couldn't keep up with that rate, so transactions queued, and each one's latency includes the time it spent waiting to start. You can see the backlog forming in the single-connection curve: about three-quarters of transactions finish quickly, then the rest stall for seconds. Adding connections made it worse, since more writers meant more fsyncs competing for the same disk.

With group commit, every curve moves left by three to five orders of magnitude. Median latency sits around 1 ms at every connection count, and the p99.9 is lower than SQLite's in every configuration, from 14 ms at one connection down to 2.4 ms at 32.

#Why more connections means lower tail latency

One result looks backwards at first: Turso's tail latency goes down as connections go up. At a single connection, p99.9 latency was 14 ms. At 32 connections, it was 2.4 ms.

Group commit explains this. With few concurrent writers, groups are small, and each transaction still pays for something close to a full fsync. With many writers, each group holds more transactions, so the cost of one fsync is shared across more of them. The busier the database, the better group commit amortizes its most expensive step.

The contrast with SQLite is even more impressive because, under the same load, SQLite's p99.9 grows from 73 ms at one connection to 1.2 s at 32, because writers back off and sleep while waiting for the single write lock.

#Where this works best

The throughput benchmark writes to different rows in each transaction and that's deliberate: it's the workload we expect to be most common, and the one where Turso's MVCC works best.

When transactions update the same row, two things go wrong. First, they contend on that row's lock in the MVStore, so they block each other. Second, at commit, all but the first fail with a write-write conflict and have to retry. We saw this with one user whose workload updated the same rows every few milliseconds, and MVCC wasn't a good fit for it (at least for now).

The sweet spot is writing a lot of data where transactions mostly touch different rows. If your writers constantly fight over the same few rows, SQLite's single-writer model may serve you just as well.

#Try it out

Install Turso locally and run BEGIN CONCURRENT on your own workload:

$curl -sSL tur.so/install | sh

MVCC is also available in tech preview on Turso Cloud. The benchmarks are open source, so run them on your hardware and tell us what you find on Discord.

If you want the rest of the story, including how MVCC works end to end, how we handle conflicts, and the hard parts of checkpointing, come to my talk, "Breaking SQLite's Single-Writer Bottleneck," at P99 CONF on October 21–22.