Skip to content

Making Background Jobs Idempotent

Every practical job system — Celery, arq, taskiq, SQS, RabbitMQ, Kafka consumers — delivers jobs at least once. A worker that crashes after doing the work but before acknowledging it gets the job redelivered; a retry after a timeout runs a job whose first attempt may have succeeded; a visibility timeout that expires mid-job hands it to a second worker while the first is still running. The job will sometimes run twice, and possibly at the same time. Measured against PostgreSQL 17 with every one of 100 credit jobs delivered three times concurrently, a naive job wrote 300 ledger rows, tripling every balance; the same job guarded by a processed-jobs row inserted in the same transaction wrote exactly 100, and reported the other 200 deliveries as duplicates. This guide covers that guard and the other techniques for jobs whose effects are not in your database.

Prerequisites

1. Give every job a stable key

Idempotency needs something that identifies "the same job" across deliveries. A generated job id is not enough if the producer can enqueue the same logical work twice (a web request retried by the client); derive the key from the business operation instead:

def job_key(kind: str, *parts: object) -> str:
    return f"{kind}:" + ":".join(str(p) for p in parts)


await queue.enqueue("credit_account", job_key("credit", order_id), account_id, amount)

credit:order-8812 means "the credit for order 8812" no matter how many times it is enqueued or delivered. Keys based on time (credit:2026-10-02T10:15) or random ids defeat the purpose. If the operation legitimately repeats — a monthly fee — the period belongs in the key: fee:account-7:2026-10.

Verify: enqueue the same logical operation from two code paths; both produce the same key.

2. Record the key in the same transaction as the effect

When the job's effect is a database write, the guard is a unique row in the same transaction:

import asyncpg

SCHEMA = """
create table if not exists processed_jobs (
    job_key text primary key,
    done_at timestamptz not null default now()
);
"""


async def credit_account(pool: asyncpg.Pool, key: str, account: int, amount: int) -> str:
    async with pool.acquire() as conn, conn.transaction():
        claimed = await conn.fetchval(
            "insert into processed_jobs(job_key) values($1) on conflict do nothing returning job_key",
            key,
        )
        if claimed is None:
            return "duplicate"                         # another delivery already did it
        await conn.execute("insert into ledger(account, amount) values($1, $2)", account, amount)
        return "applied"

Because the guard row and the effect commit or roll back together, there is no window in which one exists without the other. Two concurrent deliveries are serialised by the primary key: the second insert … on conflict do nothing waits for the first transaction to finish, then sees the conflict and returns nothing. If the first transaction rolls back, the second proceeds and does the work — exactly the behaviour you want from a retry. Measured: 300 concurrent deliveries, 100 effects, 200 duplicates.

Verify: deliver every job several times concurrently; effect counts equal job counts.

100 jobs, each delivered three times concurrently 2 horizontal bars comparing naive job with the others. 100 jobs, each delivered three times concurrently naive job 300 rows, balances tripled guarded job 100 rows, 200 duplicates skipped PostgreSQL 17, asyncpg pool of 20 connections. The guard row and the effect commit together, so the database enforces exactly one effect per key.

3. Handle effects outside your database

When the job calls an external API — charging a card, sending an email, posting to a webhook — the database transaction cannot cover it. Two techniques, in order of preference:

Pass an idempotency key to the API. Payment providers and many other APIs accept one and return the original result for a repeated key:

async def charge_order(pool, order_id: str, amount: int) -> str:
    key = job_key("charge", order_id)
    charge = await payments.create_charge(amount=amount, idempotency_key=key)
    async with pool.acquire() as conn, conn.transaction():
        await conn.execute(
            "insert into processed_jobs(job_key) values($1) on conflict do nothing", key)
        await conn.execute("update orders set charge_id=$1 where id=$2", charge.id, order_id)
    return charge.id

Check-then-act with a state machine when the API has no key: record "started" before calling, "done" after, and on a redelivery that finds "started", query the external system for whether the effect happened before calling again. That is more code and still has a window; prefer APIs that support keys, and for messages you publish yourself, use the transactional outbox described in implementing the transactional outbox pattern in asyncio.

Verify: kill the worker after the API call but before the database commit; on redelivery the API returns the original charge and no second charge exists.

How do I make this job's effect happen once? A decision on Where does the job's effect land with 3 outcomes. How do I make this job's effect happen once? Where does the job's effect land? in my database guard row, same transaction exact an API with idempotency keys pass the job key provider dedups an API without keys started/done + reconcile narrow window remains The closer the guard is to the effect, the smaller the window for duplicates.

4. Prefer naturally idempotent operations

Some operations are idempotent by construction, and rewriting a job in those terms removes the need for a guard:

# not idempotent: each run adds again
await conn.execute("update accounts set balance = balance + $1 where id = $2", amount, account)

# idempotent: setting a value, not incrementing it
await conn.execute("update profiles set avatar_url = $1 where user_id = $2", url, user_id)

# idempotent: upsert keyed on the business identity
await conn.execute(
    """insert into shipments(order_id, carrier, tracking) values($1,$2,$3)
       on conflict (order_id) do update set carrier = excluded.carrier, tracking = excluded.tracking""",
    order_id, carrier, tracking,
)

"Set to X", "upsert by natural key" and "delete if exists" can run any number of times with one outcome; "add X", "append", "send" cannot. When a job mixes both kinds, guard the whole job — one non-idempotent step makes the job non-idempotent.

Verify: run each job twice in a row in a test; the database state after two runs equals the state after one.

Which operations are safe to repeat? A grid of 6 rows by 3 columns. Which operations are safe to repeat? operation run twice gives needs a guard set a value same state no upsert by natural key same row no delete if exists same state no increment a counter double count yes append a row or message two copies yes send an email or charge two side effects yes, or an API key Rewrite steps as set or upsert where you can; guard the rest.

5. Expire the guard table, carefully

The processed-jobs table grows by one row per job. Keep it bounded, but keep rows long enough to cover the longest possible redelivery:

async def prune_processed(pool, keep_days: int = 14) -> int:
    return await pool.fetchval(
        "with d as (delete from processed_jobs where done_at < now() - make_interval(days => $1) "
        "returning 1) select count(*) from d", keep_days)

The retention must exceed the maximum time a job can sit in a queue or be retried plus any manual replay window — dead-letter replays days later included. Deleting a key and then replaying its job runs it again. Two weeks is a common, conservative choice; size it from your actual maximum retry horizon and replay policy. A periodic job that prunes in batches avoids long-running deletes on a busy table, as with the scheduler in scheduling cron jobs inside an asyncio service.

Verify: the table size is stable over weeks, and replaying a dead-lettered job within the retention window reports a duplicate.

Verification

Jobs are idempotent when:

  • Every job has a business-derived key, stable across enqueues and deliveries.
  • Database effects commit together with the guard row.
  • External calls pass idempotency keys, or are reconciled when they cannot.
  • Repeated and concurrent delivery in tests produces exactly one effect per key.

Diagnostic Hook: count "duplicate" outcomes per job type as a metric. A steady trickle is normal at-least-once behaviour; a spike means mass redelivery — a worker fleet crash, a visibility timeout shorter than the job, or a producer enqueueing in a retry loop — and is worth investigating even though nothing was double-applied.

Pitfalls & edge cases

  • Guard row committed separately from the effect. A crash between them either loses the effect or skips it forever.
  • Random job ids as keys. Producer-side duplicates get different keys and both run.
  • Checking the guard before a long external call without holding it. Two deliveries both see "not done" and both call; insert the guard in the same transaction that records the result, or use the API's key.
  • Pruning too early. A replay after the guard expired runs the job again.

Frequently Asked Questions

Why do background jobs run more than once?

Job systems deliver at least once: a worker can crash after doing the work but before acknowledging it, retries can repeat a job whose first attempt succeeded, and visibility timeouts can hand a slow job to a second worker.

How do I make a database job idempotent?

Insert a row keyed by the job's business key into a processed-jobs table with on conflict do nothing, in the same transaction as the job's writes, and skip the work when the insert returns nothing. In testing, 300 concurrent deliveries of 100 jobs produced exactly 100 effects.

How do I make a job that calls an external API idempotent?

Pass a stable idempotency key derived from the job to APIs that support one, so repeated calls return the original result. Otherwise record started and done states and reconcile with the external system before retrying.

What should the idempotency key be?

Something derived from the business operation, such as charge:order-8812, so the same logical work produces the same key whether it was enqueued once or several times. Not a random or time-based id.