Skip to content

Persisting Queue Items to Disk

asyncio.Queue lives in memory. Everything in it — items accepted from clients, work scheduled for later, events waiting for a consumer — disappears when the process does, whether by deploy, crash or out-of-memory kill. When accepted work must survive, the queue has to be written to disk before the producer is told it was accepted. Measured on Python 3.14 with a producer putting an item every 0.5 ms, a consumer taking 2 ms per item, and the process killed with SIGKILL after 3 seconds: the in-memory queue had accepted 2,766 items, processed 1,382, and lost 1,384. A SQLite-backed queue in WAL mode accepted 2,546, processed 1,278 before the crash and the remaining 1,275 after a restart — 0 lost — with 7 items processed twice, because they were in a batch that had been claimed but not yet acknowledged. Durability has a price per commit, measured on an SSD: an in-memory put cost 0.26 µs; a SQLite put committed with synchronous=FULL cost 1,254 µs (797 per second); with synchronous=NORMAL, 50 µs; committing 100 puts per transaction with FULL, 27 µs each. This guide builds the disk-backed queue.

Prerequisites

1. Decide what "accepted" must mean

A queue's durability matters exactly at the point where a producer is told its item was accepted — an HTTP 202, an acknowledgement to an upstream broker, a "saved" message to a user. If that answer is given while the item lives only in memory, a crash loses accepted work:

async def accept(request):
    await queue.put(request.payload)              # in memory only
    return Response(status=202)                   # the client now believes it is safe

Measured: when the process was killed after 3 seconds, the in-memory queue had accepted 2,766 items and processed 1,382; the other 1,384 vanished, every one of them acknowledged to its producer. The consumer was slower than the producer, so a backlog built up — which is the normal reason to have a queue at all, and exactly what makes the loss large. Either make the queue durable before acknowledging, or acknowledge only after processing.

Verify: for every queue, you know whether an accepted item may be lost on a crash, and the producer's acknowledgement is consistent with that.

SIGKILL after 3 s with a backlog A grid of 2 rows by 6 columns. SIGKILL after 3 s with a backlog queue accepted processed before crash after restart lost processed twice asyncio.Queue (memory) 2,766 1,382 - 1,384 0 SQLite WAL, claim + ack 2,546 1,278 1,275 0 7 Producer every 0.5 ms; consumer 2 ms per item; process killed with SIGKILL.

2. Store items in SQLite before acknowledging them

SQLite gives a crash-safe store in one file, with transactions and no server. Keep the connection on one dedicated thread — SQLite connections belong to the thread that created them — and run every operation there with run_in_executor:

import asyncio, sqlite3
from concurrent.futures import ThreadPoolExecutor

class DiskQueue:
    def __init__(self, path: str, synchronous: str = "NORMAL"):
        self.path, self.synchronous = path, synchronous
        self.executor = ThreadPoolExecutor(max_workers=1)      # one thread owns the connection

    async def _run(self, fn, *args):
        return await asyncio.get_running_loop().run_in_executor(self.executor, fn, *args)

    def _open(self):
        self.db = sqlite3.connect(self.path, isolation_level=None)
        self.db.execute("PRAGMA journal_mode=WAL")
        self.db.execute(f"PRAGMA synchronous={self.synchronous}")
        self.db.execute("CREATE TABLE IF NOT EXISTS q(id INTEGER PRIMARY KEY, body TEXT, claimed INTEGER DEFAULT 0)")
        self.redelivered = self.db.execute("UPDATE q SET claimed = 0 WHERE claimed = 1").rowcount

    async def open(self):
        await self._run(self._open)

    def _put_many(self, bodies):
        self.db.execute("BEGIN")
        self.db.executemany("INSERT INTO q(body) VALUES (?)", [(b,) for b in bodies])
        self.db.execute("COMMIT")

    async def put_many(self, bodies):
        await self._run(self._put_many, bodies)       # returns after the commit

Measured: with the same crash, the SQLite queue lost nothing — every acknowledged put was either processed before the crash or still in the file afterwards. The single executor thread also serialises access, so there is no locking to get wrong, and the event loop never blocks on disk. aiosqlite wraps the same pattern if you prefer a ready-made API.

Verify: killing the process at any moment and restarting loses no item whose put had returned.

3. Claim, process, then acknowledge

Consumers must not delete an item when they take it — a crash during processing would lose it. Mark items as claimed, process them, and delete them only afterwards; on start-up, release any claims left by a crashed process:

    def _claim(self, n):
        self.db.execute("BEGIN IMMEDIATE")
        rows = self.db.execute(
            "SELECT id, body FROM q WHERE claimed = 0 ORDER BY id LIMIT ?", (n,)).fetchall()
        self.db.executemany("UPDATE q SET claimed = 1 WHERE id = ?", [(r[0],) for r in rows])
        self.db.execute("COMMIT")
        return rows

    def _ack(self, ids):
        self.db.execute("BEGIN")
        self.db.executemany("DELETE FROM q WHERE id = ?", [(i,) for i in ids])
        self.db.execute("COMMIT")

Measured on restart after the crash: 1,275 items were still queued, 20 of them claimed-but-unacknowledged — the batch in progress — which the start-up reset made available again. Seven of those had already been processed before the kill and were processed a second time. That is at-least-once delivery: nothing is lost, and anything in flight at a crash may repeat. Make processing idempotent, or record completion in the same database transaction as the acknowledgement, as discussed in deduplicating work across replicas. Smaller claim batches mean fewer repeats per crash, at the cost of more transactions.

Verify: items claimed at the time of a crash are redelivered after restart, and repeated processing is harmless.

An item's life in the disk-backed queue A flow of 5 stages. An item's life in the disk-backed queue put INSERT, COMMIT, then acknowledge claim mark a batch claimed process idempotently ack DELETE batch, COMMIT restart reset claims: redeliver At-least-once: nothing lost, in-flight work may repeat.

4. Pay for durability per commit, not per item

Every commit with synchronous=FULL waits for the data to reach stable storage. That is the guarantee, and it is slow:

await queue.put_many([item])                      # one fsync per item
await queue.put_many(items[:100])                 # one fsync per 100 items

Measured on an SSD, 5,000 puts of a 150-byte payload: an in-memory asyncio.Queue put cost 0.26 µs (3.8 million per second); SQLite in WAL mode with synchronous=FULL and one commit per put, 1,254 µs (797 per second); with synchronous=NORMAL, 50 µs (20,068 per second); with FULL and 100 puts per commit, 27 µs each (37,483 per second). NORMAL in WAL mode survives a process crash — what this guide's test did — but may lose the last transactions on a power failure or kernel crash; FULL survives those too. To keep FULL and its guarantee at high rates, batch: collect producers' items for a few milliseconds and commit them together, acknowledging each producer after the commit that contains its item, using the batcher from batching queue items by size and time.

Verify: the chosen synchronous level matches the failures you must survive, and throughput at that level is measured.

Puts per second, 150-byte items, on an SSD 3 horizontal bars comparing SQLite FULL, 1 per commit with the others. Puts per second, 150-byte items, on an SSD SQLite FULL, 1 per commit 797/s SQLite NORMAL, 1 per commit 20,068/s SQLite FULL, 100 per commit 37,483/s asyncio.Queue in memory: 3,797,574/s (off this scale). WAL mode throughout. The fsync per commit is the cost of FULL; batching shares it.

5. Keep the file healthy and the queue observable

A disk-backed queue is a database, with a database's housekeeping. Delete acknowledged rows promptly so the table stays small; let WAL checkpoints run, or run PRAGMA wal_checkpoint(TRUNCATE) periodically if the WAL file grows; and back up or replicate the file if the host itself might be lost — a local queue protects against process crashes, not disk failures. Export the depth and the oldest item's age from the table:

    def _stats(self):
        depth, claimed, oldest = self.db.execute(
            "SELECT count(*), coalesce(sum(claimed), 0), min(id) FROM q").fetchone()
        return {"depth": depth, "claimed": claimed, "oldest_id": oldest}

Store an enqueue timestamp per item to report age directly, as in monitoring queue depth and item age. When several processes or hosts must share the queue, or the backlog must survive the host, move to a broker or a database server table with SELECT ... FOR UPDATE SKIP LOCKED; a SQLite file is the right tool for one process that must not lose its own work.

Verify: after a long run, the database file size stays bounded and depth and age are exported.

Does this queue need to be on disk? A decision on What happens if queued items vanish with 4 outcomes. Does this queue need to be on disk? What happens if queued items vanish? nothing: producer retries or data regenerates asyncio.Queue 0.26 us per put accepted work is lost SQLite, commit before ack 0 lost on SIGKILL durable and fast batch commits under FULL 37,483 puts/s shared by hosts or must outlive the host broker / server DB not a local file Acknowledge only what is already somewhere a crash cannot reach.

Verification

A disk-backed queue is durable and usable when:

  • Producers are acknowledged only after the commit containing their item.
  • Consumers claim, process and then acknowledge, and claims are reset on start-up.
  • Processing is idempotent, because in-flight items repeat after a crash.
  • The synchronous level and batching are chosen from the failures to survive and the throughput needed.

Diagnostic Hook: after every restart, log how many claimed items were reset. A steady, small number is the in-flight batch at each crash or deploy; a number that grows over time means consumers are claiming more than they finish — an ack path that fails silently — and those items are being processed again and again.

Pitfalls & edge cases

  • Acknowledging before persisting. Measured: 1,384 accepted items lost.
  • Deleting items when they are taken. A crash mid-processing loses them.
  • Per-item commits with FULL. Measured: 797 puts per second.
  • Sharing one SQLite connection across threads. Give it one dedicated thread.

Frequently Asked Questions

How do I make an asyncio queue survive a crash?

Store items in SQLite (WAL mode) through a single-thread executor, commit before acknowledging the producer, and delete items only after processing. In testing, a SIGKILL lost 1,384 items from asyncio.Queue and none from the SQLite queue.

Will a disk-backed queue process items twice?

Items claimed but not acknowledged at a crash are redelivered: 7 were processed twice in testing. Make processing idempotent.

How fast is a SQLite-backed queue?

On an SSD: 797 puts/s committing each with synchronous=FULL, 20,068 with NORMAL, and 37,483 with FULL and 100 puts per commit.

Is synchronous=NORMAL safe for a queue?

In WAL mode it survives process crashes, as tested here; it can lose the most recent commits on a power failure or OS crash, which FULL prevents.