Проблема: сколько веток — неизвестно заранее

В уроке про reducers мы уже запускали параллельные ветки: три поиска от START, результаты собирал operator.add. Но там было ровно три узла, прописанных руками на этапе сборки графа. А что, если число веток зависит от данных?

Классический пример: пользователь прислал текст, мы разбили его на главы и хотим просуммировать каждую главу параллельно. Сколько глав — пять? пятьдесят? — выяснится только в рантайме. Прописать пятьдесят узлов summarize_1 ... summarize_50 заранее невозможно и абсурдно.

ℹ️ Статический fan-out ≠ динамический

Статический (из урока про reducers): известное число веток, заданных рёбрами при сборке. Динамический (этот урок): число веток вычисляется во время выполнения по содержимому состояния. Первый делается обычными рёбрами, второй — через Send.

Map-reduce: map → reduce

Паттерн пришёл из распределённых вычислений (Google MapReduce, Hadoop) и состоит из двух фаз:

  • Map — применить одну и ту же операцию к каждому элементу независимо и параллельно. Элементы не знают друг о друге.
  • Reduce — собрать все частичные результаты в один итог (выбрать лучший, сложить, отранжировать).

В LangGraph фаза map реализуется динамическим fan-out через Send, а fan-in (схождение в reduce) — через узел, который запускается после завершения всех веток, и reducer, накопивший их результаты.

100%
колёсико — масштаб  ·  зажать и тянуть — перемещение
START split даёт N глав MAP · параллельно, N веток summarize #1 summarize #2 summarize #N reduce собрать итог END Send × N reducer число веток N вычисляется в рантайме; reduce запускается после завершения всех веток

Send: динамический fan-out

Send — это объект-команда «запусти узел X с вот таким состоянием». Возвращают его из функции условного ребра: вместо одной строки-имени узла функция отдаёт список Send. Сколько элементов в списке — столько параллельных веток и создастся. Каждый Send несёт свой кусочек состояния для воркера.

Send — по одному на каждый элемент
python
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 живёт в функции ребра, а не в узле

Send возвращают из функции-маршрутизатора условного ребра — там же, где в обычном случае возвращали имя узла. Узел-источник (split) лишь готовит данные; решение «разветвиться на N» принимает функция ребра.

Сбор результатов: reducer + узел reduce

Все ветки summarize работают параллельно и пишут в одно поле общего состояния. Чтобы их результаты не затёрли друг друга, это поле обязано иметь reducer (помнишь урок про operator.add?). Reducer склеивает частичные результаты по мере завершения веток.

Когда все ветки отработали, LangGraph переходит к узлу reduce (fan-in): он читает накопленный список и формирует итог. Связывается это обычным ребром summarize → reduce.

Состояние: воркеры пишут в поле с reducer'ом
python
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'])} глав"}
⚠️ Поле сбора без reducer = потеря данных

Если у summaries нет Annotated[..., operator.add], параллельные ветки начнут конфликтовать за поле, и ты получишь InvalidUpdateError либо потеряешь все результаты, кроме одного. Reducer на поле сбора — обязательное условие map-reduce, а не опция.

Полный пример

Map-reduce целиком: суммаризация N глав → сводка
python
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 выполнит ветки пачками по пять.

Типичные ошибки

Ошибка 1: поле сбора без reducer

Без Annotated[list, operator.add] на поле, куда пишут ветки, — конфликт записи и InvalidUpdateError. Поле fan-in всегда с reducer'ом.

Ошибка 2: воркер возвращает элемент, а не список

С operator.add воркер должен вернуть {"summaries": ["x"]} (список), а не {"summaries": "x"}. Иначе reducer попытается сложить список со строкой и упадёт.

Ошибка 3: Send возвращают из узла

Send отдаёт функция условного ребра, а не узел. Если попытаться вернуть список Send из обычного узла как обновление состояния — fan-out не произойдёт.

Ошибка 4: тяжёлые ветки без max_concurrency

Сотни одновременных веток с LLM-вызовами упрутся в rate-limit или память. Для тяжёлых map-операций задавай max_concurrency.

Шпаргалка

Map-reduce — всё в одном месте
python
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:

Задание: генератор и отбор идей

  1. Состояние: topic: str, ideas: list[str], drafts: Annotated[list[str], operator.add], best: str.
  2. Узел brainstorm кладёт в ideas список из 3–5 идей по теме.
  3. Через Send разверни ветку expand на каждую идею; воркер пишет развёрнутый вариант в drafts (списком из одного элемента).
  4. Узел pick (reduce) выбирает самый длинный черновик в best. Свяжи expand → pick → END.
  5. Запусти и проверь, что число веток равно числу идей, а drafts собрал их все.
  6. Со звёздочкой: добавь max_concurrency=2 и, если воркер реально зовёт LLM, убедись, что ветки идут пачками, а не все разом.

Что дальше

Это последний урок раздела «Продвинутые возможности». Теперь ты владеешь памятью, сессиями, вложенностью и динамическим параллелизмом — всем, что нужно для нетривиальных агентов. Дальше — раздел про взаимодействие с человеком.