A queue. An ordered list. First in, first out. Building a job queue sounds like a nice weekend project. You can do it in Redis. You can do it in your database. What about on an LSM tree? How hard could it be?
We did it and we found out. Zizq is a job queue server that stores all of its own data in an embedded database built on an on-disk storage engine called an LSM tree, rather than sitting on top of an external storage layer like Redis or Postgres. Probably the most well-known LSM tree database and the one that’s often used as a reference is RocksDB.
As at the time of writing, Zizq is at version 0.7, but before the first 0.1 release was shipped, the core of how the queue stores and dispatches jobs was reworked multiple times. Not because we went blindly looking for elegance (well, maybe sometimes), but rather because each iteration along the way ran into some kind of uh-oh moment we needed to move past, such as slowdowns the more work was pushed through it, workers stampeding the database, or jobs stuck waiting on acknowledgement from a worker that was no longer running.
In this post we share our experiences, the underlying implementation of Zizq and the choices we faced along the way. We’ll start with the obvious (naïve) design, see how that runs into the same issues we ran into, and look at how we tackled each problem in turn. By the end you’ll have a pretty complete picture of what happens between your app enqueueing a job and a worker running it, and why we designed it the way we did.
If you’ve used background jobs in a web framework, probably through something like Sidekiq, Celery or BullMQ, you have all the background you need. We’ll explain LSM trees, skip lists and other concepts as we go. And if you’ve ever built this sort of thing yourself, you might recognise a few lessons we learnt too.
This post is not a comprehensive overview of the implementation of every feature Zizq offers. The focus is on the core queue machinery: how jobs are pushed into the queue, how they are dispatched to workers, what happens when jobs need to wait until later —or fail— and some (important) considerations around fault tolerance.
Why not just use Redis or Postgres?

It’s a good question, and many job queues are built on Redis or a relational database like Postgres.
I guess there are two primary factors: the first is just based on real-world experience running applications at scale; the second is from the position that Zizq should “just work” out of the box. You can literally download a single binary, start it and it’s up and ready to use right away.
When data is stored in an external service, those services add moving parts to be deployed and managed correctly. If your application already operates these services, you could just share the same database, or the same Redis instance, but once you reach a certain scale you eventually learn that such dependencies are better isolated from your primary application workload.
Redis is also memory-backed, which makes for a very different failure scenario compared to storage that is on disk. Queues built on Redis also make some very real trade-offs when it comes to queue introspection, locking guarantees and features like unique jobs.
Persistent storage provides durability guarantees and allows building features like mutability and unique jobs without imposing arbitrary expiries or deadlines. Using Postgres as a job queue begins to strain at scale, as queues inherently store ephemeral data and produce a huge amount of dead records to be vacuumed. Tuning Postgres correctly to deal with this can be a painstaking exercise. From experience there’s also a real likelihood that sharing your primary application database with your job queue can degrade application runtime performance when the queue is used heavily. The inevitable path you head down at scale is to provision and managed dedicated Postgres instances for your queues.
By using an embedded storage solution there’s nothing else to provision or manage, and as a persistent key-value store we can carefully design the data layout to suit the data access patterns needed by the queue. The thinking is in terms of collections of ordered keys and values, rather than tables, rows and indexes, and those key-value collections are designed to provide the most direct and efficient access to the data. This includes the basic enqueue and dequeue logic — and because Zizq makes visibility and control of data within the queue first-class, this also includes common data access requirements, such as listing and filtering jobs by type, queue or status. Jobs in Zizq are also mutable by design, so data needs to be updated even when it sits in the middle of a queue.
All that being said, by making a call to own the storage, we also accept the need to deal with all the hurdles that go with that.
A crash course in LSM trees
At this point it’s worth taking a moment to explain LSM trees.
Zizq is written in Rust, and its storage is built on fjall, an embedded key-value store based on a log-structured merge tree, or LSM tree. By embedded we mean it runs inside the Zizq process, more like SQLite than Postgres: no separate server, no network overhead, just a compiled-in library reading and writing files on disk.

An LSM tree is built around one idea: writing to disk sequentially is far cheaper than updating data in place at non-contiguous locations. So when you write a key, the storage engine doesn’t go looking for where that key lives on disk. Instead, the write is appended to a journal file (so it survives a crash) and added to an in-memory sorted table called the memtable. When the memtable reaches some configured size, it’s written out to disk in one go as an immutable, sorted file called an SST file. Over time those files accumulate, so a background process called compaction periodically merges them together, keeping only the newest version of each key.
Reading a key involves checking the memtable first, and then checking the files on disk (newest first), until you find it. Because everything is kept sorted, it’s also cheap to scan a range of keys in order much like an index in a relational database, which is going to be very useful for a queue and was one of the motivators for picking LSM tree storage in Zizq.
The way they work makes LSM trees a natural fit for write-heavy workloads, and a job queue is certainly write-heavy. Every job is written when it’s enqueued, then again when a worker takes it, and again when it completes or fails, often all within a fraction of a second, absolutely slamming writes on a busy workload.
One final and important point on LSM trees that has major ramifications: since files on disk are immutable, deleting a key can’t physically remove its entry. Instead it appends a tombstone: a marker indicating “this key was deleted”, which hides the old value until compaction eventually cleans both of them up. Attempts to read that key see the tombstone, treat the entry as non-existent and stop searching any further. The existence of tombstones poses some downsides for range scans, but we’ll get to that later.
Building an LSM-tree based job queue in 20 lines
So how do you build a simple job queue on top of an LSM tree? The obvious approach is surprisingly small. We started with something like this.
Every job needs to be stored somewhere, so each one gets a record keyed by its ID, holding all its attributes: its type, its queue, its payload, its status. That alone gives us a way to insert and retrieve any job, but it doesn’t inherently give us a way to find the next job that can be run. For that we need a second set of keys, an index, containing one entry for every job that’s ready to run:
| Key | Value |
|---|---|
job/{id} | Job data |
ready/{priority}/{id} | (Empty) |
Values (and indeed keys) are just raw bytes, and values can be entirely blank, which is useful when key-presence alone is sufficient.
This is where the fact LSM trees keep their keys sorted comes in incredibly
handy. Because the “ready” keys start with the priority, a scan from the start
of ready/ sees the most urgent jobs first. Note the numeric priority in Zizq
is stored as its fixed-width byte representation so they are always naturally
sorted, unlike their natural string representations such as 9 and 10.
Within the same priority, the job ID provides the secondary sort order. Zizq
uses IDs that are monotonically increasing by creation time, so jobs of equal
priority are retrieved in the order in which they were enqueued.
An ordered list. First in, first out. The database does all the work. Simple.
With those two kinds of key alone, the whole queue would require about twenty lines of pseudo-code:
def enqueue(job):
with db.transaction() as tx:
tx.insert(f"job/{job.id}", job)
tx.insert(f"ready/{job.priority:05}/{job.id}", "")
def dequeue():
with db.transaction() as tx:
key = tx.first_key_with_prefix("ready/")
if key is None:
return None
job = tx.get(f"job/{id_from(key)}")
job.status = "in_flight"
tx.insert(f"job/{job.id}", job)
tx.remove(key)
return job
def ack(job_id):
with db.transaction() as tx:
tx.remove(f"job/{job_id}")
Enqueueing writes the job and its ready index key atomically, so neither can exist without the other. Dequeueing takes the first ready index key, removes it, and marks the job as in-flight, all in one transaction, so two workers can never claim the same job. Acknowledging a finished job simply deletes it.
A worker then just loops continuously: dequeue a job, run it, ack it, and if there was nothing to dequeue, sleep for a moment before checking again.
The reality is more complex than this of course. Zizq also maintains indexes
by status, by queue and by job type, so it can handle queries like “which
jobs in the emails queue have failed?” by scanning and intersecting a few
sorted ranges of keys, rather than reading every job. But the general structure
is the same: records that hold the data, and indexes made of carefully
structured keys that point at them.
Keyspaces are not like database tables
Time for a quick confession. This is where I first got working with fjall wrong…
fjall lets you split your data into keyspaces, which are separate ordered collections of keys that can each be tuned differently. Coming from relational databases, I thought of them like tables: one for jobs, one for payloads, one for error records, one for each index. Structurally it looked tidy, and it did work, but it was wrestling against how fjall works under the hood.
Every keyspace has its own memtable, but they all share one journal, and an old journal file can only be discarded once every keyspace has flushed the writes it contains. A keyspace that’s written to rarely, like the one holding error records, takes a long time to fill its memtable, so it rarely flushes, and until then the journal needs to be retained. The result was far more journal to replay on startup than there needed to be, which made restarts slower.
My mental model was wrong. The correct approach was to stop thinking in terms of tables and to treat keyspaces are flat shared collections of keys, with keys of different types separated by prefixes. Zizq now has just two keyspaces, split by how they’re accessed rather than by what they store. The data keyspace holds the records themselves (jobs, payloads, error records), and is almost always point-read one key at a time, so it uses an LSM feature called bloom filters. Bloom filters are a compact summary kept for each file on disk that can answer the question “is this key definitely not in here?” without opening the file, saving a needless disk access. The index keyspace holds every secondary index, and is almost always queried as a range scan, where bloom filters don’t help.
Within each keyspace, every key starts with a one-byte tag indicating what kind of key it is, keeping the different kinds of keys in their own neatly sorted ranges.
With all of that in place, the naïve queue implementation works. Enqueue a few thousand jobs, start some workers, and it runs beautifully.
Until it starts to get bumpy.
Why does it get slower the more jobs we process?
While bechmarking throughput with a massive volume of jobs going through the queue it was clear something wasn’t right.
We do most of our local testing on two very different machine types: powerful Macbook Pros with relatively huge amounts of RAM, lots of cores and fast SSDs, and on cheap Raspberry Pi’s with far less RAM, limited CPU, running off of a microSD card.
Dequeueing started out snappy, got gradually slower the longer it ran, then suddenly recovered, before it started slowing down all over again. An immediately recognisable sawtooth pattern. So what was going on here?
I mentioned tombstones earlier, which are written to stand in as entries for deleted records in LSM trees. Think about what dequeueing actually does to the ready index. Every dequeue removes the first key in it, and removing a key in an LSM tree doesn’t remove anything at all. It writes a tombstone. So after a thousand dequeues, the start of the ready index is now a thousand tombstones, followed by the jobs still waiting to run.
For a single point read (for a specific key) that wouldn’t be much of an issue.
Looking up a single key stops at the first tombstone it finds for that key. But
finding the next job isn’t a point read, it’s a range scan. We don’t know which
key we’re looking for, only that we want the first one that still exists. So
the scan starts at the beginning of ready/ and has to step over every
tombstone, repeatedly, before it finds a live key. Each dequeue leaves another
tombstone behind for the next dequeue to step over, compounding the performance
hit to the range scan.
Eventually compaction runs, merges the tombstones with the keys they deleted, and the graveyard is cleared away. Dequeueing is fast again. Then the tombstones start stacking up all over again.
It’s worth reflecting for a moment on how unlucky this access pattern is. A queue writes new keys at one end of the sorted list, deletes them from the other end, and always scans from the end it deletes from. Every scan begins precisely where the freshest tombstones are present, which are exactly the entries compaction hasn’t got to processing yet.
It’s also important to note that tombstones can’t simply be dropped the first time compaction sees them. Compaction doesn’t merge every file at once, so an older copy of the key might still be sitting in a file this round of compaction didn’t touch. If the tombstone disappeared first, that old copy would suddenly be visible again, and a job we’d already dequeued would come back to life. So tombstones tend to survive several rounds of compaction before there’s nothing left for them to hide and they can finally be dropped.
If stable performance is a concern, a high throughput queue is arguably close to a worst case for an LSM tree when used in this way.
In our throughput benchmarks, on a humble Raspberry Pi 5, finding the next job could start taking around 40 milliseconds by the time the tombstones had accumulated. That doesn’t sound like much but it all adds up when you consider that cost is paid on every single dequeue, by every worker. Ideally more of that time would be spent by your workers running jobs, rather than waiting for jobs.
This is essentially the same class of problem as the one faced by a relational database like Postgres as dead rows accumulate in its pages until vacuuming occurs. Any storage engine that doesn’t update data in place has to clean up after itself at some stage, and because the data that a queue stores is largely ephemeral, it generates waste readily.
We could have gone down the path trying to tune our keyspaces to lessen the effect, such as by compacting more aggressively or trying to time compaction opportunistically. But ultimately we’re fighting the storage rather than working with it. The real question was whether the ready index needed to be stored in the LSM tree at all.
What if the queue lived in memory?
If we think about at what the ready index actually is, it’s a sorted list of short keys, where an entry typically lives for a few milliseconds before it’s removed. It’s scanned constantly, always from the front. An in-memory tree structure would serve the same need well, and much faster. We experimented with this approach, embraced it and that’s what is still in place today. The LSM tree is still the source of truth. Every job record lives on disk, and so does an index of every job by its status. But the decision about which job runs next is made from a sorted set held in memory, containing just the priority and ID of each ready job. Taking the first entry from an in-memory sorted set doesn’t require any kind of cleanup. There are no tombstones, no compaction, and therefore no sawtooth over time.
Taking this approach, the sustained throughput testing we ran on a Raspberry Pi 5 showed an improvement from around 40 milliseconds to around 4 microseconds. That’s quite a substantial improvement, and crucially that performance is constant under sustained throughput.
You might rightly think keeping the queue in memory is a bit risky. There are two obvious concerns. The first is memory use, but because the set only holds keys (think of them like pointers), and not the jobs themselves, it’s much smaller than you might expect: roughly 130MB for every million jobs that are ready to run but have not yet been claimed. Given most applications should look to scale workers to meet demand, this trade-off feels reasonable. The second concern is around what happens to that in-memory data when the server process dies. We’ll get to that later.
Reducing lock contention with skip lists
The first version of this in-memory queue used Rust’s standard BTreeSet
protected through a mutex. This worked, but of course it introduces lock
contention around every call site that enqueues or dequeues jobs. We replaced
it with a lock-free skip list structure, from the
crossbeam crate. If you haven’t
encountered skip lists before, they are a great data structure, still sorted
like the b-tree but lock-free through some clever use of hierarchical linked
lists and compare-and-swap operations.
Skip lists are well suited to concurrency because inserting or removing an entry only changes a handful of pointers right next to it, which means each change can be made with a compare-and-swap. Many threads can insert and remove entries at the same time, and where no collision occurs, there is no wait penalty.
Compare-and-swap also makes claiming a job from this structure safe. Imagine two workers both look at the front of the queue and see the same job. Both try to remove it, but only one compare-and-swap can succeed. The winner gets the job, and the loser simply tries again and sees a different job the next time.
The in-memory index is just a hint
There’s an important subtlety that must be understood with this approach. The job records on disk (in the fjall LSM tree) can change while their keys are sitting in the in-memory ready index. A job might be deleted through the API, or moved to a different queue, in the instant between a worker claiming its key and actually taking it. This is ok, and acknowledged because the in-memory set is treated only as advisory. It’s great at hinting which job should run next, but it doesn’t dictate that decision. After claiming a key from memory, Zizq immediately point-reads the job back from disk, checks it’s still there and still ready to run. If the job isn’t ready, the claim is quietly discarded and the next one is tried, and so on.
The write that marks the job as in-flight is itself a compare-and-swap against the record that was just read, so if anything changed the job in the meantime, the write fails, leaves on-disk data unchanged (i.e. the transaction is aborted) and the whole operation starts over rather than working from stale data.
This same pattern was adopted in other places too. Jobs scheduled to run later live in their own in-memory set, ordered by when they’re due, and jobs currently being worked on are tracked in another, ordered by when they were taken. In each case, the LSM tree remains the source of the truth and memory holds a fast, disposable view of keys to the relevant records.
What happens to all that memory when the power goes out?
Anything held in memory is lost when the process stops, whether that’s a clean shutdown, a crash, or someone tripping over the power cable. So with an in-memory queue, how do we avoid losing jobs?
Recall that the disk (the LSM tree) is the authoritative source. Every job, and its status, lives on disk. The in-memory indexes are always derivable from the source data. If we lose them, they can be rebuilt from that source data. The real work is around making sure that’s always true, even when the process dies at the worst possible moment.
What order do writes occur, disk-first or memory-first?
Most operations in Zizq change a job on disk and update at least one of the in-memory sets. Those are two separate steps, and the process could die between one and the other, so the order in which those updates occur is important and requires careful consideration.
When a job becomes ready (because it’s enqueued, or because a scheduled job reaches its due time, or because a job is returned to the queue after a worker disconnected), the change is committed to disk first, and only then added to memory. If we did this the other way around, a worker could take a job from memory before it has been committed, or after a failed commit, only to find the job does not exist in the LSM tree.
When a job is dequeued, we go the other way. The key is removed from memory first, since that removal acts as the claim, and only then is the job marked as in-flight on disk. If the commit fails (which should be an exceptional situation in practice), the key is put back so another worker can have a go.
Either way, if the process dies between the two steps, nothing is lost. The memory is gone regardless, but the disk still holds the last state that was successfully committed. Which brings us to what happens on start up.
Server start up sequence, and recovery
When Zizq starts, the disk is the only source of job data, and there are two things that need handling before normal operation can commence…
First, before we accept any connections at all, we need to deal with any jobs the LSM tree says are in-flight. Those are jobs that were being worked on when the process terminated, but the workers that had them were logically disconnected so we need to assume these jobs will never be acknowledged if nothing else changes. Any in-flight jobs are moved back to ready, so they can be picked up again (yep, that means a job can run more than once, which we’ll come back to later).
Secondly, we need to rebuild those in-memory indexes. The fact our LSM tree stores its own keys of jobs by status makes this operation a targeted range scan. Zizq scans just the range of the status index covering ready jobs, reads each of those jobs to find its priority and queue, and inserts them into the skip list. The same is done for scheduled jobs. Both rebuilds run in parallel, in the background, so the server is already up and accepting enqueues while they happen, though workers are not yet able to claim jobs.
Until the index rebuild is finished, the queue behaves as if it is empty to all workers, so a worker that connects in this window simply waits for new work, rather than working from a half-built index. As soon as the rebuild completes, all waiting workers are notified that there’s work to do, through a broadcast channel event, and normal operation resumes. This once-off startup process is mostly negligible as it only deals with the backlog.
How do workers retrieve work?
In the naïve example we covered earlier, workers poll for work, sleeping when no more work exists. This polling and sleeping approach is pretty wasteful and either hammers the disk excessively when dozens of workers are polling, or delays jobs that are otherwise ready to be executed when workers are sleeping. There’s not really a perfect solution when a fixed sleep is involved, although other widely-adopted queues do exactly this.
In Zizq’s model, workers don’t poll at all. A worker opens a single long-lived
HTTP request (to /jobs/take), and the server streams jobs down it as they
become available. The worker specifies via a query parameter how many jobs it’s
prepared to be given at once — its prefetch limit — and the server does the
work to ensure it always has that number. The default prefetch value is 1, so
by default the worker would receive one job and just sit there forever with no
further work to do if it does nothing else but listen on the stream. It is the
worker’s responsibility to make sure that job is processed on the client side,
after which the worker sends an acknowledgement which marks the job
completed. Marking a job completed removes it from the worker’s in-flight set,
which then frees up prefetch capacity for the server to deliver another job, if
one is ready. When there’s nothing to do, the connection just sits there
quietly, until a new job is enqueued.
This moves the problem from the client to the server, but at least the server knows exactly what is going on with jobs in the queue. Each connected worker is handled by its own asynchronous task on the server, and each of those tasks now needs to know when new work could be available to deliver on the stream.
Wakey wakey, eggs & bakey! 🍳🥓
Inside the server, every change to a job is announced on a broadcast channel: a job was enqueued, a job was completed, a job failed, and so on. Every worker connection subscribes to this channel and uses it as a flow control mechanism. The first version of this was as simple as it gets: workers take jobs until there are none, or the prefetch limit is reached, then they wait for an event on the broadcast channel. Whenever anything was broadcast, every waiting connection woke up and queried the database for new jobs.
This was wasteful and you can probably see why. With fifty workers connected, every single enqueue wakes fifty connections, all of which hit the database at once, and 49 of them come back empty-handed. Moreover, it wasn’t only enqueues that would trigger this wake up. Every acknowledgement, every failure, every scheduled job etc woke them all too, even though none of those inherently makes any new work available to workers. The busier the queue got, the more of the server’s time went on workers checking for jobs that weren’t there, and all of those checks were creating contention with the work that actually mattered. We had a thundering herd situation.
An intentional race
There were two changes involved in fixing this. Firstly, stop waking for events that never result in new work for workers (e.g. scheduled jobs). Our workers still need to process some events that don’t result in new work, such as completion or failure events which are handled by running some bookkeeping (fast, in-memory), but they just go back to waiting again after processing those events.
Secondly, and more interestingly, the event broadcasting that a job is available now carries two extra fields: the name of the job’s queue, so a connection that doesn’t serve that queue can just ignore it, and a claim token.
The claim token is just a boolean flag shared by all event subscribers.
Specifically it is an AtomicBool starting out false (unclaimed). A
connection that has capacity for more work, and whose worker is ready to
receive it, tries to flip the flag from false to true with a
compare-and-swap. Only one connection can succeed. The winner of the race goes
to the database and takes more work. All other connections go back to waiting
for new events without touching the database at all.
The winner of this race doesn’t necessarily stop at one job. If it has capacity for 4 more jobs, for example, it will try to take 4 jobs (but may get less if less are available).
What if the broadcast channel fills up?
A broadcast channel holds a fixed number of recent events, and a subscriber that lags too far behind, perhaps because the server is under heavy load, can theoretically miss some. Event processing is generally in-memory unless there really is work available, so lagging behind significantly is rare, but it’s still something that needs to be considered. A missed broadcast event could leave jobs sitting in the queue, and a missed completion could leave a connection believing it’s at its prefetch limit when it isn’t, holding its worker idle forever.
In Rust, broadcast channels at least tell subscribers when they have fallen behind, so Zizq handles this by taking a slower, more careful reconciliation path when that happens. The connection checks each job it thinks it has in-flight against the database, discards any that have since finished, and then goes and looks for work it may have missed while it was lagging. Not the end of the world, and close to the polling approach when compared with the normal event-driven path, but this logic only runs when something momentarily couldn’t keep up for any reason. The behaviour is correct as a fallback and it means a missed event can never leave a worker hanging.
Let’s talk about fault tolerance
Once a job has been sent to a worker on the streaming connection, the job is in-flight, and it stays that way until the worker acknowledges it. But workers aren’t always 100% reliable. They get terminated through deploys, scaled down, killed for using too much memory, lose power, or just crash unexpectedly. What happens to the jobs they were processing? Zizq ties the concept of whether or not a job is in-flight with whether or not its worker is still connected. If the server detects that a worker is no longer connected, it moves that worker’s jobs back from in-flight to ready.
How does the server know a worker is no longer connected?
Remember that the server delivers jobs to workers over a long-lived streaming response. That response is a stream fed by the connection’s job claiming task, and the client and server are bound by this connection: if the connection closes, the HTTP server tears down the response, and the task feeding it gets a notification that nothing is listening on the other end any more. At this point the server stops streaming, stops checking for new jobs, and cleans up.
When a worker process exits or crashes, its operating system closes the connection on its behalf, so this happens almost immediately. The server doesn’t need to be told the worker is shutting down, or why. A closed connection is the signal the worker is gone.
The tricky thing with a long-lived stream, is that it can be idle for long periods of time. If no jobs are being dequeued, nothing is being written to the stream, and a connection with nothing being written to it is very hard to monitor. The Zizq server sends a small heartbeat on any stream that has sat for more than three seconds (by default) without any other data being sent. On the client side, the worker simply discards these heartbeats. They serve two purposes:
- For the server, they mean there’s always something being written, and writing to a connection is a hook for disconnect detection logic. A write to a connection the other end has already closed fails, so a disconnect gets picked up within seconds even if the close itself was dropped.
- For the worker, they’re proof that the server is still connected. Our official clients all treat 30 seconds without receiving any data, not even a heartbeat, as a dead connection, and they reconnect.
Returning in-flight jobs to the queue
Each connection’s asynchronous task that takes jobs from the database and writes them to the stream keeps a record of the jobs it has sent on the stream that have not yet been acknowledged. We call this the worker’s in-flight set. When it detects a dead connection, the task loops through that in-flight set and moves every job back from in-flight to ready, both in the LSM tree and in the skip list, broadcasting the available jobs on the event channel so other connected workers can take them.
What if the worker leaves without saying goodbye?
A crashed process is the trivial case, because the operating system is still around to close the connection. The harder case is what happens when the whole machine abruptly disappears (e.g. power loss) or the network drops out somewhere between the client and the server. This is just radio silence at this point. Nothing closes the connection, because nothing is connected to close it. As far as the server can tell, the worker has just gone extremely quiet.
The heartbeats the server sends help to detect this case quickly. The server continues delivering the heartbeats, but now there is nothing on the other end confirming receipt of them (via TCP ACKs). TCP is tolerant of transient issues like this, and it is expected that occasional packet loss occurs, so the OS continues trying to retransmit the heartbeats hoping for eventual success. If enough time passes, the operating system gives up and reports the connection as dead.
The catch is around how long “enough time” really is. On Linux, with default settings, this takes about 15 minutes.
We built a test harness that fakes a worker’s host vanishing while it has a job in-flight, and measured the time it took for the connection to be closed. The connection was closed and the in-flight job moved back to ready 943 seconds after the host “disappeared”. That’s around 15 minutes where a job is held by a worker that no longer exists, and no other worker can run it while it is still seen as in-flight.
The fix is to explicitly set a TCP socket option — TCP_USER_TIMEOUT on Linux,
with equivalents on macOS and Windows — which caps how long unacknowledged data
can sit on the TCP connection before that connection is considered dead. Zizq
now sets this to 30 seconds by default to align with the official client
libraries’ 30 second client-side disconnect detection. Now in the same harness
the job moves back to ready in just over 30 seconds.
A note on at-least-once delivery
So a job can be delivered to a worker more than once in the case of a disconnect? Yes, by design. When a worker disappears, the server can’t know how far it was through processing that job. Maybe it never started the job. Maybe it finished it, and the acknowledgement was lost on the way. We just don’t know. Returning the job to the queue is the only safe option, since the alternative risks it never running at all. Of course this means the job might run more than once.
This is at-least-once delivery, and almost every job queue implements it. In practice it means job handlers should be written to be idempotent. Running a job twice should have the same effect as running it once. This is good practice in software design in general.
Implementing scheduled jobs
Zizq lets jobs be enqueued with a time at which they become ready to run, and until that time those jobs sit in the scheduled status rather than in the ready status. This feature is also the foundation on which exponential backoff is implemented.
A scheduled job doesn’t go into the ready index (the skip list) because it cannot yet be dequeued, so no worker can take it. We need a mechanism to move scheduled jobs from scheduled to ready when they reach their scheduled run time.
The first version of this feature did this with another index in the LSM tree, keyed by the time each job was due, and a background task that scanned the index for jobs whose scheduled time had passed. You can probably already see where this leads. Each due job is removed from the front of that index as it’s promoted, so the front fills up with tombstones just like the naïve queue implementation we saw earlier. Each scan has to step over more tombstones than the last. It’s the LSM ready index problem all over again, and the solution is the same. Scheduled jobs now live in their own in-memory sorted set, ordered by the time they’re due and then by ID, rebuilt from the status index on startup like everything else.
The scheduler is an alarm clock, not a stopwatch
The job that’s due next is always the first entry in the set, which means the background task, the scheduler, never needs to poll on a fixed cycle. Rather than waking up every second to check for jobs to promote, it looks at the first entry and sleeps until exactly that time. If the earliest job is due in 40 minutes, the scheduler sleeps for 40 minutes. If there’s nothing scheduled at all, it sleeps indefinitely.
Of course, a job could be enqueued while the scheduler is asleep, and that enqueued job could be due sooner than the one it’s currently waiting for. Like the worker streaming connection’s async task, the broadcast channel is utilised here again. Every newly scheduled job is broadcast on the channel along with the time it’s due, and the scheduler, while sleeping, listens for those announcements. If the newly enqueued job is due before the scheduler’s current alarm, it resets the alarm to the earlier time and continues sleeping. Otherwise, it ignores the event and continues sleeping. Either way the scheduler hasn’t wasted time touching the database.
Promotion from scheduled to ready
The scheduler processes in batches wherever possible. When the alarm goes off, the scheduler wakes, takes a batch of jobs that are now due from the front of the set, and promotes each one to ready. Both on disk and in the in-memory skip lists.
- The entry is removed from the scheduled set.
- The job is read back from disk to check it still exists and is still scheduled. It may have been deleted or updated in the meantime, in which case it’s skipped.
- The job is marked as ready on disk, with the same compare-and-swap as before.
- Only once that is successfully committed, is it then added to the ready skip list.
- It’s announced on the broadcast channel, so a waiting worker can pick it up.
At this point the job that was once scheduled is just like any other job on the queue. Its ID will be logically lower than freshly enqueued jobs, but it sits at the front of the queue ready to run immediately. If more jobs were due than fit in one batch, the scheduler loops again until nothing else is ready at which point it enters its sleeping state again.
Job failures and retries
Jobs fail. Doing that work on the queue lessens the impacts on users because the failure happens in the background. An upstream service may be down, a database query times out, a bug was introduced in a recent deploy…
When a job fails, the worker reports the failure back to Zizq (instead of acknowledging success), along with the error message and, if it has one, a backtrace. This job will need to be retried later, and because we already built scheduled jobs, this falls out as just a small layer on top of that existing machinery.
In a single transaction, we bump the job’s attempt count, and write an error record for that attempt (one record for each failed attempt), keyed by the job’s ID and the attempt number, so the whole history of a repeatedly failing job can be listed in order with a single range scan. In the same transaction we also calculate when the job should be retried, using an exponential backoff formula, and move the job from in-flight to scheduled, due at that time.
A retry is just a scheduled job with an increased attempt count and some error history.
Can we touch the disk less often?
Everything we’ve covered so far has been about reading less, such as reducing scans and avoiding unnecessary queries. We also implemented ways to reduce write overhead and increase write throughput, which are interesting in themselves.
fjall allows one write transaction at a time, which is a big part of what makes the compare-and-swap checks we’ve been relying on so simple, but it also means every write in the system queues up behind the mutex when another thread is writing. A job queue does lots of small writes. Every enqueue is a tiny transaction. Every acknowledgement is a tiny transaction, and so on. With many clients enqueueing and many workers dequeueing and acknowledging at once, that’s a lot of tiny transactions, each incurring a small fixed transaction overhead (snapshotting and committing), one after another. 1000 small transactions take longer to process than a single transaction that does the same 1000 writes.
Can we turn many single requests into fewer single transactions? This is write coalescing, and Zizq implements it on both the server and in our official clients in similar ways.
Timer-free write coalescing
Where it can, Zizq turns multiple individual requests into one grouped request, which is invisible to the application but provides a real performance boost. Every enqueue request is handed to a single dedicated thread, which reads from an in-memory queue and waits for the first request arrive in that queue. It then grabs whatever else has queued up in the meantime, if anything, and commits everything in one transaction. Each caller still gets back just its own result, and a request to enqueue many jobs together as a bulk enqueue is still all-or-nothing, even when it shares a transaction with other requests. Acknowledgements are coalesced in the same way.
The nicest thing about this is that there’s no arbitrary timer. A lot of batching procedures wait a few milliseconds in case anything else arrives, which adds the same arbitrary latency to every request, even when the server is largely idle. In Zizq, a request that arrives when the server is quiet is committed immediately, as a batch of one. Anything that arrives while that commit is in progress waits for the next commit and is committed along with anything else that was queued in the same timeframe. These requests would have had to wait anyway due to fjall’s single-writer design, so batching them saves overall write latency by reducing the total number of writes, and consequently improves throughput. Batches only grow when requests are arriving faster than they can be committed, so nothing needs tuning; this just naturally right-sizes itself.
Batching all the way down
Our official clients apply the same idea before a request even leaves the worker. In a concurrent worker, multiple jobs completing close together all need to send acknowledgments, each of which is a separate HTTP/2 request in a hot path. The HTTP/2 request is shared and multiplexed, but still one request with many acks is cheaper than 10 requests each carrying one ack.
The worker hands each acknowledgment off to a client-side background task, which takes whatever has queued up since its last request and sends them all at once in a single bulk request. The layers cooperate here: the client turns many acks into one request, and the server turns many requests into one transaction. Under extremely high throughput, a single commit can persist the results of hundreds of jobs across dozens of workers.
We use the same approach for fsync of the journal during database commits too where a series of commits that would each need to fsync are coalesced together into a single fsync. This is a well-established pattern in mainstream databases like Postgres because fsync is a very slow operation, but fjall does not (yet) implement it natively so Zizq layers it on top.
Final thoughts
In the big picture, the complexity of where we landed is a reasonable amount higher than where we started, but the end result works really well. We have observed well in excess of 100,000 jobs per second going through Zizq and the LSM tree beneath it provides durability and the ability to design features like mutability and unique jobs cleanly. Taking a hybrid approach involving the use of small focused sets of data in memory coupled with point-reads from the LSM tree greatly improves throughput for this particular use case, without sacrificing durability and while retaining the ability to query the database arbitrarily for non hot-path use cases. The same hybrid approach could of course be applied to a standard relational database, but LSM trees are optimised for write-heavy workloads, and jobs queues are write-heavy.
Our quest to squeeze out as much juice as possible is ongoing. Always seeking to reduce I/O overhead, by batching, by coalescing, or by doing some work in memory rather than in the storage can greatly improve performance.
As a general mantra in software design, simplicity wins and we should always strive for it, but sometimes the optimal design isn’t the simplest one. These techniques add moving parts, but each part in isolation is simple in itself and the pay off is worth it.