Course Content
LangGraph Agents
7 sections · 49 lessons
How do you implement parallel work (fan-out research) and aggregation (fan-in synthesis)?
What you need to know
Python
1import operator2from typing import Annotated, TypedDict3from langgraph.types import Send45class State(TypedDict):6 question: str7 subtasks: list[str]8 findings: Annotated[list[dict], operator.add]9 failures: Annotated[list[dict], operator.add]10 report: str1112def plan(state: State) -> dict:13 return {"subtasks": decompose(state["question"])[:8]} # bound the width1415def fan_out(state: State) -> list[Send]:16 return [Send("worker", {"subtask": s}) for s in state["subtasks"]]1718def worker(inp: dict) -> dict:19 try:20 return {"findings": [{"subtask": inp["subtask"], "text": research(inp["subtask"])}]}21 except SearchUnavailable as e:22 return {"failures": [{"subtask": inp["subtask"], "error": str(e)}]}2324def synthesise(state: State) -> dict:25 ordered = sorted(state["findings"], key=lambda f: f["subtask"])26 return {"report": write_report(ordered, missing=state["failures"])}Wiring: plan → (fan_out) → worker → synthesise, with a worker RetryPolicy for transient errors and {"max_concurrency": 5} in the run config.
Production checklist
- Reducers on
findingsandfailures, orInvalidUpdateError. - Catch inside the worker for expected errors; let retries handle transient ones first.
- Join once — all
Sendworkers run in one super-step, sosynthesiseruns once; if workers are multi-step subgraphs of different lengths, adddefer=True. - Stable order — sort before synthesis so the same inputs give the same report.
- Bounded width — cap subtasks, and use
max_concurrency. - Stream progress —
stream_mode="updates"shows each worker finishing, which makes a 40-second run feel responsive.
A real-life example
An equity research team's agent answers "Compare the five largest paint companies on margins, debt and capacity plans." The planner makes 15 subtasks (5 companies times 3 topics). With max_concurrency: 5, all 15 finish in about 35 seconds; sequentially it took 3 minutes. One day the filings API was down for one company: its three workers returned failures, and the report included a clear note, "capacity data for Company D unavailable", instead of failing the whole run or inventing numbers.
Follow-up questions to expect
- "How do you stop one slow worker from delaying everything?" — Give workers a timeout (node
timeout=for async nodes, or a client timeout) and treat a timeout as a failure entry. - "Can workers be full agents?" — Yes;
Sendcan target a compiledcreate_agentsubgraph. Keep their outputs short and structured. - "What if the fan-out list is huge?" — Batch it, or process in waves with a loop, rather than launching hundreds of workers at once.