Multi-Agent Systems and Collaboration

Task Schedulers & Queues for Multi-Agent Systems


A SaaS company shipped an agent that generated onboarding documentation from a customer's codebase. It ran inside the HTTP request handler: the webhook arrived, the agent ran, the response returned. In testing the agent took 8 seconds. In production, on real repositories, it averaged 45.

Their load balancer had a 30-second idle timeout. So at second 30 the caller got a 504, and the caller — a CI system — retried. Meanwhile the original agent was still running, because nothing had told it to stop. Two agents, same repository, both writing to the same output path. On a busy Tuesday this produced 3,400 duplicate runs, roughly 40% of the day's model spend, and a documentation file whose contents depended on which of two agents finished last.

The team's first fix was to raise the load balancer timeout to 120 seconds. That stopped the 504s and introduced a worse failure: web workers now sat blocked for two minutes each, the pool of 20 exhausted at 10 concurrent requests, and healthy fast endpoints started timing out because there was no worker free to serve them.

The problem was never the timeout. It was that a 45-second job was being run inside a request that had to answer in under 30. The structural fix is to stop running it there.

Getting the 45-second agent out of the requestWebhook arrives,returns immediatelyJob enqueuedwith anidempotency keyWorker poolpulls, sizedto arrival rateAgent runs,retriessafely on failureResult stored,caller notifiedEight seconds in testing and 45 in production is the same code meeting real repositories.
A queue converts a timeout into a backlog you can see, size and drain — which is the only kind you can fix.

Why a queue and not a function call

A function call couples caller and callee in four ways at once, and a queue breaks all four.

CouplingDirect callQueue
TimeBoth must be running nowProducer enqueues; consumer runs whenever
RateCaller's rate is callee's rateBursts absorb into the backlog
FailureCallee's crash is the caller's exceptionTask stays queued and is retried
CapacityScale togetherAdd workers without touching the producer

Rewritten with a queue, the webhook handler becomes:

Python
@app.post("/webhook/generate-docs", status_code=202)   # FastAPIdef generate_docs(payload: dict):    job = queue.enqueue(run_docs_agent, payload["repo"], job_timeout=900)    return {"job_id": job.id, "status_url": f"/jobs/{job.id}"}

That returns in about 4 milliseconds. The 504s are gone, the retries are gone, and the agent now has 900 seconds instead of 30. The caller polls the status URL, or you call it back when the job finishes.

Sizing the worker pool

Queues make capacity a calculation. Little's Law says the average number of jobs in a system equals the arrival rate multiplied by the average time each spends there:

L=λWL = \lambda W

With arrivals at λ = 2 jobs per second and a service time of W = 6 seconds, you need L = 12 jobs being processed concurrently just to keep pace. Twelve workers gives you exactly 100% utilisation — which is the worst place to sit, because any variance in arrivals makes the queue grow and it never recovers. Target 70–80%:

workers=λWρ=120.8=15\text{workers} = \frac{\lambda W}{\rho} = \frac{12}{0.8} = 15

Run 15. The extra three are not waste; they are the margin that keeps queue depth near zero instead of climbing. And notice how visible the failure becomes if you run 11 instead of 12: capacity is 11/6 = 1.83 jobs per second against 2 arriving, so the backlog grows by one job every six seconds (about 0.17 per second) — 600 jobs per hour, 14,400 over a day. Nothing errors. Queue depth just climbs.

Run a queue at 100% utilisation and it never drains. Capacity planning is not "enough to keep up"; it is "enough to catch up".

Redis Queue: the small option

RQ is a few hundred lines of surface area over Redis lists. If you already run Redis, you can have a working queue in ten minutes.

Bash
pip install rq redisredis-server &              # or point at an existing instancerq worker docs default      # one worker process, listens to two queues
Python
from redis import Redisfrom rq import Queue, Retryredis_conn = Redis(host="localhost", port=6379)docs_q = Queue("docs", connection=redis_conn)def run_docs_agent(repo_url: str) -> dict:    """An ordinary function. RQ does not care that it contains an agent."""    agent = DocsAgent(model="…", tools=[read_repo, write_file])    return agent.run({"repo": repo_url})job = docs_q.enqueue(    run_docs_agent, "https://github.com/acme/api",    job_timeout=900,                        # kill the job after 15 minutes    result_ttl=86_400,                      # keep the result for a day    retry=Retry(max=3, interval=[10, 60, 300]),)job.get_status()     # 'queued' -> 'started' -> 'finished' | 'failed'job.result           # None until finished

Three arguments there deserve attention because their defaults will hurt you.

job_timeout defaults to 180 seconds. An agent that legitimately takes 400 seconds will be killed mid-run, and the symptom — jobs failing at almost exactly three minutes, with no exception from your code — is confusing the first time you meet it. Set it to roughly twice your 99th-percentile duration.

result_ttl defaults to 500 seconds. If your caller polls every ten minutes, the result is already gone and the job appears never to have existed.

Retry(interval=[10, 60, 300]) gives explicit backoff between attempts. A bare Retry(max=3) retries immediately, which is exactly wrong when the cause is a rate limit — three instant retries make the rate limiter angrier and burn three times the tokens to fail three times.

Making an agent safe to retry

Retries only help if running the task twice is harmless. Agents usually have side effects, so this needs deliberate work:

Python
def run_docs_agent(repo_url: str, request_id: str) -> dict:    key = f"docs:done:{request_id}"    if redis_conn.exists(key):        return json.loads(redis_conn.get(key))     # already done; return it    result = DocsAgent(...).run({"repo": repo_url})    # Claim and store atomically-ish: write result, then mark done.    redis_conn.set(key, json.dumps(result), ex=604_800)    publish_docs(result)                            # the side effect, once    return result

The request_id comes from the caller — a webhook delivery ID, say — not generated inside the task, or every retry gets a fresh one and the guard does nothing. This is the single most common mistake in queued agent work.

Celery: the full option

Celery gives you routing, scheduling, and composition. That is worth real complexity when you need it and pure overhead when you do not.

Python
from celery import Celeryfrom celery.exceptions import SoftTimeLimitExceededapp = Celery("agents",             broker="redis://localhost:6379/0",             backend="redis://localhost:6379/1")app.conf.update(    task_acks_late=True,              # ack AFTER the task finishes    worker_prefetch_multiplier=1,     # do not hoard tasks per worker    task_time_limit=900,              # hard kill    task_soft_time_limit=840,         # raises SoftTimeLimitExceeded first    task_routes={        "agents.research.*": {"queue": "research"},        "agents.gpu.*":      {"queue": "gpu"},    },)@app.task(bind=True, max_retries=3,          autoretry_for=(TimeoutError, RateLimitError),          retry_backoff=True, retry_jitter=True)def research_agent(self, topic: str) -> dict:    agent = ResearchAgent()    try:        return agent.run(topic)    except SoftTimeLimitExceeded:        # the agent keeps its findings on itself as it goes        return {"partial": True, "notes": agent.notes_so_far}

The two settings at the top matter more than everything else on the page. task_acks_late=True means a task is acknowledged after it completes, so a worker killed mid-task returns the task to the queue instead of losing it — with the default of early acknowledgement, a worker crash silently drops in-flight work. worker_prefetch_multiplier=1 stops a worker grabbing four tasks and sitting on three of them while it runs the first; with long agent tasks, prefetching creates the strange situation where one worker holds a queue of work while another is idle.

The soft time limit is the agent-specific one. A hard kill loses everything; SoftTimeLimitExceeded is raised inside your code 60 seconds earlier, giving the agent time to return whatever it has found so far. For a research agent that has gathered eight of twelve sources, partial results are worth a great deal more than nothing.

Composing agent workflows

Celery's canvas gives you three composition primitives, and they map cleanly onto multi-agent shapes.

Python
from celery import chain, group, chord# chain: sequential pipeline, each result feeds the nextpipeline = chain(    fetch_agent.s("https://acme.com/report.pdf"),    extract_agent.s(),    summarise_agent.s(),)pipeline.apply_async()# group: fan out, run concurrently, collect a list of resultsfan_out = group(research_agent.s(t) for t in                ["pricing", "competitors", "regulation"])# chord: a group plus a callback that runs once, when ALL members finishworkflow = chord(fan_out)(synthesise_agent.s())workflow.get(timeout=600)
PrimitiveShapeMulti-agent equivalent
chainA → B → CA handoff pipeline
groupA, B, C concurrentlyParallel specialists with no join
chord(A, B, C) → DFan out to workers, then a synthesiser

The chord's callback receives a list of results in the order the group was defined, not the order they completed — which is what makes it usable, but also means one slow member holds the callback. Give the chord a timeout, and consider whether a partial synthesis from three of four agents is better than nothing.

One sharp edge: if any member of a chord fails, the callback does not run at all — the chord ends in an error, and the three good results never reach the synthesiser. Set max_retries on every group member and make sure a permanently failing member returns an error result rather than raising, so the callback still fires with three good results and one recorded failure.

Routing agents to their own queues

The task_routes block above is the feature that most justifies Celery in a multi-agent system. Different agents have different resource profiles, and mixing them in one queue means the scarce resource gates everything.

Bash
celery -A agents worker -Q research -c 12 -n research@%h   # I/O bound: manycelery -A agents worker -Q gpu      -c 2  -n gpu@%h        # GPU bound: fewcelery -A agents worker -Q default  -c 4  -n default@%h

Twelve concurrent research agents are fine because they spend their time waiting on HTTP. Two GPU agents is the limit because there are two GPUs. In a single shared queue you would have to set concurrency to 2 for everyone, cutting research throughput by a factor of six for no reason.

RQ or Celery

DimensionRQCelery
BrokersRedis onlyRedis, RabbitMQ, SQS, others
Setup effortMinutesHours, and a config file you will revisit
CompositionJob dependencies (depends_on)chain, group, chord, chunks, map
SchedulingVia rq-schedulerBuilt in (celery beat)
RoutingQueue name at enqueue timeRule-based routing, priorities, exchanges
Monitoringrq-dashboardFlower, events, plus broker tooling
Windows supportPoor (fork-based)Not officially supported
Failure surfaceSmall enough to read the sourceLarge; misconfiguration is a real risk

Choose RQ when you have one or two agent task types, a single Redis, and a team that does not want to learn a framework. Choose Celery when you need per-agent queues with different concurrency, scheduled runs, or fan-out-then-join as a first-class construct. The wrong reason to choose Celery is "we might need it later" — the configuration surface is where the outages come from, and acks_late being wrong by default has cost more teams more work than any missing feature in RQ.

Monitoring queue health

Four numbers tell you almost everything, and one of them is usually missing from dashboards.

MetricHealthyWhat a bad value means
Queue depthNear zero, spikyRising steadily: capacity is below arrival rate
Oldest job ageUnder a few multiples of service timeHigh with low depth: a poison job is stuck at the head
Active worker countEquals what you deployedLower: workers are crashing or being OOM-killed
Failed-registry sizeFlatGrowing: a systematic failure, not bad luck

Oldest job age is the one people omit, and it is the most diagnostic. Depth alone cannot distinguish a healthy burst from a stuck head. A depth of 3,000 with an age of 40 seconds means a spike is draining normally. A depth of 120 with an age of 3 hours means one job at the front is failing and being retried forever while everything behind it starves.

Python
import timefrom rq import Queue, Workerfrom rq.registry import FailedJobRegistry, StartedJobRegistrydef queue_health(conn, name="docs") -> dict:    q = Queue(name, connection=conn)    jobs = q.get_jobs(0, 1)                    # just the head of the queue    oldest_age = (time.time() - jobs[0].enqueued_at.timestamp()) if jobs else 0.0    return {        "depth": len(q),        "oldest_age_s": round(oldest_age, 1),        "running": len(StartedJobRegistry(name, connection=conn)),        "failed": len(FailedJobRegistry(name, connection=conn)),        "workers": len([w for w in Worker.all(connection=conn)                        if name in w.queue_names()]),    }

Queue depth tells you how much work is waiting. Oldest-job age tells you whether anything is moving. A short queue with an old head is a stuck job, and depth alone will never show it to you.

Alert on rates of change, not absolute values. "Depth above 1,000" fires on every legitimate burst and gets muted within a week. "Depth increased for 10 consecutive minutes" fires only when capacity is genuinely short. Similarly, "oldest age above 15 minutes" catches the stuck head that depth alone would never reveal.

For agent-specific health, add two more: tokens per job and retry rate. A sudden rise in tokens per job usually means an agent is looping — running its tool loop to the step cap rather than converging — and that shows up in the bill long before it shows up in the failure count. A retry rate above about 10% means you are treating a systematic error as transient.

Named failure modes

Early acknowledgement. The default in Celery is to acknowledge on receipt. A worker killed mid-task — an OOM, a deploy, a spot instance reclaim — loses the task silently. Symptom: jobs disappear during deploys. Fix: task_acks_late=True plus idempotent tasks.

The default timeout. RQ's 180 seconds and Celery's unlimited default are both wrong for agents. Too short kills legitimate long runs; unlimited lets one hung agent hold a worker forever. Set both a soft and a hard limit.

Non-idempotent retries. The agent publishes documentation, then fails on the notification step, then retries the whole task and publishes again. Symptom: duplicate outputs after any transient error. Fix: a claim key derived from a caller-supplied ID.

Retrying non-transient errors. A malformed payload raises ValidationError, gets retried three times with backoff, and fails identically each time, having spent three times the tokens over five minutes. Fix: autoretry_for should list only transient exception types, and everything else should fail immediately to a dead-letter queue.

Prefetch hoarding. Default prefetch multiplied by long tasks means one worker holds a stack of jobs while another idles. Symptom: uneven worker utilisation with a non-empty queue. Fix: worker_prefetch_multiplier=1.

Result backend growth. Agent results are large — full transcripts, source documents. Store them for a week each at a few hundred kilobytes, times thousands of jobs, and your Redis is your database. Fix: store results in object storage and put a URL in the result payload.

The invisible poison job. One task fails, is retried, fails, forever. Depth looks fine. Fix: cap attempts, move exhausted tasks to a dead-letter queue, and alert on oldest-job age.

What this means when you build

Draw the line between "request" and "work" early. Anything an agent does that could take more than a couple of seconds belongs behind a queue, and the endpoint's job is to enqueue, return an identifier, and get out of the way. Retrofitting this later means changing your API contract, which means changing every caller — the CI system in the opening story had to be modified too, and that took longer than the queue did.

Make idempotency a property of the task signature. If a task takes a caller-supplied identifier as its first parameter and guards on it, retries are free. If it does not, every retry is a potential duplicate side effect, and you will find out about it from a customer.

Give each agent type its own queue with its own concurrency, even if today they all run on one worker. Adding a queue name at enqueue time costs nothing now and saves a migration when the GPU agent arrives and cannot share a pool with twelve HTTP-bound research agents.

And put the four health metrics on a dashboard on day one, with oldest-job age given equal prominence to depth. Both of the expensive failures in this lesson — the growing backlog that nobody noticed and the poison job at the head of a short queue — are invisible in every other metric and obvious in these two.