Course Content
Multi-Agent Systems and Collaboration
4 sections · 12 lessons
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.
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.
| Coupling | Direct call | Queue |
|---|---|---|
| Time | Both must be running now | Producer enqueues; consumer runs whenever |
| Rate | Caller's rate is callee's rate | Bursts absorb into the backlog |
| Failure | Callee's crash is the caller's exception | Task stays queued and is retried |
| Capacity | Scale together | Add workers without touching the producer |
Rewritten with a queue, the webhook handler becomes:
1@app.post("/webhook/generate-docs", status_code=202) # FastAPI2def generate_docs(payload: dict):3 job = queue.enqueue(run_docs_agent, payload["repo"], job_timeout=900)4 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:
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%:
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.
pip install rq redisredis-server & # or point at an existing instancerq worker docs default # one worker process, listens to two queues1from redis import Redis2from rq import Queue, Retry34redis_conn = Redis(host="localhost", port=6379)5docs_q = Queue("docs", connection=redis_conn)67def run_docs_agent(repo_url: str) -> dict:8 """An ordinary function. RQ does not care that it contains an agent."""9 agent = DocsAgent(model="…", tools=[read_repo, write_file])10 return agent.run({"repo": repo_url})1112job = docs_q.enqueue(13 run_docs_agent, "https://github.com/acme/api",14 job_timeout=900, # kill the job after 15 minutes15 result_ttl=86_400, # keep the result for a day16 retry=Retry(max=3, interval=[10, 60, 300]),17)1819job.get_status() # 'queued' -> 'started' -> 'finished' | 'failed'20job.result # None until finishedThree 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:
1def run_docs_agent(repo_url: str, request_id: str) -> dict:2 key = f"docs:done:{request_id}"3 if redis_conn.exists(key):4 return json.loads(redis_conn.get(key)) # already done; return it56 result = DocsAgent(...).run({"repo": repo_url})78 # Claim and store atomically-ish: write result, then mark done.9 redis_conn.set(key, json.dumps(result), ex=604_800)10 publish_docs(result) # the side effect, once11 return resultThe 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.
1from celery import Celery2from celery.exceptions import SoftTimeLimitExceeded34app = Celery("agents",5 broker="redis://localhost:6379/0",6 backend="redis://localhost:6379/1")78app.conf.update(9 task_acks_late=True, # ack AFTER the task finishes10 worker_prefetch_multiplier=1, # do not hoard tasks per worker11 task_time_limit=900, # hard kill12 task_soft_time_limit=840, # raises SoftTimeLimitExceeded first13 task_routes={14 "agents.research.*": {"queue": "research"},15 "agents.gpu.*": {"queue": "gpu"},16 },17)1819@app.task(bind=True, max_retries=3,20 autoretry_for=(TimeoutError, RateLimitError),21 retry_backoff=True, retry_jitter=True)22def research_agent(self, topic: str) -> dict:23 agent = ResearchAgent()24 try:25 return agent.run(topic)26 except SoftTimeLimitExceeded:27 # the agent keeps its findings on itself as it goes28 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.
1from celery import chain, group, chord23# chain: sequential pipeline, each result feeds the next4pipeline = chain(5 fetch_agent.s("https://acme.com/report.pdf"),6 extract_agent.s(),7 summarise_agent.s(),8)9pipeline.apply_async()1011# group: fan out, run concurrently, collect a list of results12fan_out = group(research_agent.s(t) for t in13 ["pricing", "competitors", "regulation"])1415# chord: a group plus a callback that runs once, when ALL members finish16workflow = chord(fan_out)(synthesise_agent.s())17workflow.get(timeout=600)| Primitive | Shape | Multi-agent equivalent |
|---|---|---|
chain | A → B → C | A handoff pipeline |
group | A, B, C concurrently | Parallel specialists with no join |
chord | (A, B, C) → D | Fan 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.
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@%hTwelve 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
| Dimension | RQ | Celery |
|---|---|---|
| Brokers | Redis only | Redis, RabbitMQ, SQS, others |
| Setup effort | Minutes | Hours, and a config file you will revisit |
| Composition | Job dependencies (depends_on) | chain, group, chord, chunks, map |
| Scheduling | Via rq-scheduler | Built in (celery beat) |
| Routing | Queue name at enqueue time | Rule-based routing, priorities, exchanges |
| Monitoring | rq-dashboard | Flower, events, plus broker tooling |
| Windows support | Poor (fork-based) | Not officially supported |
| Failure surface | Small enough to read the source | Large; 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.
| Metric | Healthy | What a bad value means |
|---|---|---|
| Queue depth | Near zero, spiky | Rising steadily: capacity is below arrival rate |
| Oldest job age | Under a few multiples of service time | High with low depth: a poison job is stuck at the head |
| Active worker count | Equals what you deployed | Lower: workers are crashing or being OOM-killed |
| Failed-registry size | Flat | Growing: 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.
1import time2from rq import Queue, Worker3from rq.registry import FailedJobRegistry, StartedJobRegistry45def queue_health(conn, name="docs") -> dict:6 q = Queue(name, connection=conn)7 jobs = q.get_jobs(0, 1) # just the head of the queue8 oldest_age = (time.time() - jobs[0].enqueued_at.timestamp()) if jobs else 0.09 return {10 "depth": len(q),11 "oldest_age_s": round(oldest_age, 1),12 "running": len(StartedJobRegistry(name, connection=conn)),13 "failed": len(FailedJobRegistry(name, connection=conn)),14 "workers": len([w for w in Worker.all(connection=conn)15 if name in w.queue_names()]),16 }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.