Building a Delay Queue for Scheduled Items¶
Retries with backoff, reminder emails, "check this payment again in 30 seconds", rate-limited resubmissions — a lot of work is not "do this now" but "do this at time T". Spawning a task that sleeps until T works for a handful of items and scales badly: each sleeping task holds a coroutine frame and its locals, nothing can see or cancel the backlog as a whole, and everything is lost at restart. A delay queue holds items in time order and releases each one when it becomes due. A 30-line heap-based version delivered 1,000 items with random delays up to 200 ms to four consumers at a median of 0.54 ms and a worst case of 1.08 ms after their due time, and correctly woke early when an item with an earlier deadline was inserted while it waited. A Redis sorted-set variant claimed 100 due jobs across three workers with no duplicates, and survives restarts.
Prerequisites¶
- Python 3.11+, stdlib only; the durable variant uses
redis.asyncio. - Queue basics, from Async Queue Management.
- Timer handles, from using loop.call_later and timer handles.
1. Order items by due time in a heap¶
A min-heap keyed on the due time gives O(log n) inserts and O(1) access to the next item. Consumers sleep until the earliest due time — or until an earlier item is inserted:
import asyncio
import heapq
import itertools
import time
class DelayQueue:
def __init__(self) -> None:
self._heap: list[tuple[float, int, object]] = []
self._seq = itertools.count() # tie-breaker: never compare items
self._changed = asyncio.Event()
def put(self, item, delay: float) -> None:
due = time.monotonic() + delay
heapq.heappush(self._heap, (due, next(self._seq), item))
if self._heap[0][0] == due: # new earliest item: wake sleepers
self._changed.set()
async def get(self):
while True:
if not self._heap:
self._changed.clear()
await self._changed.wait()
continue
due, _, item = self._heap[0]
wait = due - time.monotonic()
if wait <= 0:
heapq.heappop(self._heap)
return item
self._changed.clear()
try:
async with asyncio.timeout(wait):
await self._changed.wait() # woken early by an earlier insert
except TimeoutError:
pass
def __len__(self) -> int:
return len(self._heap)
The sequence number in each heap entry matters: without it, two items with the same due time are compared directly, which raises TypeError for dicts and other unorderable payloads. time.monotonic() rather than wall-clock time keeps the queue immune to clock adjustments.
Verify: schedule items with random delays; each get() returns them in due order, never before their due time.
2. Wake on earlier inserts, not on every insert¶
The put() only sets the event when the new item is the new head of the heap. Waking consumers on every insert would make them recompute their sleep for nothing; not waking them at all would make an item due in 50 ms wait behind a sleep for an item due in an hour. Measured: with an item due at 1 s already queued, inserting one due 50 ms later released the new item at 0.10 s — exactly when it was due.
The event-based wake is cheap and correct for many consumers on one loop: all of them re-check the head, one wins the heappop, the others go back to sleep. With thousands of consumers that thundering herd would cost something; with a handful — the normal case — it is negligible.
Verify: insert progressively earlier items while a consumer sleeps; each is delivered at its own due time.
3. Measure delivery precision¶
Delivery precision is set by the event loop's timer resolution and by how busy the loop is. Measure lateness — actual delivery minus due time — under realistic load:
async def measure(q: DelayQueue, n: int = 1000) -> list[float]:
lateness: list[float] = []
async def consumer():
while True:
item, due = await q.get_with_due() # variant of get() that also returns due
lateness.append((time.monotonic() - due) * 1000)
workers = [asyncio.create_task(consumer()) for _ in range(4)]
for i in range(n):
q.put(i, random.uniform(0, 0.2))
await asyncio.sleep(0.3)
for w in workers:
w.cancel()
return sorted(lateness)
On an idle loop: p50 0.54 ms, p99 1.06 ms, max 1.08 ms across 1,000 items. On a busy loop, lateness grows by however long the loop is blocked, which is why a delay queue is a good canary for loop lag — the technique in measuring event loop lag in production.
Verify: lateness p99 stays within a few milliseconds under normal load; a rise matches loop-lag alerts.
4. Make it durable with a Redis sorted set¶
An in-memory delay queue loses everything at restart — unacceptable for retries of payments or scheduled notifications. A Redis sorted set with the due time as the score is the standard durable version. The claim must be atomic so two workers never take the same item; a short Lua script does it in one round trip:
import time
import redis.asyncio as aioredis
CLAIM = """
local due = redis.call('ZRANGEBYSCORE', KEYS[1], '-inf', ARGV[1], 'LIMIT', 0, ARGV[2])
if #due > 0 then redis.call('ZREM', KEYS[1], unpack(due)) end
return due
"""
async def schedule(r: aioredis.Redis, job_id: str, delay: float) -> None:
await r.zadd("delayed", {job_id: time.time() + delay})
async def worker(r: aioredis.Redis, handle) -> None:
claim = r.register_script(CLAIM)
while True:
due = await claim(keys=["delayed"], args=[time.time(), 10])
if not due:
await asyncio.sleep(0.05) # poll interval = worst-case extra delay
continue
for job_id in due:
await handle(job_id.decode())
Verified on Redis 7.4: three workers claimed 100 scheduled jobs between them, every job exactly once. Wall-clock time.time() is necessary here because the scores are shared between machines; keep their clocks synchronised. Claimed-but-unfinished jobs are lost if a worker crashes mid-handle — for at-least-once delivery, move claimed ids to a processing set with a lease and re-queue expired leases, the pattern behind building a durable job queue on Postgres with asyncio.
Verify: run several workers against one sorted set; no job id is handled twice, and jobs scheduled before a restart are handled after it.
5. Choose in-memory, Redis, or a job system¶
# retries inside one request: in-memory is enough
retry_queue.put(request, delay=backoff(attempt))
# anything that must survive a deploy: durable
await schedule(redis, f"reminder:{user_id}", delay=3600)
In-memory delay queues suit short, disposable delays — client-side retries, debounced re-checks — where losing the backlog on restart only costs a little latency. Anything business-relevant belongs in durable storage. And once you need retries, uniqueness, visibility into the backlog and a UI, a job system with delayed jobs (arq's _defer_by, Celery's countdown) is less code to own, as compared in choosing between Celery, arq and taskiq.
Verify: for each delayed workload, the storage choice matches the cost of losing the backlog.
Verification¶
The delay queue is correct when:
- Items are released in due order and never early.
- Earlier inserts wake sleeping consumers, so a new short delay is not stuck behind a long one.
- Lateness is measured and stays within a few milliseconds on a healthy loop.
- Durable variants never double-deliver, and survive restarts.
Diagnostic Hook: export queue length and the age of the head item past its due time (zero when healthy). A head that is overdue and growing means consumers are saturated or the loop is blocked; for the Redis variant, also export ZCOUNT delayed -inf now, the number of due-but-unclaimed jobs, which should hover near zero.
Pitfalls & edge cases¶
- Heap entries without a tie-breaker. Equal due times compare payloads and raise
TypeError. - Wall-clock time in a single-process queue. Clock jumps reorder or stall items; use
monotonic(). - One task per delayed item. Fine for tens, a memory and visibility problem for thousands.
- Polling Redis too slowly. The poll interval adds to every item's delay; too fast, and idle workers load Redis.
Frequently Asked Questions¶
How do I implement a delay queue in asyncio?
Keep items in a heap ordered by due time, with a sequence number as a tie-breaker. Consumers wait until the head is due, using an asyncio.Event set by put() when a new item becomes the earliest so sleepers re-check.
Is a delay queue better than asyncio.sleep in a task per item?
For more than a handful of items, yes: one structure holds the backlog, it can be inspected and bounded, and consumers are a fixed pool rather than one sleeping task per item.
How precise is an asyncio delay queue?
On an idle loop, 1,000 items were delivered with a median lateness of 0.54 ms and a maximum of 1.08 ms. On a busy loop, lateness grows with event loop lag.
How do I make delayed jobs survive a restart?
Store them in durable storage such as a Redis sorted set scored by due time, and claim due items atomically with a Lua script or a database transaction so two workers never take the same job.
Related¶
- Async Queue Management — up to the topic overview.
- Retrying failed background jobs with backoff — the most common producer of delayed items.
- Concurrent Execution & Worker Patterns — the section overview.