为什么数据采集场景绕不开 fan-out/fan-in
数据采集大概是 Go 并发模式出现频率最高的场景之一。一批 URL 要抓取,一批分页要拉完,一批 Kafka 分区要消费,一个上游表要全量同步。这些任务往往彼此独立,单个任务又主要在等待网络 I/O,于是很自然地会想到用多个 goroutine 同时干。fan-out/fan-in 就是最常见的写法:任务从一个 channel 分发到多个 worker,各 worker 处理完再汇聚成一个 channel。思路不复杂,但落到工程里,channel 关闭、worker 数量、错误传播、优雅退出,几乎每个细节都会把你卡住。

如果只是把任务列表拆给几个 goroutine,那其实不算完整的 fan-out/fan-in。更准确地说,fan-out 阶段是把工作流拆成可并行执行的子任务;fan-in 阶段是等所有子任务完成后,把结果合并到原始流程。在 Go 的 channel 模型里,这个语义通常表现为:上游一个 task channel,下游一个 result channel,中间是 N 个 worker。这个模型的价值不在于并发本身,而在于它把任务的产生、执行、结果收集三个阶段解耦了。
Fan-out/fan-in 的工程边界在哪
但解耦也带来了新的问题。比如,task channel 到底由谁来关闭?如果生产者和消费者都有多个,关闭时机稍有偏差,就会 panic 或死锁。再比如,worker 数量开多少合适?开少了浪费机器,开多了把目标服务压垮,在数据采集场景里这直接决定了采集是否成功。
先说 channel 关闭。Go 有一个约定:只能在发送方关闭 channel,接收方不要关闭,否则可能 panic。在 fan-out 阶段,如果只有生产者发送,由生产者 defer close 是最简单的。如果有多个生产者并发发送,就不能由其中一个 close,必须通过 WaitGroup 或类似机制等所有发送方结束后统一 close。fan-in 阶段也一样,多个 worker 都在向 result channel 发送,关闭它就得等所有 worker 退出。很多人写到这里就直接上一个 sync.WaitGroup,然后等 wg.Wait() 完再 close(resultCh),这没问题,但要注意如果有 worker 提前 return,WaitGroup 的计数器必须正确维护,代码一改就容易漏。
第二个边界是并发度。goroutine 是非常轻量的,但也不是无限制的。每个 goroutine 至少会有几 KB 的栈空间,遇到阻塞还会有更多开销。更关键的是,如果任务队列无界,goroutine 数量可能跟随任务量线性增长。采集场景里常见的一个事故就是数据库同步任务被拆成几百万条,然后每条开一个 goroutine,最终进程内存直接爆掉。所以并发度必须被显式控制,常见的做法是固定 worker 数量或者用带缓冲的信号量。
第三个边界是错误处理。数据采集业务里,失败是常态。超时、限流、响应体不合法、远端连接重置,这些都会发生。如果你把 fail-fast 的语义套到这上面,一个数据源失败就取消整批任务,那么采集任务很可能永远无法完成。反过来,如果对错误完全视而不见,数据质量又会出大问题。这里没有银弹,需要团队根据数据的重要程度决定:是失败即停止,还是记录后继续,还是在所有任务结束后统一重试。
一个相对靠谱的代码骨架
先给出一个常用于批量采集的代码骨架。它使用 errgroup 管理 goroutine,用固定 worker 数量限制并发,并把结果收集到一个 slice 中。注意,这只是骨架,不是可以直接上生产的模板。
type Result struct {
URL string
Data []byte
}
func collect(ctx context.Context, baseURLs []string, workerNum int) ([]Result, error) {
g, ctx := errgroup.WithContext(ctx)
taskCh := make(chan string)
// producer
g.Go(func() error {
defer close(taskCh)
for _, url := range baseURLs {
select {
case taskCh <- url:
case <-ctx.Done():
return nil
}
}
return nil
})
var mu sync.Mutex
results := make([]Result, 0, len(baseURLs))
// consumers
for i := 0; i < workerNum; i++ {
g.Go(func() error {
for {
select {
case <-ctx.Done():
return nil
case url, ok := <-taskCh:
if !ok {
return nil
}
data, err := fetch(ctx, url)
if err != nil {
return err
}
mu.Lock()
results = append(results, Result{URL: url, Data: data})
mu.Unlock()
}
}
})
}
if err := g.Wait(); err != nil {
return nil, err
}
return results, nil
}
这里有两个地方值得留意。第一,生产者通过 select 发送任务,一旦 ctx 被取消就返回,避免关闭 channel 时阻塞。第二,worker 从 taskCh 接收任务时,也检查了 ctx.Done(),这样在整体取消时,worker 不会继续领任务。但要注意,如果 fetch(ctx, url) 本身不响应 ctx,worker 还是会被卡住。所以在真实的采集代码里,fetch 内部必须使用可取消的 HTTP 请求或带超时的数据库查询。
另外,这个实现把错误语义简化为 fail-fast:任何一个 worker 返回错误,errgroup 都会取消 context,整个函数返回错误。对于需要部分成功结果的采集任务,需要把 results 收集改成通过 result channel 流式发送,聚合器独立处理每条结果和错误。
不同实现方式怎么选
很多团队其实不需要自己从零写这套逻辑。Go 生态里已经有不少选择:最原生的就是 goroutine + channel,稍微成熟一点可以用 errgroup,再复杂一点可以接入完整的批量任务框架。下面用一张表来对比几种常见实现。
| 实现方案 | 并发控制 | 错误处理 | 取消支持 | 适用场景 |
|---|---|---|---|---|
| goroutine per task | 不可控 | 自行收集 | 需自己实现 | 任务量小且稳定 |
| 自建 worker pool | 可控 worker 数 | 自定义结果/错误 channel | 需额外封装 context | 长期运行、需要限流背压的采集进程 |
| errgroup + 固定 worker | 固定 goroutine 数 | 首个错误触发取消 | 内置 context 传播 | 一次批量采集,整体成功/失败可接受 |
从表里能看出一个倾向:如果任务是一次性的、可整体失败并重跑的,errgroup 是最省心。如果采集服务是常驻进程,比如实时抓取队列消息,那么自建 worker pool 更合适,因为它能让你精确控制生命周期和背压。至于 goroutine per task,我通常只建议用在任务数量很少、且目标服务不会被打死的内部工具里。
采集场景里最容易踩的几个坑
结合我带过的和见过的采集系统,这几个坑出现的频率最高:
- 结果 channel 关闭太早或太晚。太早,下游接收方还没拿完数据;太晚,发送方可能一直阻塞,因为没人接收。常见解法是用 WaitGroup 或 errgroup 保证所有发送结束后再 close。
- worker 数量被设置成跟任务数量一样大。这等于没有 fan-out,还会带来大量 goroutine 切换。正确的做法是先估算目标系统能承受的并发量,再反过来定 worker 数。
- 忽略 panic 隔离。某个数据源返回了无法解析的内容,worker 内部 panic,如果外层没有 recover,整个进程直接退出。在采集服务里,这几乎是不可接受的。
- 把 channel buffer 当成性能优化方案。buffer 只能吸收短时抖动,如果生产速度长期大于消费速度,内存会持续增长。更重要的是,采集任务应该主动对远端做限流,而不是依赖 channel 缓冲去冷却。
这些问题都不是 fancy 的技巧,但它们几乎会出现在每一个采集服务的故障复盘里。
什么时候不要用 fan-out/fan-in
没有哪个模式是万能的。fan-out/fan-in 适合任务之间没有强依赖、且天然可以乱序处理的场景。如果你的任务有先后顺序,比如分页接口必须上一页返回了 next token 才能请求下一页,那强行并发只会让代码复杂,还不能提速。即便可以乱序,下游如果要求严格按照原始顺序写入,fan-in 阶段必须带上序号做重排,这需要额外的 buffer,复杂度和内存都会上升。
还有一种情况是任务量特别稳定且耗时极短。比如从内存缓存里批量读取一批 key,单个操作只需要几十微秒,开 goroutine 的调度开销可能比任务本身还大。这种场景更适合单线程循环。
所以更准确的判断标准不是“能不能并发”,而是“并发带来的收益能否覆盖工程复杂度”。在数据采集里,网络 I/O 通常占大头,并发收益明显,也因此成为 fan-out/fan-in 最适合的土壤。
落地时可以关注这几件事
如果你正准备在采集系统里落地这个模式,我的建议是先不要急着封装框架。从最简单的固定 worker 数开始,确认以下几点:
- 任务是否可以重复执行?如果可以,错误处理走 fail-fast 或者失败重跑都可以。如果不可以,就要做更细粒度的状态记录。
- 目标服务有没有 QPS 限制?有的话,worker 数量应设置成限制允许的上限,而不是机器的核数。
- 进程退出后如何恢复?尽量把任务 ID 记下来,重启后跳过已完成的,或者容忍重复采集后做去重。
完成这些再考虑要不要引入消息队列、任务调度、分布式执行。fan-out/fan-in 只是在单机内做并发的机制,当采集任务超出单机容量时,问题就变成分布式调度了,那是另一个话题。
回到一开始说的,数据采集场景找上 Go,很大程度上是因为标准库的并发原语足够直接。fan-out/fan-in 是其中使用频率很高的一个模式,但它不是几行 channel 代码那么简单。真正的难点在于对任务边界的理解,对错误语义的选择,以及对目标系统能力的敬畏。把这几个问题想清楚,你写出来的采集程序会稳定得多。
原创文章,作者:fudengji,如若转载,请注明出处:https://fudengji.cn/article/706/