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的输出扇出给web和kb。merge有两条入边 → 引擎要把web、kb的输出扇入成一份,再喂给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 个槽位。也就是说,web 和 kb 拿到的是同一个对象引用。这是刻意为之的性能取舍(避免无谓深拷贝),但它给你留了一颗地雷:
⚠️ 扇出的共享引用陷阱
并行下游拿到的是同一个值引用。如果
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) > 1func 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被调用的次数和时机差异。再故意让web和kb返回带相同 key 的 map,看mergeValues的duplicated key报错——然后想想:如果你要的就是”多路结果按来源分区”,该怎么改 key 命名来避免冲突?
本章小结
- 并行有两端:扇出(一份输出喂多个下游)和扇入(多个上游汇一个下游),分别归约到
copyItem和mergeValues。 - 扇出
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 三姐妹如何把一张跑到一半的图精确地”接着跑”。