Day 10 / 共 60 天 · 阶段2 状态与数据流

多字段状态与并行写

前三天把单个 reducer / 单个通道拆透了。今天把它们放到"多字段 + 多节点并行"的真实场景里:当多个节点在同一步同时往同一个字段写时,到底发生什么?为什么有的字段能合并、有的却报错?我们对照 LastValueBinaryOperatorAggregateupdate,再看 _add_schema 怎么把多份 schema 的同名字段"对齐"成一个通道,最后用官方 bench/pydantic_state.py 的真实并行图收尾。

📍 你在 60 天里的位置(阶段2:状态与数据流 · 共 6 天)
阶段1 入门 D07 reducers D08 add_messages D09 累加通道 D10 并行写 D11 输入输出 D12 Pydantic 阶段3 控制流
💡 先用一个类比兜住今天 想象三位同事同时往同一块公告栏(一个字段)贴便条。如果这块栏的规矩是"只能贴一张,多了不知道听谁的"(LastValue),三张一起来它就直接罢工报错,逼你说清规矩;如果规矩是"来几张都往下摞"(累加通道),它就把三张按合并规则叠成一张。而 LangGraph 保证:不是谁先贴谁覆盖谁,而是先把这一轮所有便条收齐,再一次性按规矩合并——这就是超步(super-step)的批量语义。今天讲透这套并发写规则。
L01

痛点:并行写为什么危险

🤔 痛点很多图会扇出(fan-out):一个节点后面接多个并行节点(比如同时查数据库、查缓存、查日志),它们各自产出结果,最后扇入(fan-in)汇到下游。问题是——如果这几个并行节点都想往同一个字段写(比如都往 results 追加、或都想更新 status),它们几乎是"同时"完成的,谁的写入算数?会不会互相覆盖丢数据?顺序乱不乱?在多线程里这是经典的"竞态"噩梦。
💡 本质:LangGraph 用"超步 + 通道合并"根除竞态不让节点直接改共享状态。节点只是"提交写入请求",本超步内所有请求先被收集成一个列表,等这一步的节点全跑完,才由每个通道的 update(该字段的所有写入) 一次性合并。于是"并发写"退化成了"给通道传一个列表"——竞态被通道的合并逻辑吸收了。字段能不能接住多值,全看它配的是哪种通道。
大白话节点之间从不"抢着改同一个变量",而是各写各的小纸条投进信箱,超步结束时信箱管理员(通道)按规矩统一处理。没有共享可变状态,就没有竞态。这也是它叫 Pregel(图计算的 BSP 模型)的原因,阶段4 会深入。
L02

超步语义:update 收到的是"一个列表"

前几天你已经反复见到线索:无论 LastValue.update 还是 BinaryOperatorAggregate.update,参数都叫 values(复数、列表)。这不是巧合——它就是"本超步内该字段收到的所有写入"。

控制流:一个超步内,并行写如何被"先收集后合并" 节点A 写 x=+3 节点B 写 x=+5 节点C 写 x=+2 收集本步写入 x → values=[3,5,2] 通道 update(values) 累加通道 → 10 LastValue → 报错! 关键:节点并行跑完 → 写入攒成列表 → 一次性交给通道合并 字段能否接住多值,取决于它是哪种通道(L03/L04 对照)
图注:同一超步的并行写被批量收集成 values 列表,再整体喂给通道。合并逻辑决定结果。
"超步(super-step)"= 图执行的一轮。一轮里所有该跑的节点并行跑,跑完统一把写入合并进通道,再进入下一轮。这是阶段4 Pregel 引擎的核心,这里先建立直觉:并行写不是即时生效的,是本轮末尾批量合并的。
L03

LastValue:并行写直接报错(这是特性)

D07 见过 LastValue.update,现在从"并行写"视角重看它(channels/last_value.py:56):

# channels/last_value.py:56
    def update(self, values):
        if len(values) == 0:
            return False
        if len(values) != 1:                              # ★本步收到超过 1 个写入
            msg = create_error_message(
                message=f"At key '{self.key}': Can receive only one value per step. "
                        "Use an Annotated key to handle multiple values.",
                error_code=ErrorCode.INVALID_CONCURRENT_GRAPH_UPDATE,
            )
            raise InvalidUpdateError(msg)                 # 直接抛错,并教你怎么修
        self.value = values[-1]
        return True
len(values) != 1 → raise★核心:LastValue(没写 reducer 的字段)只接受本步恰好 1 个写入。两个并行节点同时写它,values 长度是 2,它无法决定"覆盖成谁",于是抛 InvalidUpdateError
错误码 INVALID_CONCURRENT_GRAPH_UPDATE专门的错误码,表明这是"并发更新冲突",不是随便一个异常。可观测系统能据此归类。
提示文案报错还顺手教你:"Use an Annotated key to handle multiple values."——想并行写?给这个字段加个 reducer(变成累加通道)就行。
💡 设计取舍①:为什么并行写覆盖字段是"报错"而不是"随便留一个"? 最省事的实现是"多个值就留最后一个"(悄悄丢掉其余)。但那样会静默丢数据:你两个节点都算出了结果,其中一个被无声吞掉,你还以为都保存了。LangGraph 选择响亮地报错,逼你显式表态:"这个字段到底该覆盖(那就别并行写它)还是该合并(那就给它加 reducer)"。把"歧义"变成"编译/运行期的明确错误",而不是"运行时的隐秘数据丢失"——这是它对正确性的执念,和 D08"删不存在 ID 就报错"一脉相承。
L04

累加通道:从容接住多值

换成有 reducer 的字段(累加通道),同样的并行写就能合并。重看 D09 的 updatechannels/binop.py:123):

# channels/binop.py:123(聚焦多值循环)
    def update(self, values):
        if not values:
            return False
        if self.value is MISSING:
            self.value = values[0]
            values = values[1:]
        seen_overwrite = False
        for value in values:                 # ★本步的每个并行写入,依次叠进来
            ...
            if not seen_overwrite:
                self.value = self.operator(self.value, value)
        return True
for value in values和 LastValue 最本质的区别:它不嫌 values 多,反而正好遍历它们、逐个用 operator 合并。三个并行节点写 [3,5,2],它累加成 10;三个节点各返回一段消息列表,add_messages 把它们按 ID 合并成一条完整历史。
合并顺序回顾 D09 L07 的边界:这里是按 values 顺序叠的,所以并行字段的 reducer 最好可交换(加、并集、字典合并),避免顺序影响结果。
字段写法底层通道同一步 3 个并行写
status: strLastValueInvalidUpdateError(收到 3 个值)
total: Annotated[int, operator.add]累加✓ 三值相加
logs: Annotated[list, operator.add]累加✓ 三段列表拼接
messages: Annotated[list, add_messages]累加✓ 按 ID 合并去重
💡 一句话决策"这个字段会被多个并行节点写吗?"——会,就必须给它加 reducer(否则 LastValue 报错);不会,用默认 LastValue 挺好(还能帮你抓到"意外的并发写"这种 bug)。
L05

_add_schema:多份 schema 的同名字段如何对齐成一个通道

多字段状态里还有个隐蔽问题:一个字段名可能出现在多个 schema 里(状态 schema、输入 schema、某节点的私有输入 schema……D11 细讲)。这些"同名字段"必须映射到同一个通道对象,否则并行写就乱套了。负责登记、对齐的是 _add_schemagraph/state.py:342):

# graph/state.py:342
    def _add_schema(self, schema, /, allow_managed=True):
        if schema not in self.schemas:
            _warn_invalid_state_schema(schema)
            channels, managed, type_hints = _get_channels(schema)   # D07 的读取逻辑
            ...
            self.schemas[schema] = {**channels, **managed}
            for key, channel in channels.items():
                if key in self.channels:                    # 这个字段名之前登记过?
                    if self.channels[key] != channel:       # ★类型不一致?
                        if isinstance(channel, LastValue):
                            pass                            # 新的是 LastValue → 宽容放过
                        else:
                            raise ValueError(
                                f"Channel '{key}' already exists with a different type"
                            )
                else:
                    self.channels[key] = channel            # 首次见到 → 登记
if schema not in self.schemas同一个 schema 只处理一次(去重),避免重复登记。
key in self.channels这个字段名之前已被别的 schema 登记过通道了。现在要检查两次的通道类型是否一致
self.channels[key] != channel★靠通道的 __eq__ 判断(D09 L06 讲过累加通道怎么比、lambda 为什么宽容)。不一致意味着"同一个字段在两处定义成了不同合并规则"——这是矛盾。
isinstance(channel, LastValue): pass★宽容特例:如果新出现的是个默认 LastValue(很可能是某个只声明了类型、没写 reducer 的 schema),就放过、保留已登记的那个。因为 LastValue 是"没特别要求"的默认值,不该覆盖掉别处更明确的 reducer。
else: raise否则(两个都是明确的、不同的通道)直接报错——不允许一个字段名对应两套冲突的合并逻辑。
💡 设计取舍②:为什么对 LastValue 冲突"宽容"、对其它冲突"严格"? LastValue 是"我没意见、用默认"的信号;而写了 reducer 的通道是"我明确要这样合并"的信号。当两者撞在同一字段名上,让"有明确意见的"赢、"没意见的"让路是最符合直觉的。但如果两个字段都"有明确意见"却又不同(比如一处 operator.add、另一处 add_messages),那就是真矛盾,必须报错。宽容处理"默认 vs 明确",严格处理"明确 vs 明确"——分寸拿捏得很细。
L06

真实多字段并行图:bench/pydantic_state.py

官方性能基准里有个真实的多字段、含并行分支的图(bench/pydantic_state.py:249)。它同时演示了多个 reducer 字段 + 扇出扇入:

# bench/pydantic_state.py:249(结构节选)
builder = StateGraph(State)
builder.add_edge(START, "one")
builder.add_node("one",  partial(read_write, "messages", ["trigger_events", "primary_issue_medium"]))
builder.add_edge("one", "two")
builder.add_node("two",  partial(read_write, "trigger_events", ["autoresponse", "issue"]))
builder.add_edge("two", "three")     # ★扇出:two 后面同时接 three 和 four
builder.add_edge("two", "four")
builder.add_node("three", ...)
builder.add_node("four",  ...)
builder.add_node("five",  ...)
builder.add_edge(["three", "four"], "five")   # ★扇入:three 和 four 都完成后才进 five
builder.add_edge("five", "six")

而它的 State 里,不同字段配了不同 reducer(bench/pydantic_state.py:14):

# bench/pydantic_state.py:14(字段节选)
class State(BaseModel):
    messages: Annotated[list, operator.add] = Field(default_factory=list)        # 列表累加
    trigger_events: Annotated[list, operator.add] = Field(default_factory=list)  # 列表累加
    primary_issue_medium: Annotated[str, lambda x, y: y or x] = Field(default="email")  # "有新值用新的"
    autoresponse: Annotated[dict | None, lambda _, y: y] = Field(default=None)   # 永远覆盖
    slack_participants: Annotated[dict, operator.or_] = Field(default_factory=dict)  # 字典合并
    bot_id: str | None = Field(default=None)                                     # 没 reducer → LastValue
two → three / two → four扇出:two 之后 threefour同一超步并行跑。它们分别写不同字段(three 写 relevant_rules,four 写 categorizations/responses/memory_docs),所以不冲突。
["three","four"] → five扇入:five 用列表形式声明"等 three 和 four 都完成才触发"。这就是"同步屏障"(阶段5 讲 NamedBarrierValue)。
不同字段不同 reducer★这才是"多字段状态"的精髓:一个 State 里,每个字段按自己的语义各配一种合并规则——列表累加、字典合并、覆盖、默认覆盖,混搭自如。
lambda x, y: y or x一个务实的小 reducer:"有新值 y 就用 y,否则保留旧值 x"。注意它恰好 2 个位置参数,满足 D07 _is_field_binop 的签名校验。
📝 关键结论在这张图里,threefour 并行,但它们写的是不同字段,所以哪怕都用 LastValue 也不冲突。只有当多个并行节点写同一字段时,才必须用累加通道。这是最容易被误解的点——"并行"本身不危险,"并行写同一个 LastValue 字段"才危险。
L07

边界:输出通道过滤 + 今日小结

多字段状态还牵涉"哪些字段算最终输出"。编译时 output_channels过滤掉托管值graph/state.py:1256):

# graph/state.py:1256
        output_channels = (
            "__root__"
            if len(self.schemas[self.output_schema]) == 1
            and "__root__" in self.schemas[self.output_schema]
            else [
                key
                for key, val in self.schemas[self.output_schema].items()
                if not is_managed_value(val)      # ★托管值不算对外输出
            ]
        )
这里先埋个伏笔:输出到底包含哪些字段,由 output_schema 决定;而"托管值"这类框架内部字段会被排除。D11 就专门讲输入/输出 schema 与私有中间字段——今天先知道"不是所有字段都会出现在最终输出里"。
⚠️ 最常见的三个并行写误区"我节点里 return 了,就立刻生效"——错,本超步末尾才合并(L02)。② "并行节点写同名 LastValue 字段"——直接 InvalidUpdateError,要么别并行写它、要么加 reducer(L03)。③ "给并行字段用了顺序敏感的 reducer"——结果随执行顺序漂移、不可复现(D09 L07)。记住这三条,90% 的状态诡异问题能自查。

🧠 今天你应该能回答

  • 并行写为什么在 LangGraph 里不产生竞态?(节点不改共享状态,写入先收集成列表,超步末尾由通道统一合并)
  • update(values) 的 values 是什么?(本超步该字段收到的所有写入)
  • 并行写同一个 LastValue 字段会怎样?为什么这么设计?(报 InvalidUpdateError;宁可响亮报错也不静默丢数据)
  • 怎么让一个字段支持并行写?(加 reducer,变成累加通道)
  • _add_schema 遇到同名字段通道类型冲突时,什么时候宽容、什么时候报错?(新的是 LastValue 就让路;两个都明确且不同则报错)
  • "并行"本身危险吗?(不,"并行写同一 LastValue 字段"才危险;写不同字段无碍)

✋ 10 分钟动手

# 1. 读并行冲突的两处 update 对照
sed -n '56,67p' libs/langgraph/langgraph/channels/last_value.py
sed -n '123,144p' libs/langgraph/langgraph/channels/binop.py

# 2. 读通道对齐逻辑
sed -n '342,372p' libs/langgraph/langgraph/graph/state.py

# 3. 读真实多字段并行图
sed -n '249,297p' libs/langgraph/bench/pydantic_state.py
明日预告 · Day 11:今天末尾提到 output_schema 会过滤字段。明天讲清楚输入 schema / 输出 schema / 状态 schema 三者的关系——怎么让图"只接收某些字段、只对外暴露某些字段",中间那些"私有工作字段"怎么藏起来,以及托管值为什么禁止出现在输入输出里。
← Day 09 累加通道 Day 11 · 输入/输出 schema →