目录 · 第 22 / 28 章
EinoPart IV · 支撑 ADK 的编排引擎

22扇出与扇入:并行的两端

扇出靠 copyItem(值共享引用、流走 copy-on-read);扇入靠 mergeValues(map 键并集、流按到达序);Pregel「任一」vs DAG「齐活 + skip」。

compose/graph_run.go:1020schema/stream.go:837compose/values_merge.go:39internal/merge.go:62schema/stream.go:538compose/pregel.go:55compose/dag.go:128compose/field_mapping.go:31

并行有两端,而我们总只盯着一端

一说到”并行”,大多数人脑子里浮现的是”同时跑好几个节点”。但那只是并行的一端。一条并行的路径,从”分岔”到”汇合”其实横跨两个截然不同的机制:

  • 扇出(fan-out):一个节点的输出,要同时喂给多个下游。问题是——这份输出该怎么”分身”?值和流的分身方式一样吗?
  • 扇入(fan-in):多个上游的输出,要汇到同一个下游。问题是——这些输出怎么”合并”?map 怎么合、流怎么合?下游又该等几个上游到齐才开跑?

前两章我们分别讲过 Pregel/DAG 两套引擎、以及流的 copy-on-read。这一章把它们收拢到”并行”这个统一视角下:扇出扇入这两端,在 compose 里各自归约到一个核心函数——扇出是 copyItem,扇入是 mergeValues——而”下游等几个上游”则由触发模式(Pregel vs DAG)拍板。看完这一章,你再画数据流图时,就不会再把”分岔”和”汇合”当成理所当然的箭头了。

先动手感受一下,再逐一拆开:

最小实战:一次并行研究

下面这个图,把一个问题同时丢给两个检索器,再把两路结果汇总。它同时触发了扇出和扇入:

g := compose.NewGraph[string, *Report]()
_ = g.AddLambdaNode("split", compose.InvokableLambda(splitQuery))
_ = g.AddLambdaNode("web", compose.InvokableLambda(searchWeb)) // 分支 A
_ = g.AddLambdaNode("kb", compose.InvokableLambda(searchKB)) // 分支 B
_ = g.AddLambdaNode("merge", compose.InvokableLambda(mergeReports))
_ = g.AddEdge("split", "web") // ┐ 扇出:split 的输出同时喂给 web 和 kb
_ = g.AddEdge("split", "kb") // ┘
_ = g.AddEdge("web", "merge") // ┐ 扇入:web、kb 汇到 merge
_ = g.AddEdge("kb", "merge") // ┘
  • split 有两条出边 → 引擎要把 split 的输出扇出webkb
  • merge 有两条入边 → 引擎要把 webkb 的输出扇入成一份,再喂给 merge

这两件事,用户一行都没写——都藏在引擎里。下面我们把两端各自的核心函数挖出来。

扇出:copyItem——值共享引用,流走复制

扇出的全部秘密,浓缩在一个短得惊人的函数 copyItem(compose/graph_run.go:1020):

func copyItem(item any, n int) []any {
if n < 2 {
return []any{item} // 只有一个下游,原样返回
}
ret := make([]any, n)
if s, ok := item.(streamReader); ok {
ss := s.copy(n) // ← 流:裂成 n 个独立可读的 reader
for i := range ret {
ret[i] = ss[i]
}
return ret
}
for i := range ret {
ret[i] = item // ← 普通值:n 个槽位塞的是同一个 item(共享引用!)
}
return ret
}

请盯住最后那个 for:对普通值,copyItem 不做任何深拷贝——它把同一个 item 塞进 n 个槽位。也就是说,webkb 拿到的是同一个对象引用。这是刻意为之的性能取舍(避免无谓深拷贝),但它给你留了一颗地雷:

⚠️ 扇出的共享引用陷阱

并行下游拿到的是同一个值引用。如果 web 分支在内部就地修改了这个共享对象(比如往传进来的 map 里塞 key、给 slice 追加元素),kb 分支会看到被污染的数据,产生难以复现的竞态 bug。规则很简单:扇出的下游只读,不要就地改共享输入;要改,先自己拷一份。这也是为什么 Eino 的很多节点函数倾向于”读入参、产新值”的纯函数风格。

而当输出是时,共享引用根本不成立——一条 stream 只能被消费一次,两个下游各读各的必然打架。所以 copyItem 判断出 streamReader 后改走 s.copy(n)。这里的复制不是把数据抄 n 份,而是上一章讲过的 copy-on-read:n 个副本共享一条隐藏的单链表 cpStreamElement,每个副本各持一个游标,靠 peek 惰性拉取(schema/stream.go:837):

func (p *parentStreamReader[T]) peek(idx int) (t T, err error) {
elem := p.subStreamList[idx]
// sync.Once:第一个读到该 chunk 的副本负责真正 Recv 并缓存,
// 其余副本沿链表读同一份缓存,绝不重复拉取上游。
elem.once.Do(func() {
t, err = p.sr.Recv()
elem.item = streamItem[T]{chunk: t, err: err}
if err != io.EOF {
elem.next = &cpStreamElement[T]{}
p.subStreamList[idx] = elem.next
}
})
// ...沿 subStreamList[idx] 前移游标
}

于是扇出的两种情形有了统一而不同的答案:值靠共享引用(零拷贝、须只读),流靠 copy-on-read(零重复拉取、各读各的)。 一个 copyItem 把两种语义收在一处。

扇入:mergeValues——map 键并集,流按到达序

扇入这一端,核心是 mergeValues(compose/values_merge.go:39)。它先看第一个值的类型,再决定怎么合:

// caller 保证 len(vs) > 1
func mergeValues(vs []any, opts *mergeOptions) (any, error) {
v0 := reflect.ValueOf(vs[0])
t0 := v0.Type()
if fn := internal.GetMergeFunc(t0); fn != nil { // map / 注册过合并函数的类型
return fn(vs)
}
if s, ok := vs[0].(streamReader); ok { // 多条流
// ...类型校验后
ms := s.merge(ss) // 合成一条 multiStreamReader
return ms, nil
}
return nil, fmt.Errorf("(mergeValues) unsupported type: %v", t0)
}

分两种主力情形。第一种是 map,走 mergeMap(internal/merge.go:62):

func mergeMap(typ reflect.Type, vs []any) (any, error) {
merged := reflect.MakeMap(typ)
for _, v := range vs {
iter := reflect.ValueOf(v).MapRange()
for iter.Next() {
key, val := iter.Key(), iter.Value()
if merged.MapIndex(key).IsValid() {
return nil, fmt.Errorf("(values merge map) duplicated key ('%v') found", key)
}
merged.SetMapIndex(key, val) // ← 键的并集
}
}
return merged.Interface(), nil
}

这是键并集:把多个上游的 map 合成一个大 map。注意那句 duplicated key 检查——如果两个上游产出了同一个 key,合并会返回 error(而不是后者覆盖前者)。这是刻意的严格:扇入的语义是”各上游贡献不同的字段”,键冲突通常意味着图接错了线,与其静默覆盖埋雷,不如当场报错。

📝 为什么键冲突要报错而不是覆盖

“后写覆盖先写”在单机 map 里是常识,但在并行扇入里是灾难:两个上游谁先到是不确定的(见下面的到达序),若允许覆盖,合并结果就依赖调度时序,变成不可复现的随机行为。报错等于强制你在设计期就把字段划分清楚——这与第 19 章 DAG”齐活才跑”的严格气质一脉相承。

第二种是多条流,合成一个 multiStreamReader,它的 recv 决定了到达顺序(schema/stream.go:538):

func (msr *multiStreamReader[T]) recv() (T, error) {
for len(msr.nonClosed) > 0 {
// 用 reflect.Select / receiveN 从"当前有数据"的那条流里取
chosen, item, ok := receiveN(msr.nonClosed, msr.sts)
if ok {
return item.chunk, item.err // ← 谁先有数据先返回谁
}
// ...摘掉已关闭的流
}
return t, io.EOF
}

多条流的扇入不是”先 A 后 B”的拼接,而是按 chunk 的到达顺序交错:哪条流此刻有数据,就先吐哪条的 chunk。这对流式 Agent 很关键——你不想干等最慢的那条流读完才开始输出。

谁来决定”扇入节点该等几个上游”

mergeValues 回答了”到齐之后怎么合”,但还有个前置问题:扇入节点要等几个上游到齐才开跑? 这正是第 19 章 Pregel/DAG 的差异在”扇入”场景下的直接体现。

Pregel 的 pregelChannel.get(compose/pregel.go:55)只要有值就就绪:

func (ch *pregelChannel) get(isStream bool, name string, edgeHandler *edgeHandlerManager) (
any, bool, error) {
if len(ch.Values) == 0 {
return nil, false, nil
}
defer func() { ch.Values = map[string]any{} }()
values := make([]any, len(ch.Values))
names := make([]string, len(ch.Values))
i := 0
for k, v := range ch.Values {
resolvedV, err := edgeHandler.handle(k, name, v, isStream)
if err != nil {
return nil, false, err
}
values[i] = resolvedV
names[i] = k
i++
}
if len(values) == 1 {
return values[0], true, nil
}
// merge
mergeOpts := &mergeOptions{
streamMergeWithSourceEOF: ch.mergeConfig.StreamMergeWithSourceEOF,
names: names,
}
v, err := mergeValues(values, mergeOpts)
if err != nil {
return nil, false, err
}
return v, true, nil
}
func (ch *pregelChannel) get(...) (any, bool, error) {
if len(ch.Values) == 0 {
return nil, false, nil // 一个上游都没到 → 不就绪
}
// ...有值就就绪;若来了多个上游,内部会调 mergeValues 合并
}

任一上游到货,扇入节点就能跑——若同一轮里到了多个上游,才调 mergeValues 合并。这也是为什么 Pregel 能容忍环:返回边送来的那一个值,足以再次激活节点。

DAG 的 dagChannel.get(compose/dag.go:128)则相反,要等所有前驱都有交代:

for _, state := range ch.ControlPredecessors {
if state == dependencyStateWaiting {
return nil, false, nil // 还有前驱在等 → 不就绪
}
}
for _, ready := range ch.DataPredecessors {
if !ready {
return nil, false, nil // 还有数据前驱没到 → 不就绪
}
}

所有前驱齐活才跑,而没被选中的分支靠 skip 传播明确”跳过”,避免下游永远干等。回到开头那个并行研究的例子:如果你用默认的 Pregel,merge 可能在 web 先到时就先跑一次;如果你要的是”web 和 kb 都到齐再汇总”,就该用 AllPredecessor(DAG)。同一张图的拓扑,选不同触发模式,扇入行为完全不同。

更细的扇入:Workflow 的字段级映射

Graph 的扇入是”整份输出合并”;而 Workflow 把粒度降到字段级,靠 FieldMapping(compose/field_mapping.go:31):

type FieldMapping struct {
fromNodeKey string // 来自哪个上游节点
from string // 上游输出的哪个字段/键
to string // 映射到下游输入的哪个字段/键
customExtractor func(input any) (any, error)
}

有了它,你能声明”把上游 A 的 Summary 字段 + 上游 B 的 Score 字段,分别接到下游的两个入参上”。这就是为什么 Workflow 强制跑在 DAG 模式(第 20 章):字段级扇入必须先等所有相关上游到齐,才能把各字段拼成下游完整的输入结构。从 Graph 的”整份合并”到 Workflow 的”按字段接线”,扇入的表达力被进一步细化了,但底层依旧是”等齐 + 合并”这套骨架。

🔑 本章的设计钥匙

并行的两端各归约到一个核心函数:扇出 = copyItem(值共享引用、流 copy-on-read),扇入 = mergeValues(map 键并集、流按到达序)。而”扇入节点等几个上游”这个正交问题,被下沉给可替换的 channel 策略(Pregel 任一 / DAG 齐活+skip)。数据怎么复制、怎么合并,与何时触发,是三个解耦的关注点——正因为解耦,同一套 copyItem/mergeValues 才能同时服务 Pregel、DAG、Workflow 三种编排,只是”何时触发”和”合到什么粒度”不同。

💡 动手

把上面的并行研究图跑两遍:第一遍用默认 Pregel,第二遍在编译时传 compose.WithNodeTriggerMode(compose.AllPredecessor)。观察 merge 被调用的次数和时机差异。再故意让 webkb 返回带相同 key 的 map,看 mergeValuesduplicated key 报错——然后想想:如果你要的就是”多路结果按来源分区”,该怎么改 key 命名来避免冲突?

本章小结

  • 并行有两端:扇出(一份输出喂多个下游)和扇入(多个上游汇一个下游),分别归约到 copyItemmergeValues
  • 扇出 copyItem(compose/graph_run.go:1020):普通值共享引用(零拷贝,下游须只读),流走 copy-on-read(schema/stream.go:837,共享单链表 + 各自游标,零重复拉取)。
  • 扇入 mergeValues(compose/values_merge.go:39):map 走键并集(internal/merge.go:62,键冲突返回 error 而非覆盖),多条流合成按到达序交错的 multiStreamReader(schema/stream.go:538)。
  • 扇入节点等几个上游由触发模式决定:Pregel(compose/pregel.go:55)任一到货即跑、容忍环;DAG(compose/dag.go:128)齐活才跑、skip 传播。
  • Workflow 用 FieldMapping(compose/field_mapping.go:31)把扇入细化到字段级,并因此强制 DAG。
  • 设计钥匙:复制/合并/触发是三个解耦的关注点,让一套 copyItem/mergeValues 复用于 Pregel/DAG/Workflow。

至此 compose 引擎的”并行”两端都讲清了。下一章我们回到中断这条主线,看图引擎自己的 Checkpoint 与续跑:Interrupt 家族如何定位、gob 如何序列化状态、Resume 三姐妹如何把一张跑到一半的图精确地”接着跑”。

源码

正在读取完整文件…