Fan-out / fan-in¶
Run independent workers in parallel, then combine what they found.
Fan-out sends the same task to several workers at once, each with a different angle (here benefits, drawbacks and risks), and fan-in combines what they return into one answer. The parallelism is asyncio.gather in your code, not a model decision. A worker that fails doesn't sink the run: the others' work is kept, the failure is reported, and the answer is built from what succeeded.
Use it when
- The subtasks are independent: different perspectives, sources or chunks of a document.
- Wall-clock time matters.
- One worker failing shouldn't sink the whole answer.
Look elsewhere when
- Each step needs the one before it:
pipeline. - Which workers run depends on the input:
routerorplanner_executor.
What it shows
- The fan-out is
asyncio.gather— ordinary Python, not an LLM decision return_exceptions=True, so one failed worker doesn't discard its siblings' paid-for work; failures are reported in the output (perspectives_failed), andAllWorkersFailedErroris raised only if nothing succeeded- One
RunUsageshared across the parallel runs, soUSAGE_LIMITSbounds all of them - Trace labels
<name>.workerand<name>.aggregator
See it run: sample_run.md is a recorded run against a real model: what each agent was asked, which tools it called, and what it returned.
run_fan_out returns a RunResult: .output is the summary and which perspectives were used or
failed, and .steps holds the workers that succeeded (in completion order) and then the aggregator.
Adapt it by changing PERSPECTIVES or replacing them with your own subtasks.
examples/fan_out/test_example.py tests the parallelism and the failure cases.
Source¶
All of it is in examples/fan_out/.
"""Fan-out / fan-in: run workers in parallel, then combine their results.
Use this pattern when:
- Independent subtasks can run at the same time (different perspectives, sources, chunks)
- Wall-clock time matters, or one worker's failure shouldn't sink the whole answer
- A final step must reconcile what the workers produced
topic ─┬→ worker (pros) ─┐
├→ worker (cons) ─┼→ aggregator → answer
└→ worker (risks) ─┘
The fan-out is `asyncio.gather` — ordinary Python, not an LLM decision.
"""
from __future__ import annotations
import asyncio
from dataclasses import dataclass
from pydantic import BaseModel
from pydantic_ai import Agent
from pydantic_ai.capabilities import RaiseContentFilterError
from pydantic_ai.usage import UsageLimits
from agent.config import settings
from agent.logging import agent_label, configure_logging, get_logger
from agent.runs import Flow, RunResult
logger = get_logger(__name__)
LABEL = agent_label(__name__) # names this agent's run spans in Logfire traces
# Parallel runs all draw on one Flow's shared usage, so this bounds the whole fan-out. Size
# request_limit for (workers + aggregator), with headroom for retries.
USAGE_LIMITS = UsageLimits(
request_limit=20, total_tokens_limit=200_000, cost_limit=settings.cost_limit
)
PERSPECTIVES = ["benefits", "drawbacks", "risks"]
@dataclass
class FanOutDeps:
"""Runtime dependencies shared by the workers and the aggregator."""
pass
class Finding(BaseModel):
result: str
worker_agent: Agent[FanOutDeps, Finding] = Agent(
settings.model,
name=f"{LABEL}.worker",
output_type=Finding,
deps_type=FanOutDeps,
capabilities=[RaiseContentFilterError()],
instructions="You analyze a topic from one given perspective, in two or three sentences.",
)
class FanOutOutput(BaseModel):
result: str
perspectives_used: list[str]
perspectives_failed: list[str]
class Summary(BaseModel):
result: str
aggregator_agent: Agent[FanOutDeps, Summary] = Agent(
settings.model,
name=f"{LABEL}.aggregator",
output_type=Summary,
deps_type=FanOutDeps,
capabilities=[RaiseContentFilterError()],
instructions=(
"You are given findings on one topic from several perspectives. Combine them into one "
"balanced summary. Do not mention perspectives that were not provided."
),
)
class AllWorkersFailedError(Exception):
"""Every parallel worker failed, so there is nothing to aggregate."""
async def run_fan_out(user_input: str, deps: FanOutDeps | None = None) -> RunResult[FanOutOutput]:
"""Analyze `user_input` from several perspectives in parallel, then summarize.
One failing worker is tolerated — the summary is built from the rest and the failure is
reported in `perspectives_failed`. If every worker fails there is nothing to summarize.
Returns:
A RunResult: `.output` is the FanOutOutput; `.steps` holds the workers that succeeded
(in completion order) and then the aggregator.
Raises:
AllWorkersFailedError: When no worker produced a finding.
"""
if deps is None:
deps = FanOutDeps()
flow = Flow(USAGE_LIMITS) # shared across the parallel runs, so it bounds all of them
# return_exceptions=True keeps one failure from cancelling its siblings and discarding
# work already paid for.
outcomes = await asyncio.gather(
*(
flow.run(worker_agent, f"Perspective: {perspective}\nTopic: {user_input}", deps=deps)
for perspective in PERSPECTIVES
),
return_exceptions=True,
)
findings: dict[str, str] = {}
failed: list[str] = []
for perspective, outcome in zip(PERSPECTIVES, outcomes, strict=True):
if isinstance(outcome, BaseException):
# Let cancellation and Ctrl-C through; only treat ordinary failures as "a worker failed".
if not isinstance(outcome, Exception):
raise outcome
logger.warning(
"Worker failed", extra={"perspective": perspective, "error": str(outcome)}
)
failed.append(perspective)
else:
findings[perspective] = outcome.output.result
if not findings:
raise AllWorkersFailedError(f"All {len(PERSPECTIVES)} workers failed for: {user_input!r}")
material = "\n\n".join(f"{name}:\n{text}" for name, text in findings.items())
summary = await flow.run(aggregator_agent, f"Topic: {user_input}\n\n{material}", deps=deps)
return flow.finish(
FanOutOutput(
result=summary.output.result,
perspectives_used=list(findings),
perspectives_failed=failed,
)
)
if __name__ == "__main__":
configure_logging()
print(asyncio.run(run_fan_out("Adopting a monorepo")).output)
"""Workers run in parallel, one failure is tolerated, and total failure is reported."""
import asyncio
import pytest
from pydantic_ai.messages import ModelResponse, ToolCallPart, UserPromptPart
from pydantic_ai.models.function import AgentInfo, FunctionModel
from examples.fan_out.agent import (
PERSPECTIVES,
AllWorkersFailedError,
aggregator_agent,
run_fan_out,
worker_agent,
)
def prompt_of(messages) -> str:
return "\n".join(str(p.content) for p in messages[-1].parts if isinstance(p, UserPromptPart))
def summarizer(seen: list[str]):
def model_fn(messages, info: AgentInfo) -> ModelResponse:
seen.append(prompt_of(messages))
return ModelResponse(parts=[ToolCallPart(info.output_tools[0].name, {"result": "summary"})])
return FunctionModel(model_fn)
def workers(fail_on: set[str] = frozenset(), stats: dict | None = None):
"""Each worker answers `finding about <perspective>`, or raises if told to fail."""
async def model_fn(messages, info: AgentInfo) -> ModelResponse:
perspective = next(p for p in PERSPECTIVES if f"Perspective: {p}" in prompt_of(messages))
if stats is not None:
stats["running"] += 1
stats["peak"] = max(stats["peak"], stats["running"])
await asyncio.sleep(0.01) # yield so the other workers can start
stats["running"] -= 1
if perspective in fail_on:
raise RuntimeError(f"{perspective} worker is down")
fields = {"result": f"finding about {perspective}"}
return ModelResponse(parts=[ToolCallPart(info.output_tools[0].name, fields)])
return FunctionModel(model_fn)
async def test_every_perspective_is_analyzed_and_summarized():
seen: list[str] = []
with worker_agent.override(model=workers()), aggregator_agent.override(model=summarizer(seen)):
result = await run_fan_out("a topic")
output = result.output
assert output.result == "summary"
assert output.perspectives_used == PERSPECTIVES
assert output.perspectives_failed == []
# Three workers (in completion order), then the aggregator.
assert [step.agent for step in result.steps] == ["fan_out.worker"] * 3 + ["fan_out.aggregator"]
assert all(f"finding about {p}" in seen[0] for p in PERSPECTIVES)
async def test_the_workers_really_run_at_the_same_time():
stats = {"running": 0, "peak": 0}
with (
worker_agent.override(model=workers(stats=stats)),
aggregator_agent.override(model=summarizer([])),
):
await run_fan_out("a topic")
assert stats["peak"] == len(PERSPECTIVES)
async def test_a_failed_worker_is_reported_and_the_rest_still_count():
seen: list[str] = []
with (
worker_agent.override(model=workers(fail_on={"risks"})),
aggregator_agent.override(model=summarizer(seen)),
):
result = await run_fan_out("a topic")
output = result.output
assert output.perspectives_failed == ["risks"]
assert len(result.steps) == 3 # the two workers that succeeded, and the aggregator
assert output.perspectives_used == ["benefits", "drawbacks"]
assert "finding about risks" not in seen[0]
async def test_when_every_worker_fails_nothing_is_summarized():
def must_not_run(messages, info):
raise AssertionError("the aggregator ran with no findings")
with (
worker_agent.override(model=workers(fail_on=set(PERSPECTIVES))),
aggregator_agent.override(model=FunctionModel(must_not_run)),
pytest.raises(AllWorkersFailedError),
):
await run_fan_out("a topic")
async def test_cancellation_is_not_mistaken_for_a_failed_worker():
"""A cancelled worker is not "a worker failed": the cancellation must propagate."""
async def cancelled(messages, info: AgentInfo) -> ModelResponse:
raise asyncio.CancelledError
with (
worker_agent.override(model=FunctionModel(cancelled)),
aggregator_agent.override(model=summarizer([])),
pytest.raises(asyncio.CancelledError),
):
await run_fan_out("a topic")
"""Live check: the workers run in parallel and the aggregator combines them. `pytest -m eval`."""
import pytest
from evals.trace import traced_run
from examples.fan_out import agent as module
from examples.live_support import assert_every_agent_ran, run_as_script
pytestmark = pytest.mark.eval
async def test_every_perspective_is_analyzed_and_summarized():
traced = await traced_run(module.run_fan_out, "Adopting a monorepo")
output = traced.result.output
assert output.perspectives_used == module.PERSPECTIVES
assert output.perspectives_failed == []
assert len(output.result.split()) >= 15
labels = [step.agent for step in traced.result.steps]
assert labels[-1] == "fan_out.aggregator" # the aggregator runs after every worker
assert labels[:-1] == ["fan_out.worker"] * len(module.PERSPECTIVES)
# One request per step, plus any output retries.
assert traced.result.usage.requests >= len(module.PERSPECTIVES) + 1
assert_every_agent_ran(module, traced.agents_ran)
async def test_the_demo_script_runs():
assert "perspectives_used=" in await run_as_script("examples.fan_out.agent")
Recorded run · gemini-3.1-flash-lite · 4 steps · $0.0006
Recorded 2026-10-07 with google:gemini-3.1-flash-lite · 4 steps · 813 tokens · $0.0006 · 3.1 s.
Model output varies between runs. Regenerate with uv run python scripts/record_example.py fan_out.
Input
Adopting a monorepo
Steps
1. fan_out.worker
137 tokens · $0.0001
Prompt
Perspective: risks Topic: Adopting a monorepo
Output
{
"result": "Adopting a monorepo introduces significant risks regarding tooling complexity and the potential for a single point of failure in the build pipeline. Furthermore, tight coupling between projects can lead to slower CI/CD cycles and increased coordination overhead as the codebase grows."
}
2. fan_out.worker
143 tokens · $0.0001
Prompt
Perspective: drawbacks Topic: Adopting a monorepo
Output
{
"result": "Adopting a monorepo can lead to significant performance bottlenecks, as scaling tooling and CI/CD pipelines becomes increasingly difficult with a growing codebase. Furthermore, it often introduces complex dependency management challenges and risks compromising developer autonomy, as changes in one area can inadvertently disrupt unrelated projects."
}
3. fan_out.worker
136 tokens · $0.0001
Prompt
Perspective: benefits Topic: Adopting a monorepo
Output
{
"result": "Adopting a monorepo enhances developer productivity by facilitating easier code sharing and atomic commits across related projects. This centralized approach also streamlines dependency management and ensures consistent versioning, which significantly reduces the integration challenges typically found in multi-repository setups."
}
4. fan_out.aggregator
397 tokens · $0.0003
Prompt
Topic: Adopting a monorepo
benefits: Adopting a monorepo enhances developer productivity by facilitating easier code sharing and atomic commits across related projects. This centralized approach also streamlines dependency management and ensures consistent versioning, which significantly reduces the integration challenges typically found in multi-repository setups.
drawbacks: Adopting a monorepo can lead to significant performance bottlenecks, as scaling tooling and CI/CD pipelines becomes inc …
Output
{
"result": "Adopting a monorepo offers a centralized approach that can enhance developer productivity by simplifying code sharing, enabling atomic commits, and streamlining dependency management with consistent versioning. These features help reduce the integration challenges often present in multi-repository environments.\n\nHowever, this architecture also presents significant drawbacks and risks. As the codebase grows, organizations may encounter performance bottlenecks and complex dependency management issues. The reliance on a large, shared structure can introduce a single point of failure within build pipelines and lead to slower CI/CD cycles due to increased coordination overhead and tight coupling. Additionally, these complexities may compromise developer autonomy, as changes in one area risk inadvertently disrupting unrelated projects and straining tooling capabilities."
}
Result
run_fan_out(...).output
{
"result": "Adopting a monorepo offers a centralized approach that can enhance developer productivity by simplifying code sharing, enabling atomic commits, and streamlining dependency management with consistent versioning. These features help reduce the integration challenges often present in multi-repository environments.\n\nHowever, this architecture also presents significant drawbacks and risks. As the codebase grows, organizations may encounter performance bottlenecks and complex dependency management issues. The reliance on a large, shared structure can introduce a single point of failure within build pipelines and lead to slower CI/CD cycles due to increased coordination overhead and tight coupling. Additionally, these complexities may compromise developer autonomy, as changes in one area risk inadvertently disrupting unrelated projects and straining tooling capabilities.",
"perspectives_used": [
"benefits",
"drawbacks",
"risks"
],
"perspectives_failed": []
}