Skip to content

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

topic ─┬→ worker (benefits)  ─┐
       ├→ worker (drawbacks) ─┼→ aggregator → answer
       └→ worker (risks)     ─┘

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), and AllWorkersFailedError is raised only if nothing succeeded
  • One RunUsage shared across the parallel runs, so USAGE_LIMITS bounds all of them
  • Trace labels <name>.worker and <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.

uv run python scripts/add_agent.py fan_out --name analysis

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)
title = "Fan-out / fan-in"
pattern = "fan_out"
summary = "Run workers in parallel with asyncio.gather, tolerate a failed worker, then aggregate their findings."
smoke_input = "Adopting a monorepo"

[entrypoint]
deps = "FanOutDeps"
run = "run_fan_out"
"""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": []
}