Проблема: сколько веток — неизвестно заранее
В уроке про reducers мы уже запускали параллельные ветки: три поиска от START, результаты собирал operator.add. Но там было ровно три узла, прописанных руками на этапе сборки графа. А что, если число веток зависит от данных?
Классический пример: пользователь прислал текст, мы разбили его на главы и хотим просуммировать каждую главу параллельно. Сколько глав — пять? пятьдесят? — выяснится только в рантайме. Прописать пятьдесят узлов summarize_1 ... summarize_50 заранее невозможно и абсурдно.
Статический (из урока про reducers): известное число веток, заданных рёбрами при сборке. Динамический (этот урок): число веток вычисляется во время выполнения по содержимому состояния. Первый делается обычными рёбрами, второй — через Send.
Map-reduce: map → reduce
Паттерн пришёл из распределённых вычислений (Google MapReduce, Hadoop) и состоит из двух фаз:
- Map — применить одну и ту же операцию к каждому элементу независимо и параллельно. Элементы не знают друг о друге.
- Reduce — собрать все частичные результаты в один итог (выбрать лучший, сложить, отранжировать).
В LangGraph фаза map реализуется динамическим fan-out через Send, а fan-in (схождение в reduce) — через узел, который запускается после завершения всех веток, и reducer, накопивший их результаты.
Send: динамический fan-out
Send — это объект-команда «запусти узел X с вот таким состоянием». Возвращают его из функции условного ребра: вместо одной строки-имени узла функция отдаёт список Send. Сколько элементов в списке — столько параллельных веток и создастся. Каждый Send несёт свой кусочек состояния для воркера.
from langgraph.types import Send # (в старых версиях: from langgraph.constants import Send)
def dispatch_to_summaries(state: OverallState):
"""Возвращает СПИСОК Send — по одной ветке на главу."""
return [
Send("summarize", {"chapter": chapter}) # узел "summarize" + его личное состояние
for chapter in state["chapters"]
]
# Подключаем как условное ребро: из узла-источника, функция, список целей
builder.add_conditional_edges("split", dispatch_to_summaries, ["summarize"])
Два важных момента про Send("summarize", {...}):
- Первый аргумент — имя узла-воркера (один и тот же для всех веток).
- Второй аргумент — состояние именно для этой ветки. Оно может иметь другую схему, чем общее состояние графа: воркеру обычно нужен только его элемент (
{"chapter": ...}), а не весь список.
Send возвращают из функции-маршрутизатора условного ребра — там же, где в обычном случае возвращали имя узла. Узел-источник (split) лишь готовит данные; решение «разветвиться на N» принимает функция ребра.
Сбор результатов: reducer + узел reduce
Все ветки summarize работают параллельно и пишут в одно поле общего состояния. Чтобы их результаты не затёрли друг друга, это поле обязано иметь reducer (помнишь урок про operator.add?). Reducer склеивает частичные результаты по мере завершения веток.
Когда все ветки отработали, LangGraph переходит к узлу reduce (fan-in): он читает накопленный список и формирует итог. Связывается это обычным ребром summarize → reduce.
from typing import Annotated, TypedDict
import operator
class OverallState(TypedDict):
chapters: list[str]
summaries: Annotated[list[str], operator.add] # ← сюда стекаются ветки
digest: str # итог reduce
# Состояние ОДНОЙ ветки — своя маленькая схема
class ChapterState(TypedDict):
chapter: str
def summarize(state: ChapterState) -> dict:
# обрабатываем один элемент; возвращаем список из одного результата
return {"summaries": [f"конспект главы: {state['chapter'][:30]}..."]}
def reduce(state: OverallState) -> dict:
# все ветки завершились — собираем итог
return {"digest": f"Сводка из {len(state['summaries'])} глав"}
Если у summaries нет Annotated[..., operator.add], параллельные ветки начнут конфликтовать за поле, и ты получишь InvalidUpdateError либо потеряешь все результаты, кроме одного. Reducer на поле сбора — обязательное условие map-reduce, а не опция.
Полный пример
from typing import Annotated, TypedDict
import operator
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
# --- Состояние ---
class OverallState(TypedDict):
text: str
chapters: list[str]
summaries: Annotated[list[str], operator.add]
digest: str
class ChapterState(TypedDict):
chapter: str
# --- Узлы ---
def split(state: OverallState) -> dict:
chapters = state["text"].split("\n\n") # бьём текст на главы
return {"chapters": chapters}
def summarize(state: ChapterState) -> dict: # MAP-воркер
return {"summaries": [f"конспект: {state['chapter'][:40]}"]}
def reduce(state: OverallState) -> dict: # REDUCE
return {"digest": " | ".join(state["summaries"])}
# --- Функция fan-out: по одному Send на главу ---
def to_summaries(state: OverallState):
return [Send("summarize", {"chapter": c}) for c in state["chapters"]]
# --- Сборка ---
builder = StateGraph(OverallState)
builder.add_node("split", split)
builder.add_node("summarize", summarize)
builder.add_node("reduce", reduce)
builder.add_edge(START, "split")
builder.add_conditional_edges("split", to_summaries, ["summarize"]) # MAP
builder.add_edge("summarize", "reduce") # fan-in
builder.add_edge("reduce", END)
graph = builder.compile()
result = graph.invoke({"text": "Глава 1...\n\nГлава 2...\n\nГлава 3...",
"chapters": [], "summaries": [], "digest": ""})
print(result["digest"])
# конспект: Глава 1... | конспект: Глава 2... | конспект: Глава 3...
# Все три главы обработаны параллельно, reduce собрал их в сводку
Сколько бы глав ни прислали — три или триста — граф не меняется: to_summaries создаст ровно столько Send, сколько глав, все они выполнятся параллельно, а reduce подождёт всех и соберёт итог.
Send против статических рёбер
| Критерий | Статические параллельные рёбра | Send (динамический) |
|---|---|---|
| Число веток | известно при сборке графа | вычисляется в рантайме |
| Состояние ветки | общее состояние графа | своё на каждую ветку |
| Как задаётся | add_edge(START, "a") ×N |
[Send("a", s) for ...] |
| Типичный случай | несколько фиксированных источников | обработать каждый элемент списка |
Send разворачивает все ветки сразу. Если каждая ветка дёргает LLM, сотня глав = сотня одновременных запросов и риск rate-limit. Ограничить число одновременно выполняемых веток можно через config: graph.invoke(inp, {"max_concurrency": 5}) — LangGraph выполнит ветки пачками по пять.
Типичные ошибки
Без Annotated[list, operator.add] на поле, куда пишут ветки, — конфликт записи и InvalidUpdateError. Поле fan-in всегда с reducer'ом.
С operator.add воркер должен вернуть {"summaries": ["x"]} (список), а не {"summaries": "x"}. Иначе reducer попытается сложить список со строкой и упадёт.
Send отдаёт функция условного ребра, а не узел. Если попытаться вернуть список Send из обычного узла как обновление состояния — fan-out не произойдёт.
Сотни одновременных веток с LLM-вызовами упрутся в rate-limit или память. Для тяжёлых map-операций задавай max_concurrency.
Шпаргалка
from langgraph.types import Send
from typing import Annotated, TypedDict
import operator
# 1. Поле сбора — ОБЯЗАТЕЛЬНО с reducer
class State(TypedDict):
items: list
results: Annotated[list, operator.add] # сюда стекаются ветки map
final: str
# 2. MAP-воркер: обрабатывает ОДИН элемент, возвращает список
def worker(s) -> dict:
return {"results": [process(s["item"])]}
# 3. Fan-out: функция ребра возвращает СПИСОК Send
def fan_out(state):
return [Send("worker", {"item": x}) for x in state["items"]]
builder.add_conditional_edges("split", fan_out, ["worker"]) # MAP
builder.add_edge("worker", "reduce") # fan-in
builder.add_edge("reduce", END)
# 4. Контроль параллелизма (тяжёлые ветки)
graph.invoke(inp, {"max_concurrency": 5})
# Правила:
# • Send — из функции ребра, не из узла
# • число Send = число параллельных веток (динамически)
# • у каждого Send своё состояние (своя схема)
# • поле сбора с operator.add; воркер возвращает СПИСОК
Практическое задание
Построй динамический fan-out:
Задание: генератор и отбор идей
- Состояние:
topic: str,ideas: list[str],drafts: Annotated[list[str], operator.add],best: str. - Узел
brainstormкладёт вideasсписок из 3–5 идей по теме. - Через
Sendразверни веткуexpandна каждую идею; воркер пишет развёрнутый вариант вdrafts(списком из одного элемента). - Узел
pick(reduce) выбирает самый длинный черновик вbest. Свяжиexpand → pick → END. - Запусти и проверь, что число веток равно числу идей, а
draftsсобрал их все. - Со звёздочкой: добавь
max_concurrency=2и, если воркер реально зовёт LLM, убедись, что ветки идут пачками, а не все разом.
Что дальше
Это последний урок раздела «Продвинутые возможности». Теперь ты владеешь памятью, сессиями, вложенностью и динамическим параллелизмом — всем, что нужно для нетривиальных агентов. Дальше — раздел про взаимодействие с человеком.