AutoGen Essentials

Course Content

AutoGen Essentials

7 sections · 28 lessons

How do you do parallel work in AutoGen (e.g., multiple researchers) and merge outputs reliably?


Four Goa researchers: in a chat or fanned outInside one SelectorGroupChat• Turns run one after another• Each reads the others' output• About 48 seconds per plan• Contexts mix, tokens multiplyasyncio.gather plus a merger• Fresh agent per sub-question• Same Finding schema with sources• About 16 seconds per plan• Merger has a stated conflict rule
Parallelism is the easy half; the plan is only reliable if the merger knows what to do when two workers disagree.

What you need to know

Sequential by design

In RoundRobinGroupChat and SelectorGroupChat, one agent speaks at a time and everyone sees its message. That is good for discussion, bad for speed: three 10-second research turns take 30 seconds.

Option 1: fan out with asyncio

Python
import asynciofrom pydantic import BaseModelfrom autogen_agentchat.agents import AssistantAgentclass Finding(BaseModel):    topic: str    facts: list[str]    sources: list[str]    confidence: str   # "high" | "medium" | "low"def make_researcher(name: str) -> AssistantAgent:    return AssistantAgent(name, model_client=client, tools=[web_search],                          max_tool_iterations=4, output_content_type=Finding)async def research(subquestions: list[str]) -> list[Finding | Exception]:    limit = asyncio.Semaphore(3)    async def one(i: int, q: str):        async with limit:            result = await make_researcher(f"r{i}").run(task=q)            return result.messages[-1].content     # a Finding object    return await asyncio.gather(*(one(i, q) for i, q in enumerate(subquestions)),                                return_exceptions=True)
  • Each call builds a fresh agent, so contexts never mix; do not run one agent instance on two tasks at once.
  • output_content_type=Finding makes the agent return a StructuredMessage whose content is a validated Finding.
  • The semaphore caps concurrency at 3, protecting your tokens-per-minute limit.
  • return_exceptions=True means one failure does not cancel the others.

Option 2: GraphFlow fan-out

GraphFlow with DiGraphBuilder lets one node fan out to several and join them again (activation_condition="all" waits for every branch, "any" for the first). It keeps everything inside one team with one termination condition, but it is experimental and all messages still land in one shared thread.

Merging reliably

  • Same shape from every worker. Merging JSON with a sources field is easy; merging three essays is where duplicates and contradictions come from.
  • Partition up front. "Flights", "hotels", "weather" — not "research Goa" three times.
  • Conflict rule. Tell the merger what to do when workers disagree: prefer the primary or more recent source, or show both as an open question. Never average two facts.
  • Missing parts. If a worker failed, the merged answer says what is missing.
  • Carry citations through, so each merged claim can be checked.

A real-life example

A travel-planning team answers "Plan a long weekend in Goa, 3 to 5 Oct, for 4 friends, ₹15,000 each." The sequential version asked flights, hotels, weather and events agents in turn: about 48 seconds.

The parallel version fans out four researchers with the Finding schema, then a merger agent builds the plan. Time dropped to about 16 seconds. The first merged plans had a problem: the hotels researcher said "Baga beach shacks open" and the events researcher said "most shacks reopen after 15 Oct". The merger picked one at random. The fix was a merge rule: "When facts conflict, prefer the source with the later date and mention the conflict." The merger now writes "Some beach shacks may still be closed; check before booking," with both links.

One researcher timing out (about 3% of runs) now produces "Weather: not available" instead of a hung request.

Follow-up questions to expect

  • "Why not just put four researchers in a SelectorGroupChat?" — They would run one after another, each reading the others' output, which costs more time and more tokens and mixes their contexts.
  • "How do you handle rate limits when fanning out?" — A semaphore or queue caps concurrent calls; add retries with backoff for 429 errors in the model client or around the call.
  • "What does autogen-core add here?" — Its runtime lets agents publish to a topic and handle messages concurrently, even across processes, for larger fan-out systems.