Skip to content

Deduplicating Items in an asyncio Queue

Many queues carry work that is about something — re-index user 42, refresh product 7, recompute tenant 3's quota — and the same key is often enqueued many times before a consumer reaches it. Processing every copy wastes capacity and, under load, makes the queue grow faster than consumers can drain it. A deduplicating queue keeps at most one pending entry per key and updates its payload when the key is enqueued again. In a test, 10,000 puts drawn from 1,000 keys produced 1,000 queued items and 9,000 coalesced updates, and the consumers processed each key once, with the latest payload. This guide builds it as a thin wrapper over asyncio.Queue and settles the question that makes it subtle: what happens when a key is enqueued while it is being processed.

Prerequisites

1. Index pending work by key

Queue the key, keep the payload in a dict, and skip the enqueue if the key is already pending:

import asyncio
from collections.abc import Hashable


class DedupQueue:
    def __init__(self, maxsize: int = 0) -> None:
        self._q: asyncio.Queue = asyncio.Queue(maxsize)
        self._pending: dict[Hashable, object] = {}
        self.coalesced = 0

    async def put(self, key: Hashable, payload) -> bool:
        if key in self._pending:
            self._pending[key] = payload        # latest payload wins, no new queue entry
            self.coalesced += 1
            return False
        self._pending[key] = payload
        await self._q.put(key)
        return True

    async def get(self) -> tuple[Hashable, object]:
        key = await self._q.get()
        return key, self._pending.pop(key)

    def task_done(self) -> None:
        self._q.task_done()

    async def join(self) -> None:
        await self._q.join()

    def qsize(self) -> int:
        return self._q.qsize()

The check and the dict update happen with no await between them, so on one event loop two producers cannot both decide a key is new. Wrapping rather than subclassing keeps the dedup logic independent of asyncio.Queue internals; maxsize, join() and backpressure all come from the inner queue.

Verify: enqueue the same key 100 times before any consumer runs; qsize() is 1 and the consumer receives the last payload.

10,000 puts over 1,000 keys 3 horizontal bars comparing puts with the others. 10,000 puts over 1,000 keys puts 10,000 coalesced into pending keys 9,000 queued and processed 1,000 Each key is processed once, with its latest payload, however often it was enqueued while waiting.

2. Decide what "latest payload wins" means for your data

Overwriting the payload is right when the payload is a description of desired state — "user 42's profile changed, here is the new version" — because only the newest version matters. It is wrong when each payload is a separate event — "charge £5", "charge £3" — because merging them by overwriting loses one.

For events, merge instead of overwrite:

class MergingDedupQueue(DedupQueue):
    def __init__(self, merge, maxsize: int = 0) -> None:
        super().__init__(maxsize)
        self._merge = merge

    async def put(self, key, payload) -> bool:
        if key in self._pending:
            self._pending[key] = self._merge(self._pending[key], payload)
            self.coalesced += 1
            return False
        return await super().put(key, payload)


# e.g. accumulate counter increments per key
q = MergingDedupQueue(merge=lambda a, b: a + b)

Merging works for commutative updates — counters, sets of changed fields, max timestamps. If events must be processed individually and in order, deduplication is the wrong tool; use per-key ordering from processing queue items in order per key.

Verify: for each queue, write down whether payloads are state or events; event queues use a merge function or no deduplication.

3. Handle keys enqueued while in flight

Once a consumer has taken a key, the key is no longer pending, so a new put() for it enqueues a fresh entry. That is usually correct — the state changed after the consumer read it, so it must be processed again — but it allows two consumers to work on the same key concurrently: one on the old payload, one on the new.

If concurrent processing of one key is unsafe, track in-flight keys and defer new puts until the current processing finishes:

class SerialDedupQueue(DedupQueue):
    def __init__(self, maxsize: int = 0) -> None:
        super().__init__(maxsize)
        self._in_flight: set = set()
        self._deferred: dict = {}

    async def get(self):
        key, payload = await super().get()
        self._in_flight.add(key)
        return key, payload

    async def put(self, key, payload) -> bool:
        if key in self._in_flight:
            self._deferred[key] = payload           # run again after the current pass
            self.coalesced += 1
            return False
        return await super().put(key, payload)

    async def done(self, key) -> None:
        self._in_flight.discard(key)
        if key in self._deferred:
            await super().put(key, self._deferred.pop(key))
        self.task_done()

Consumers call await q.done(key) instead of task_done(). The deferred payload is re-queued before the unfinished counter is decremented, so join() cannot return while a deferred follow-up is still owed.

Verify: with consumers that assert no other consumer holds the same key, a burst of puts for one hot key never trips the assertion, and the final processed payload is the last one put.

A key updated while a consumer is processing it A sequence of 6 messages between 3 participants. A key updated while a consumer is processing it producer dedup queue consumer get(): A, payload v1 (A in flight) put(A, v2): deferred put(A, v3): deferred payload now v3 done(A) re-queue A with v3 get(): A, v3 Two updates during processing became one follow-up pass with the newest payload.

4. Keep the pending map bounded

The dict grows with the number of distinct keys waiting, which under a backlog can be large. The inner queue's maxsize bounds it — every pending key has exactly one queue entry — so set maxsize and let put() apply backpressure to producers when it is reached:

q = DedupQueue(maxsize=50_000)          # at most 50k distinct keys waiting

Coalesced puts never block, because they do not add a queue entry: a producer re-enqueueing hot keys keeps flowing even when the queue is full of other keys. That is a useful property during incidents, when the same few keys are being hammered — the queue does not fill up with copies of them. Backpressure behaviour in general is covered in bounded asyncio queue with backpressure under load.

Verify: under a sustained burst, len(q._pending) never exceeds maxsize, and producers of new keys block while producers of pending keys do not.

5. Deduplicate across processes when needed

The in-process queue deduplicates within one worker. If several workers consume from a shared broker, each sees its own copies. Push the deduplication to the shared layer:

# Redis: a set of pending keys plus a list as the queue
async def enqueue(r, key: str, payload: str) -> bool:
    async with r.pipeline(transaction=True) as p:
        p.hset("pending:payload", key, payload)          # latest payload wins
        p.sadd("pending:keys", key)
        added = (await p.execute())[1]
    if added:
        await r.lpush("work", key)                        # only the first put queues work
    return bool(added)


async def dequeue(r) -> tuple[str, str] | None:
    item = await r.brpop("work", timeout=5)
    if item is None:
        return None
    key = item[1].decode()
    async with r.pipeline(transaction=True) as p:
        p.hget("pending:payload", key)
        p.hdel("pending:payload", key)
        p.srem("pending:keys", key)
        payload, *_ = await p.execute()
    return key, payload.decode()

The same pattern exists natively in many job systems as a job id or uniqueness key — arq's _job_id, for instance, prevents enqueueing a job whose id is already queued. The broader trade-offs are in Background Jobs & Task Queues.

Verify: two workers enqueueing the same key concurrently produce one queue entry in Redis.

How should repeated work for one key be handled? A decision on What does each payload represent with 3 outcomes. How should repeated work for one key be handled? What does each payload represent? the latest desired state dedup, latest wins process once commutative increments dedup with merge sum, union, max individual ordered events no dedup per-key ordering Deduplication is a statement about the data; decide it per queue.

Verification

The deduplicating queue works when:

  • Pending entries never exceed one per key, verified under bursts.
  • Processed payloads are the latest (or the merged result) for each key.
  • In-flight rules are explicit: concurrent same-key processing is either allowed deliberately or prevented.
  • Coalescing is measured and the pending map is bounded.

Diagnostic Hook: export the coalesced counter alongside puts and processed items. A coalescing ratio near zero means the dedup layer is not buying anything and can go; a ratio near one during an incident shows a few keys being hammered, and the top coalesced keys — logged periodically — name them.

Pitfalls & edge cases

  • Deduplicating events. Overwriting a payload that represented a payment or a message loses data.
  • Unhashable keys. Normalise to a string or tuple before enqueueing.
  • Forgetting done(key) in the serial variant. The key stays in flight and later puts are deferred forever.
  • Cross-process assumptions. An in-process dedup queue does nothing about duplicates in other workers.

Frequently Asked Questions

How do I prevent duplicate items in an asyncio Queue?

Queue keys rather than payloads, keep payloads in a dict keyed by the same key, and skip the enqueue when the key is already pending. In testing, 10,000 puts over 1,000 keys became 1,000 queue entries.

What happens if a key is enqueued while it is being processed?

With a simple dedup queue it is enqueued again and may be processed concurrently with the in-flight copy. If that is unsafe, track in-flight keys and defer new puts until processing finishes, then re-queue once with the latest payload.

Should a deduplicating queue keep the first or the latest payload?

The latest, when payloads describe desired state. For events, merge them with a commutative function, or do not deduplicate at all.

Does an asyncio dedup queue work across several worker processes?

No. Deduplicate in the shared broker, for example with a Redis set of pending keys, or use the job system's unique job id feature.