Go中用函数切片实现流水线调度器,核心是封装tasks []func(context.Context) error、concurrency int、timeout time.Duration和errCh chan error四字段,通过worker pool+context取消+fail-fast错误通道保障顺序、错误中断与资源回收,避免裸调goroutine导致的竞态、泄漏与失控。

Go 里用函数切片实现流水线调度器,核心不是堆砌 goroutine,而是控制执行顺序、错误传播和资源回收——否则并发反而会放大问题。
为什么不能直接 for range 启动 goroutine 执行 []func()
看似简单:遍历函数切片,每个丢进 goroutine。但实际立刻暴露三个硬伤:
- 无法保证执行顺序 ——
funcA和funcB可能乱序完成,流水线断裂 - 错误无法中断后续任务 —— 某个
func()panic 或返回 error,其余仍继续跑 - 没有上下文取消机制 —— 一旦启动就收不回,超时或主动中止时 goroutine 泄漏
所以必须把 []func() error 包裹进有状态的调度结构里,而非裸奔调用。
type Pipeline 必须封装哪些字段?
一个最小可用的流水线结构体,至少要带这四样:
立即学习“go语言免费学习笔记(深入)”;
-
tasks []func(context.Context) error—— 函数签名强制带context.Context,否则无法响应取消 -
concurrency int—— 控制并行度,不是“全量并发”,而是类似 worker pool 的节流阀 -
timeout time.Duration—— 全局超时,用于初始化context.WithTimeout -
errCh chan error—— 单一错误通道,首个非 nil error 就该终止全部任务(fail-fast)
注意:concurrency 不是并发数上限,而是当前活跃 worker 数;若设为 1,则退化为串行流水线,但保留了统一错误处理和上下文能力。
如何安全地并行执行并保持 fail-fast?
关键不在“怎么启 goroutine”,而在“怎么等 + 怎么停”。参考以下逻辑骨架:
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
func (p *Pipeline) Run() error {
ctx, cancel := context.WithTimeout(context.Background(), p.timeout)
defer cancel()
<pre class="brush:php;toolbar:false;">errCh := make(chan error, 1) // 缓冲 1,防 goroutine 阻塞
done := make(chan struct{})
// 启动 worker pool
for i := 0; i < p.concurrency; i++ {
go func() {
for {
select {
case <-done:
return
default:
}
// 从任务队列取一个,执行
if len(p.tasks) == 0 {
return
}
task := p.tasks[0]
p.tasks = p.tasks[1:]
if err := task(ctx); err != nil {
select {
case errCh <- err:
default:
}
cancel() // 触发所有 task ctx.Done()
return
}
}
}()
}
// 等待任一错误或全部完成
select {
case err := <-errCh:
return err
case <-time.After(p.timeout):
return context.DeadlineExceeded
}}
这里容易踩的坑:
- 别用
range p.tasks遍历后启动 —— 切片被共享,多个 goroutine 同时操作p.tasks[0]和p.tasks = p.tasks[1:]会竞态 -
errCh必须带缓冲,否则第一个 error 就卡住 goroutine,cancel 无法广播出去 - 不要在 goroutine 内直接
return err—— 外层Run()拿不到,必须走 channel 或 shared var
任务函数怎么写才真正适配这个流水线?
不是随便写个 func() { ... } 就能塞进去。必须满足:
- 签名是
func(context.Context) error,且内部所有阻塞操作(HTTP、DB、time.Sleep)都接受该 context - 主动检查
ctx.Err() != nil并提前退出,避免无谓耗时 - 不自行 recover panic —— 让 panic 被调度器捕获并转为 error,否则流水线静默失败
示例任务:
func fetchUser(ctx context.Context) error {
req, _ := http.NewRequestWithContext(ctx, "GET", "https://api.example.com/user", nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
return fmt.Errorf("fetch user failed: %w", err)
}
defer resp.Body.Close()
// ... 处理 body
return nil
}如果某个任务耗时长但又不能取消(比如 legacy Cgo 调用),它会拖慢整个 pipeline,此时应单独拆出、不放进该调度器 —— 流水线不是万能胶,它只对可中断、可组合的任务有效。
最常被忽略的一点:concurrency 值设多少,不取决于 CPU 核数,而取决于下游服务的连接池大小或 API 限流阈值。设高了不是更快,是更快地触发 429 或连接拒绝。

















