Handling Parallel Image Generation in Go

Go 缓冲通道实践工作池处理并发图片生成

前言

学习 Go 时我一直想找个机会尝试并发场景,而一个最基础的 Worker Pool(Thread Pool)就是很好的入门模式:“让固定数量的 worker 在后台运行,等待分配给它们工作”。

刚好我最近做了一个小工具 goog🔗,实际使用示例则放在 goog-demo🔗

这个场景很适合练习并发,因为每张图片彼此独立,不需要等前一张完成才能生成下一张,但也不能无限制地全部同时运行,因为后台会启动 headless Chrome tab 截图,一次性跑完并不现实。

所以我最后采用的做法是:「用 goroutine 并行处理任务,利用缓冲通道的特性实现信号量的效果,限制并行计算的数量」。

真实案例分析

举例来说,我的博客有上千篇文章的 Open Graph 预览图片要生成,通常每张图片用 satori🔗(基于 JavaScript 的图片生成库)会需要 800ms 左右,总体会耗费数十分钟来生成所有图片。如果能善用 Go(编译型语言、原生并发支持、轻量运行环境)的优势,我很好奇还能挖掘出多少硬件潜能?于是我尝试打造一个 CLI 工具来生成图片:

  1. 读取 HTML template
  2. Go template🔗 填入文章标题、描述、分类等变量
  3. 通过 chromedp🔗 打开 headless Chrome
  4. 把 HTML 注入页面
  5. 截一张 1200x630 的图片
  6. 写入图片文件

单张图片可以这样生成:

Terminal window
goog \
--var "tag=教程" \
--var "title=如何在 Go 中生成 OG 图片" \
--var "description=学习如何用 chromedp 自动生成社交卡片" \
--var "site=example.com" \
--out tutorial-og.png

Template 则是一个普通的 HTML 文件,只是可以被注入 Go Template 变量:

<body>
<div class="card">
<div class="tag">{{.tag}}</div>
<div class="title">{{.title}}</div>
<div class="description">{{.description}}</div>
<div class="footer">
<div class="dot"></div>
<span>{{.site}}</span>
</div>
</div>
</body>

但真正有用的是批处理模式,例如从 images.json 一次读取多个任务:

[
{
"vars": {
"tag": "教程",
"title": "Go 入门",
"description": "学习 Go 编程语言的基础知识。",
"site": "example.com"
},
"out": "out/getting-started.png"
},
{
"vars": {
"tag": "深入解析",
"title": "Go 中的并发",
"description": "Goroutines、channels,以及并发编程模式。",
"site": "example.com"
},
"out": "out/concurrency.png"
}
]

执行时再指定 worker 数量:

Terminal window
goog --config images.json --workers 4

--workers 4 的意思不是总共只有四张图片会被处理,而是同一时间最多只允许四个图片生成任务在运行

goog 的核心实现

用 goroutine 并行处理任务,利用缓冲通道的特性实现信号量的效果,限制并行计算的数量
generator.go
func (g *Generator) Generate(ctx context.Context, jobs []ImageJob) error {
if len(jobs) == 0 {
return fmt.Errorf("没有要处理的图片任务")
}
sem := make(chan struct{}, g.workers)
var wg sync.WaitGroup
var mu sync.Mutex
var errs []error
start := time.Now()
for i, job := range jobs {
wg.Add(1)
sem <- struct{}{} // 获取名额
go func(idx int, j ImageJob) {
defer wg.Done()
defer func() { <-sem }() // 释放名额
if err := g.processJob(ctx, j); err != nil {
mu.Lock()
errs = append(errs, fmt.Errorf("任务 %d (%s): %w", idx, j.Out, err))
mu.Unlock()
log.Printf("任务 %d 失败(%s):%v", idx, j.Out, err)
} else {
log.Printf("[%d/%d] 已保存 %s", idx+1, len(jobs), j.Out)
}
}(i, job)
}
wg.Wait()
elapsed := time.Since(start)
fmt.Printf("\n%s 内生成了 %d/%d 张图片\n", elapsed.Round(time.Millisecond), len(jobs)-len(errs), len(jobs))
if len(errs) > 0 {
return fmt.Errorf("%d 个任务失败", len(errs))
}
return nil
}

工作池解决什么问题?

Go 的并行计算可以轻松通过开启一个 goroutine 来实现:

for _, job := range jobs {
go processJob(job)
}

如果一次把上千个图片生成任务全部丢进 goroutine 会是灾难,操作混合了文件 I/O、浏览器资源、内存和外部程序协作,会导致资源耗尽,所以工作池的概念在于「让并行计算执行有上限」。

信号量是什么?

信号量(Semaphore)是实践工作池的手段之一,可以想成一个有固定数量通行证的柜台,如果 workers = 4,就代表同时只有四张通行证。每个任务开始前要先拿一张,任务结束后归还。当四张都被拿走,第五个任务就会在原地等待,直到有人完成并释放通行证。在 Go 里可以用缓冲通道自然实现:

sem := make(chan struct{}, workers)
sem <- struct{}{} // acquire:拿一张通行证
<-sem // release:归还通行证

这里用 struct{}{} 是因为我们不在乎 channel 里面的值,只在乎 buffer 里已经放了几个元素。空 struct 不占额外数据意义,很适合拿来表示“一个名额”。当 channel buffer 满了,sem <- struct{}{} 就会阻塞,直到其他 goroutine 从 sem 读出一个值。

WaitGroup:等待所有工作完成

sync.WaitGroup 负责让主流程知道所有 goroutine 都完成了。

wg.Add(1)
go func() {
defer wg.Done()
// 执行工作
}()
wg.Wait()

如果没有 wg.Wait(),主程序可能在 goroutine 还没完成之前就往下执行,甚至直接结束。我把 wg.Done() 放在 defer 里,是为了确保不管任务成功、失败或中途 return,都一定会通知 WaitGroup:这个工作已经结束。

缓冲通道:限制同时执行数量

真正控制并发数的是这一行:

sem := make(chan struct{}, g.workers)

以及每次启动 goroutine 前的 acquire:

sem <- struct{}{}

还有 goroutine 结束时的 release:

defer func() { <-sem }()

这个做法有个细节:sem <- struct{}{} 放在 go func(...) 之前,所以主 goroutine 在创建新 goroutine 之前就会先尝试取得名额。

也就是说,如果 workers = 4,前四个 job 会顺利启动。到了第五个 job,主 goroutine 会卡在 sem <- struct{}{},直到前面某个 job 结束并释放名额。

这种写法不需要真的预先创建四个长期存活的 worker goroutine,而是通过信号量让「短生命周期 goroutine」最多同时存在指定数量。

Mutex:保护错误列表

多个 goroutine 可能同时失败,并且同时把错误 append 到 errs

errs = append(errs, err)

append 会修改 slice header 和底层数组,不是 thread-safe 的操作。因此这里需要 sync.Mutex

mu.Lock()
errs = append(errs, fmt.Errorf("job %d (%s): %w", idx, j.Out, err))
mu.Unlock()

这段代码的目的不是让图片生成变慢,而是保护共享数据。图片生成本身仍然并发执行,只有写入错误列表时会短暂排队。

共享无头浏览器

一开始最容易写出的版本,是每张图片都启动一个新的 Chrome。这样隔离性最好,但成本很高,goog 最后采用的是共享 browser context:

allocCtx, allocCancel := chromedp.NewExecAllocator(
context.Background(),
append(chromedp.DefaultExecAllocatorOptions[:],
chromedp.Flag("disable-gpu", true),
chromedp.Flag("no-sandbox", true),
)...,
)
browserCtx, browserCancel := chromedp.NewContext(allocCtx)
if err := chromedp.Run(browserCtx); err != nil {
allocCancel()
return nil, fmt.Errorf("启动浏览器失败:%w", err)
}

然后每个 job 打开自己的 tab context:

tabCtx, tabCancel := chromedp.NewContext(g.browserCtx)
defer tabCancel()
tabCtx, timeoutCancel := context.WithTimeout(tabCtx, 30*time.Second)
defer timeoutCancel()

这样做的好处是 Chrome 只需要启动一次,但每个截图任务仍然有自己的页面上下文。配合 worker pool,就可以避免同时打开太多 tab。

总结

这次写 goog 最大的收获是:并发不是把所有事情同时丢出去,而是要知道哪里该并行、哪里该限流。Semaphore-based worker pool 不一定是所有 worker pool 问题的答案,但很适合这种「任务列表已知、每个任务独立、外部资源昂贵、需要限制同时执行数」的 CLI 批处理场景。

无头浏览器截图反而单张图片渲染更慢

goog 生成一张图片需要 1 秒左右的时间,反而比前面提到的 satori 800ms 左右生成更慢,原因是相较于用 JSX 渲染 SVG 的方案,它背后驱动了整个浏览器进行渲染截图,但也有更大的渲染弹性。

虽然两者的技术选型与方向不同,但我还是粗略地把我的博客从 Satori 生成🔗替换成 goog 生成🔗 后,看到整体图片渲染至少有 3 倍的速度差距(18 分钟 > 6 分钟)。所以 goog 比 satori 厉害吗?不一定,但有一些有趣的特点:

  1. 可以用任何网页 Template 渲染,而不是通过 JSX 渲染 SVG
  2. 可并行计算
  3. 现成的 GitHub Action 集成与 Markdown Frontmatter 解析

延伸阅读