Skip to content

Redelivering Items with Visibility Timeouts

Once a worker takes an item from asyncio.Queue, the item exists only in that worker. If the worker hangs, is cancelled, or gives up after a timeout, the item is gone, and nothing records that it was never finished. Message brokers solve this with a visibility timeout: a taken item becomes invisible for a while and reappears unless the worker acknowledges it. The same mechanism works in process. Measured on Python 3.14 with 500 items, 10 workers, 2% of items that hung until a 2-second per-item timeout and 10% that were slow but finished in 0.6 s: a plain asyncio.Queue processed 490 items and silently lost the 10 hung ones. An acknowledging queue with a 0.5 s visibility timeout also processed 490, redelivered the hung items until they reached a 3-attempt limit and moved all 10 to a dead-letter list — but it also redelivered the slow items while they were still being processed, so 54 items ran twice and 54 acknowledgements arrived too late. Extending visibility from a heartbeat while an item was being worked on brought duplicates to 0. This guide builds that queue.

Prerequisites

1. See where taken items disappear

With a plain queue, every exit from the worker that does not finish the item — a timeout, an exception, a cancellation — drops it:

async def worker(queue):
    while True:
        item = await queue.get()
        try:
            async with asyncio.timeout(2.0):
                await handle(item)
        except TimeoutError:
            pass                       # the item is now gone; nobody will retry it

Measured: 490 of 500 items processed, and the 10 that hung — 2% of the input — were never retried or recorded; the run's only trace of them was the difference between two counts. Re-putting the item in an except block helps with exceptions you anticipate, but not with a worker that is cancelled, a task that dies on an unexpected error, or one that is stuck in an await that never times out. A visibility timeout covers all of these with one rule: an item stays owned only while someone keeps acknowledging or extending it.

Verify: every item put into the queue is eventually either acknowledged or recorded as failed.

2. Give each taken item a receipt and a deadline

get() returns a receipt along with the item and records a deadline; ack(receipt) removes the item for good; a reaper task returns items whose deadline passes without an acknowledgement:

class AckQueue:
    def __init__(self, visibility: float, max_attempts: int = 3):
        self.visibility, self.max_attempts = visibility, max_attempts
        self.ready: asyncio.Queue = asyncio.Queue()
        self.inflight: dict[int, tuple] = {}                 # receipt -> (item, attempts, deadline)
        self.deadlines: list = []                            # heap of (deadline, receipt)
        self.receipts = itertools.count()
        self.dead: list = []
        self.redelivered = 0

    def put(self, item, attempts: int = 0):
        self.ready.put_nowait((item, attempts))

    async def get(self):
        item, attempts = await self.ready.get()
        receipt = next(self.receipts)
        deadline = time.monotonic() + self.visibility
        self.inflight[receipt] = (item, attempts + 1, deadline)
        heapq.heappush(self.deadlines, (deadline, receipt))
        return receipt, item

    def ack(self, receipt) -> bool:
        return self.inflight.pop(receipt, None) is not None   # False: it had already been redelivered

    async def reaper(self, interval: float = 0.05):
        while True:
            now = time.monotonic()
            while self.deadlines and self.deadlines[0][0] <= now:
                deadline, receipt = heapq.heappop(self.deadlines)
                entry = self.inflight.get(receipt)
                if entry is None or entry[2] != deadline:
                    continue                                   # acknowledged or extended
                del self.inflight[receipt]
                item, attempts, _ = entry
                if attempts >= self.max_attempts:
                    self.dead.append(item)
                else:
                    self.redelivered += 1
                    self.put(item, attempts)
            await asyncio.sleep(interval)

Workers acknowledge only after the item's work is done, and simply do nothing on failure — the reaper brings the item back. Measured: the 10 hung items were redelivered after each failed attempt and, after the third, moved to the dead-letter list; no item was lost.

Verify: after a run, every input item is in exactly one of: acknowledged, dead-lettered, or still queued.

500 items, 10 workers, 2% hang, 10% take 0.6 s A grid of 3 rows by 6 columns. 500 items, 10 workers, 2% hang, 10% take 0.6 s queue processed lost silently dead-lettered processed twice late acks asyncio.Queue 490 10 - 0 - AckQueue, 0.5 s visibility 490 0 10 54 54 AckQueue, 0.5 s + heartbeat 490 0 4 (others still retrying) 0 0 Per-item timeout 2 s; max 3 attempts; 10 s run.

3. Extend visibility while work is in progress

A visibility timeout shorter than an item's processing time redelivers items that are not stuck, merely slow — and the original worker then finishes and acknowledges too late:

receipt, item = await queue.get()
heartbeat = asyncio.create_task(keep_visible(queue, receipt))
try:
    async with asyncio.timeout(2.0):
        await handle(item)
    queue.ack(receipt)
finally:
    heartbeat.cancel()

async def keep_visible(queue, receipt):
    while True:
        await asyncio.sleep(queue.visibility / 3)
        queue.extend(receipt)                       # push the deadline forward

Measured without the heartbeat: the 10% of items that took 0.6 s outlived the 0.5 s visibility, were redelivered to another worker, and 54 items were processed twice, each with a late acknowledgement. With the heartbeat extending every third of the timeout, there were no duplicates and no late acknowledgements; the hung items still came back, because the per-item timeout cancelled the heartbeat along with the work. The visibility timeout then measures how long a dead worker holds an item, not how long work may take — which lets it be short. Brokers do the same with ChangeMessageVisibility, as described in consuming SQS queues with aioboto3.

Verify: under load with slow items, ack() never returns False for an item that was being worked on.

A slow item with and without a heartbeat A sequence of 6 messages between 4 participants. A slow item with and without a heartbeat worker A queue reaper worker B get: receipt, deadline +0.5 s 0.5 s: no ack, redeliver get: same item 0.6 s: ack -> False (late) with heartbeat: extend every 0.17 s 0.6 s: ack -> True Measured: 54 duplicates without the heartbeat, 0 with it.

4. Bound attempts and dead-letter what keeps failing

An item that fails every time — malformed data, a bug triggered by one input, an upstream that will never answer for it — would otherwise cycle forever, consuming a worker for a full timeout on every lap. Count attempts in the queue, not in the worker, so that every path back to the queue increments them:

if attempts >= self.max_attempts:
    self.dead.append(item)                  # record it with enough context to investigate
    DEAD_LETTERS.inc()
else:
    self.put(item, attempts)

Measured: with a limit of 3 attempts and no heartbeat, all 10 hung items reached the dead-letter list within the 10-second run; with the heartbeat, each attempt lasted the full 2-second per-item timeout plus the visibility, so 4 had been dead-lettered and the rest were on their final attempts when the run ended. Record the item with its last error and attempt count, and decide who reviews the list, as in implementing a dead-letter queue with asyncio. Add a delay before re-putting — exponential in the attempt count — when failures are likely to be transient, so a struggling upstream is not hit again immediately.

Verify: an item that always fails is dead-lettered after the attempt limit, with its attempt count and last error.

5. Remember the limits of an in-process queue

Redelivery in process protects against workers that hang, crash, are cancelled or give up — not against the process itself dying, which takes the in-flight table and the ready queue with it. For that, combine the acknowledging design with a disk-backed store, as in persisting queue items to disk, or use a broker. Within the process, export the numbers that show whether the timeout is set well:

def stats(self) -> dict:
    return {
        "ready": self.ready.qsize(),
        "inflight": len(self.inflight),
        "redelivered_total": self.redelivered,
        "dead_lettered_total": len(self.dead),
    }

Frequent redeliveries with few dead letters mean the visibility timeout is too short for normal work, or the heartbeat is missing; items stuck in flight for many visibility periods mean the reaper is not running; and every redelivery means work may run twice, so handlers must be idempotent.

Verify: redelivery and dead-letter counts are exported, and handlers tolerate running twice.

Does this queue need visibility timeouts? A decision on What happens if a worker never finishes an item with 4 outcomes. Does this queue need visibility timeouts? What happens if a worker never finishes an item? acceptable to drop asyncio.Queue 10 hung items lost silently must be retried AckQueue + visibility timeout 0 lost work can be slow heartbeat extends visibility 54 duplicates to 0 keeps failing max attempts, dead-letter 10 dead-lettered To survive the process itself, persist or use a broker.

Verification

Redelivery works as intended when:

  • Workers acknowledge only finished items, and every other exit lets the item return.
  • A reaper returns expired items, counting attempts in the queue.
  • Heartbeats extend visibility during legitimate long work, so acknowledgements are never late.
  • Repeat failures are dead-lettered after a bounded number of attempts, and handlers are idempotent.

Diagnostic Hook: count late acknowledgements — ack() returning False. Each one is an item that was redelivered while its first worker was still busy and has therefore run twice; any non-zero rate means the visibility timeout is shorter than real processing time or the heartbeat has stopped.

Pitfalls & edge cases

  • A plain queue for work that must finish. Measured: 10 items lost without a trace.
  • Visibility shorter than processing time. Measured: 54 items processed twice.
  • Attempts counted in the worker. Crashes and cancellations would not increment them.
  • Expecting in-process redelivery to survive a crash. Persist the queue for that.

Frequently Asked Questions

What is a visibility timeout?

The time a taken item stays reserved for one worker; if it is not acknowledged by then, it returns to the queue for another worker. It turned 10 silently lost items into 10 dead-lettered ones in testing.

How do I stop slow items being processed twice?

Extend the visibility from a heartbeat task while the item is being worked on; that took duplicates from 54 to 0 in testing.

How many times should an item be retried?

A small fixed number, counted in the queue: 3 attempts here, after which the item went to a dead-letter list.

Does an in-process ack queue survive a crash?

No: it protects against stuck or failing workers, not the process dying. Persist items or use a broker for that.