Task Decomposition at Scale
Learning Outcomes
- Classify a complex goal into functional, data, temporal, or hybrid decomposition strategies
- Model sub-task dependencies as a DAG and compute parallel versus sequential stages
- Build a supervisor that dynamically decomposes work based on complexity estimates
- Pass context between sub-tasks idempotently so retries and partial failures stay safe
- Reassemble distributed sub-results into one coherent output with conflict handling
Lesson Plan
| Segment | Duration | Topic |
|---|---|---|
| Intro | 3 min | Why decomposition is the load-bearing wall of orchestration |
| Concepts | 6 min | The four decomposition strategies |
| Build | 8 min | A typed plan: tasks, dependencies, the DAG |
| Build | 8 min | Dynamic decomposition by complexity |
| Build | 8 min | Scheduling: topological waves and parallelism |
| Build | 7 min | Result reassembly and conflict resolution |
| Wrap-up | 5 min | Depth trade-offs, takeaways, preview |
Before You Begin
Pre-work:
- Complete Lesson 2: Orchestration Patterns — know the supervisor pattern
- Complete Lesson 3: Communication Between Agents — decomposition produces the messages agents exchange
- Be comfortable with planning from the Agentic AI course (Course 04, Lesson 5)
Shopping List:
- Python 3.11+ with
pip install anthropic - An
ANTHROPIC_API_KEYexported in your environment - A graph mental model: DAGs, topological ordering, fan-out/fan-in
A single agent fails complex goals for the reasons in Lesson 1: finite context, no parallelism, one prompt trying to be five specialists at once. Decomposition cuts a goal into units small enough that each fits one agent's context and skill set. Four axes to cut along:
| Strategy | Cut along… | Example | Shape |
|---|---|---|---|
| Functional | Skill / role | extractor → validator → summarizer | Pipeline / supervisor |
| Data | Partition of input | 10,000 invoices into shards | Fan-out / fan-in |
| Temporal | Phase over time | plan → build → test → review | Sequential gated stages |
| Hybrid | Two+ axes at once | per-region × per-doc-type | A DAG, not a line |
Functional is splitting a monolith into microservices by responsibility; data is sharding; temporal is a staged pipeline where each phase gates the next. Real systems are almost always hybrid — fan out searches by source (data), route each through specialist analyzers (functional), merge in a synthesis phase (temporal).
Decomposition output should be a plan — a serializable structure, not free text. A plan is a set of tasks plus dependency edges; the edges form a directed acyclic graph (DAG), and the absence of cycles is what lets you schedule it.
from dataclasses import dataclass, field
@dataclass(frozen=True)
class Task:
id: str # stable, content-derived
agent: str # which specialist handles this
instruction: str # self-contained unit of work
depends_on: tuple = () # ids that must finish first
est_tokens: int = 4000 # complexity estimate, drives budgeting
@dataclass
class Plan:
goal: str
tasks: dict = field(default_factory=dict) # id -> Task
def validate(self) -> None:
for t in self.tasks.values():
for dep in t.depends_on:
if dep not in self.tasks:
raise ValueError(f"{t.id} needs unknown {dep}")
assert_acyclic(self.tasks) # DFS that raises on a back-edge
The acyclicity check is non-negotiable: a cycle is a deadlock a workflow engine runs until your timeout budget is exhausted. Implement assert_acyclic as a DFS coloring nodes white/gray/black that raises the moment it revisits a gray (in-progress) node. Run validate() before any agent fires, and treat the Plan as the contract the supervisor produces, a guardrail inspects, and the scheduler consumes.
At scale you want the supervisor agent to decompose dynamically. Have it emit the plan as structured JSON via tool use, never as prose you parse with a regex. Define a Task JSON Schema (the same fields as Step 2's dataclass, with agent as an enum of your real specialists), wrap it in a tasks array as one tool, and force the model to call it:
from anthropic import Anthropic
client = Anthropic()
def decompose(goal: str, plan_tool: dict) -> list[dict]:
msg = client.messages.create(
model="claude-sonnet-4-6", max_tokens=2048,
tools=[plan_tool],
tool_choice={"type": "tool", "name": "emit_plan"},
system=("Decompose the goal into the smallest set of independent "
"sub-tasks. Prefer data-parallel fan-out. Add a depends_on "
"edge only when a task truly needs another's output."),
messages=[{"role": "user", "content": goal}],
)
block = next(b for b in msg.content if b.type == "tool_use")
return block.input["tasks"]
Forcing tool_choice guarantees valid JSON in your schema; the enum on agent stops the supervisor inventing a specialist you never wired up. Add two guardrails: cap plan size (reject decompositions over ~50 tasks — a runaway planner is a runaway bill) and bound recursion depth to 2–3 levels, since each level of further decomposition multiplies cost.
est_tokens is a guess. Treat it as a prior for budgeting and routing, then reconcile against real usage. Lesson 9 turns these into hard token budgets per agent.Given a validated DAG, the scheduler repeatedly asks: what can run right now? A task is ready when all dependencies have completed. Group ready tasks into waves dispatched concurrently — a level-order (Kahn's algorithm) topological sort, fan out each wave then fan in before the next.
import asyncio
def waves(plan: Plan) -> list[list[str]]:
remaining, done, out = dict(plan.tasks), set(), []
while remaining:
ready = [tid for tid, t in remaining.items()
if set(t.depends_on) <= done]
if not ready:
raise RuntimeError("no progress — cycle or missing dep")
out.append(ready)
for tid in ready:
del remaining[tid]
done.update(ready)
return out
async def run_plan(plan: Plan) -> dict:
ctx = {}
for wave in waves(plan):
async def run_one(tid):
t = plan.tasks[tid]
inputs = {d: ctx[d] for d in t.depends_on} # only what it needs
return tid, await execute_agent(t, inputs)
pairs = await asyncio.gather(*(run_one(t) for t in wave),
return_exceptions=True)
ctx.update({k: v for k, v in pairs if not isinstance(v, Exception)})
return ctx
Each agent receives only its declared dependencies' outputs, not the whole accumulated context — a fan-out task over invoice #4,712 never sees invoice #11's analysis, which keeps every prompt small and free of cross-contamination. Because ids are stable, wrap execute_agent in a cache lookup so a retried wave re-runs only the tasks that failed.
asyncio.gather raises on the first exception and cancels the rest. Pass return_exceptions=True so one failed shard does not vaporize the nine that succeeded — full recovery is Lesson 7, but design for it now.Decomposition is half the job; results must come back together, and the reassembly strategy mirrors how you cut the work.
| Decomposition | Reassembly pattern | Watch for |
|---|---|---|
| Functional (pipeline) | Last stage's output is the answer | Lost intermediate evidence |
| Data (fan-out) | Concatenate / aggregate / reduce | Duplicates, conflicts |
| Temporal (phases) | Final phase plus audit trail | Stale data from earlier phases |
Data fan-out is where reassembly gets interesting, because independent agents can disagree — two shards might extract a different "total" for the same record. Run a deterministic merge that unions records and resolves what it can (dedupe, sum, take-latest), then escalate genuine semantic conflicts to a reconciliation agent with just the conflicting values and their provenance:
> You are a reconciliation agent. Two extractors disagree on a field.
> Field: invoice_total
> Value A: 4820.00 (source: shard_03, page 2)
> Value B: 4280.00 (source: shard_07, page 2)
> Return the correct value and a one-line justification. If you cannot
> determine it, return needs_human: true.
The needs_human: true exit is deliberate. Reassembly is where eventual consistency bites — agents ran at different times against possibly different context, so forcing a confident merge on ambiguous data manufactures plausible-looking wrong answers. A clean escalation path is cheaper than a confidently incorrect total.
More decomposition is not strictly better. Each sub-task carries fixed overhead: a fresh prompt, re-sent context, a round-trip, a slice of concurrency budget. Too coarse and an agent chokes on a task too big for its context; too fine and coordination overhead dominates the real work. Encode the stop rule explicitly — a should_decompose(task, depth) predicate that returns False once depth hits a max_depth of 3 or task.est_tokens falls under a single-agent threshold (say 6,000), and only otherwise splits further.
| Signal | Lean coarser | Lean finer |
|---|---|---|
| Task fits one context window | yes | — |
| Sub-tasks embarrassingly parallel | — | yes (fan out) |
| Strong inter-task dependencies | yes | — |
| Need independent retry of pieces | — | yes |
| Different specialist tools per piece | — | yes |
Every dependency edge is a synchronization point and a place a failure can stall the DAG; every wider wave trades latency for throughput against your rate limit. Depth is a tuning knob between coordination cost and single-agent overload — no free lunch, only the trade-off that fits this workload's SLOs.
Questions & Answers
done set, and a failed task never enters it. So failures quarantine their downstream subtree, which stalls rather than corrupts. What to do with the stalled subtree (retry, fallback agent, degraded path, escalate) is Lesson 7. For now, ensure gather uses return_exceptions=True so the rest of the wave completes and you capture which task failed.depends_on edges. Anything an agent can change or that is branch-specific travels through edges; anything global and read-only is ambient. Mixing them causes the blackboard contention covered in Lesson 3.should_decompose). Needing more than three levels usually signals the goal should split into separate workflows that hand off, not one mega-DAG.Key Takeaways
- Four cuts, usually combined — functional (by skill), data (by partition), temporal (by phase), and hybrid. The cut sets your parallelism and your failure surface.
- The plan is a typed DAG, not prose — emit a validated, acyclic, inspectable structure with stable task ids so it can be approved, cached, and replayed.
- Decompose dynamically with forced tool use — size the plan via a JSON tool schema, then cap plan size and recursion depth to bound cost.
- Schedule in topological waves — a task is ready when its dependencies finish; dispatch each wave concurrently, pass each agent only its declared inputs, and bound wave width.
- Reassembly mirrors decomposition — pipelines take the last stage, fan-outs reduce with explicit conflict handling, and semantic conflicts escalate to a reconciliation agent or a human.
- Depth is a trade-off, not a virtue — the right plan is the minimum decomposition that keeps every task inside one agent's context and competence.
Next Steps: Lesson 5: The Anthropic Agent SDK