Skip to content

Chaining and Grouping Background Jobs

Background work often comes in steps — resize five images, then publish the album — and the obvious way to write it is a parent job that enqueues the children and waits for their results. Inside a worker with a limited number of concurrent jobs, that pattern can wait forever, and the usual alternative — children counting down to a final step — miscounts when the queue delivers a job twice. Measured with taskiq 0.13.0 and a Redis stream broker, ten album jobs each spawning five 200 ms resize jobs: with the worker limited to 10 concurrent jobs, ten parents that awaited their children's results occupied all ten slots, the children never started, and 0 albums finished before the 30-second timeouts expired. With 20 slots the same code finished in 1.4 s — it worked only because there was room. Rewriting it as fan-out with fan-in — each child records itself in a Redis set and the last one enqueues the next step — published all 10 albums in 1.06 s with 10 slots. And under simulated at-least-once delivery, with 10% of children delivered twice, fan-in by counter fired the final step early for 26 of 100 albums; fan-in by set fired it early for 0. This guide composes jobs safely.

Prerequisites

1. See how waiting parents deadlock a worker

A parent job that waits for its children holds a worker slot for the whole time its children run — and the children need slots too:

@broker.task
async def album_waiting(album: int) -> str:
    children = [await resize.kiq(album * 100 + i) for i in range(5)]
    results = [await c.wait_result(timeout=30) for c in children]   # holds this slot meanwhile
    return f"album {album}: {len(results)} images"

Measured with --max-async-tasks 10: ten parents filled all ten slots, every one of them waiting on children that could not be scheduled, and no album finished — the parents' 30-second wait_result timeouts expired and their results came back as errors. With 20 slots, the same ten albums finished in 1.4 s. That is the danger: the pattern works whenever there is spare capacity, in development and at low load, and locks up when a burst of parents arrives, exactly when throughput matters. The same deadlock appears with any bounded pool whose jobs wait on other jobs from the same pool, as in walking trees concurrently with bounded fan-out.

Verify: no job waits on the results of other jobs that run in the same worker pool.

10 albums x 5 resize jobs of 200 ms A grid of 3 rows by 4 columns. 10 albums x 5 resize jobs of 200 ms pattern worker slots albums finished time parent awaits children's results 10 0 of 10 (deadlock) 30 s timeouts parent awaits children's results 20 10 of 10 1.4 s fan-out, set-based fan-in, enqueue next 10 10 of 10 1.06 s taskiq 0.13.0, RedisStreamBroker, one worker process.

2. Chain by enqueueing the next step, not by waiting

The fix for a sequence is continuation: each job, when it finishes, enqueues the next step with whatever it produced, and returns:

@broker.task
async def fetch_report(report_id: int) -> None:
    data = await download(report_id)
    key = await store_blob(data)
    await render_report.kiq(report_id, key)             # next step; this job ends now

@broker.task
async def render_report(report_id: int, blob_key: str) -> None:
    pdf = await render(await load_blob(blob_key))
    await email_report.kiq(report_id, await store_blob(pdf))

No job holds a slot while another runs, so a chain needs one slot at a time however long it is. Pass references to stored data between steps rather than large payloads, because every argument is serialised into the queue. If a step fails and is retried, it must not enqueue the next step twice — make "enqueue next" part of what the idempotency check covers, or let the next step deduplicate on its arguments.

Verify: a chain of N steps completes with a worker limited to one concurrent job.

3. Group with fan-out and record completions idempotently

For "run these in parallel, then continue", fan out the children, and let each record its completion in a shared place; the child that completes the set enqueues the next step:

@broker.task
async def album_fan_out(album: int) -> None:
    for i in range(5):
        await resize_and_report.kiq(album, i, 5)

@broker.task
async def resize_and_report(album: int, image: int, total: int) -> None:
    await resize(album, image)
    added = await redis.sadd(f"album:{album}:done", image)        # 1 only the first time
    if added and await redis.scard(f"album:{album}:done") == total:
        if await redis.set(f"album:{album}:finalised", 1, nx=True):
            await publish_album.kiq(album)                         # exactly once

Measured with 10 slots: all 10 albums published in 1.06 s, against a deadlock for the waiting version. The set makes completion idempotent — a child that runs twice adds nothing the second time — and the SET NX finaliser makes sure only one child enqueues the next step even if two notice completion at once. Give the keys an expiry so finished groups do not accumulate in Redis.

Verify: each group's next step is enqueued exactly once, after all its children have completed.

Fan-out and fan-in for one album A sequence of 5 messages between 4 participants. Fan-out and fan-in for one album album job resize jobs Redis set publish job kiq x5, then end SADD image (x5) SCARD == 5 and SET finalised NX kiq publish_album duplicate delivery: SADD -> 0, stop No job waits for another; completion is recorded, not awaited.

4. Do not count completions with a counter

The tempting fan-in is a counter: each child increments, and whoever reaches the total continues. Under at-least-once delivery, a child that runs twice increments twice:

if await redis.incr(f"album:{album}:count") == total:    # wrong under redelivery
    await publish_album.kiq(album)

Measured in a simulation of 100 albums of five children each, with 10% of children delivered twice — 552 deliveries for 500 children: the counter reached five before all five distinct images had finished for 26 albums, firing the next step early; the set-based version fired early for none. Both published each album once, but a quarter of the counter-based albums were published with images still missing. Count distinct members, not events, whenever deliveries can repeat — which, for every broker in choosing between Celery, arq and taskiq, they can.

Verify: a test that delivers some children twice still continues only after every distinct child has finished.

Albums whose next step fired before all children finished 2 horizontal bars comparing INCR counter reaches 5 with the others. Albums whose next step fired before all children finished INCR counter reaches 5 26 of 100 SADD set reaches 5 members 0 of 100 552 deliveries for 500 children; simulated at-least-once delivery. Duplicates count twice in a counter and once in a set.

5. Handle failed children and stuck groups

A group whose child fails permanently never completes. Decide what that means — fail the group, continue with what succeeded, or retry — and make it visible:

@broker.task
async def resize_and_report(album: int, image: int, total: int) -> None:
    try:
        await resize(album, image)
        key = f"album:{album}:done"
    except PermanentError:
        key = f"album:{album}:failed"                   # recorded, still counts toward the group
    await redis.sadd(key, image)
    finished = await redis.scard(f"album:{album}:done") + await redis.scard(f"album:{album}:failed")
    if finished == total and await redis.set(f"album:{album}:finalised", 1, nx=True):
        await finish_album.kiq(album)                    # inspects done vs failed

Retries happen before a failure is recorded as permanent, as in retrying failed background jobs with backoff. Separately, give each group a deadline: a periodic sweep can find groups whose keys are older than the expected completion time and neither finalised nor failed — a child lost to a bug, a queue purge, a mis-deploy — and report or re-drive them. Without the sweep, a stuck group is invisible, because nothing is waiting for it.

Verify: a group with one permanently failing child finalises with that failure recorded, and stuck groups are reported by a sweep.

How should these jobs be composed? A decision on What shape is the work with 4 outcomes. How should these jobs be composed? What shape is the work? a sequence each job enqueues the next 1 slot at a time parallel, then one step fan-out + set fan-in + NX finaliser 10 albums in 1.06 s children may fail record failures in the group it still finalises parent awaits children avoid in a bounded pool 0 of 10 at 10 slots Record completion; never wait for it inside a worker.

Verification

Job composition is safe when:

  • No job waits on other jobs' results inside the same worker pool.
  • Sequences continue by enqueueing the next step, idempotently.
  • Groups fan in through a set of distinct completions and a single NX finaliser.
  • Failed children are recorded as part of the group, and a sweep reports groups that never finalise.

Diagnostic Hook: export, per job type, how long jobs spend between start and finish while not using CPU. Jobs whose time is dominated by waiting on other jobs — wait_result in a profile, or durations that track their children's — are the waiting-parent pattern, and are a deadlock away from the next traffic burst.

Pitfalls & edge cases

  • Parents awaiting children in the same pool. Measured: 0 of 10 albums with 10 slots.
  • "Works with more slots." Measured: the same code finished in 1.4 s with 20 — until a bigger burst.
  • Counter-based fan-in. Measured: 26 of 100 albums continued early.
  • Groups without a deadline. A lost child leaves them stuck forever, silently.

Frequently Asked Questions

Why do my taskiq jobs that wait for other jobs hang?

A parent awaiting children's results holds a worker slot; with all slots taken by waiting parents, children cannot start. Ten parents with 10 slots finished 0 albums in testing.

How do I run jobs in parallel and then continue in a task queue?

Fan out the children, have each add itself to a Redis set when done, and let the child that completes the set enqueue the next step behind a SET NX guard; 10 albums published in 1.06 s this way.

Why not use a counter for fan-in?

Redelivered jobs increment it twice: with 10% duplicate deliveries, a counter fired the next step early for 26 of 100 groups, while a set never did.

How do I chain background jobs?

Have each job enqueue the next step when it finishes, passing references to stored data, so no job waits on another.