第 9 章 · goroutine、channel 与并发解析日志
本章目标:掌握 goroutine 启动与调度;使用 channel 在协程间传递数据;用 sync.WaitGroup 等待一组任务完成;用 select 实现多路复用与超时;用 sync.Mutex 保护共享计数器;在 toolkit-go 中实现 并发解析 access 日志 雏形;对照 java-dev ch10 线程模型与 fastapi-web 异步任务概念。
学时建议:4~5 小时(含 2 小时跟练)
前置:完成 go-dev ch08(包管理与 go mod);toolkit-go 已 go mod init example.com/toolkit-go。
9.1 场景说明:单线程太慢,并发来帮忙
运营团队每日收到 access.log.jsonl(JSON Lines 格式),单行示例:
{"time":"2026-08-19T10:00:01Z","method":"GET","path":"/api/v1/products/go-handbook","status":200,"duration_ms":45,"request_id":"req-demo-001"}
{"time":"2026-08-19T10:00:02Z","method":"POST","path":"/api/v1/orders","status":500,"duration_ms":3200,"request_id":"req-demo-002"}
单 goroutine 逐行解析 100 万行日志,在笔记本上可能需要 数十秒。toolkit-go 的目标是把文件 分块投递 到多个 worker,经 channel 汇总慢请求(duration_ms >= 阈值),为 ch10 Worker Pool 与 ch15 毕业项目 打基础。
┌─────────────┐ jobs chan ┌──────────────┐ results chan ┌─────────────┐
│ Reader │ ───────────► │ Worker × N │ ──────────────► │ Aggregator │
│ (主 goroutine)│ │ (解析 JSON) │ │ (Mutex 计数) │
└─────────────┘ └──────────────┘ └─────────────┘
| 对比维度 | Java Thread(java-dev ch10) | Go goroutine |
|---|---|---|
| 创建成本 | 较高(MB 级栈) | 极轻(KB 级,runtime 动态扩栈) |
| 调度 | OS 线程 1:1 或池化 | M:N 由 Go runtime 调度 |
| 通信 | 共享变量 + synchronized/Lock | channel 优先(CSP 模型) |
| 数量级 | 数百~数千 | 数万~数十万 常见 |
Go 并发口诀:不要通过共享内存来通信,而要通过通信来共享内存(Do not communicate by sharing memory; share memory by communicating)。
9.2 启动 goroutine
package main
import (
"fmt"
"time"
)
func main() {
// 匿名函数 goroutine
go func() {
fmt.Println("async from goroutine")
}()
// 命名函数 goroutine
go parseOneLine(`{"path":"/health","duration_ms":10}`)
time.Sleep(100 * time.Millisecond) // 教学演示;生产环境用 WaitGroup
}
func parseOneLine(line string) {
fmt.Println("parsing:", line[:20], "...")
}
| 要点 | 说明 |
|---|---|
go fn() | 启动新 goroutine,不等待 fn 结束 |
| main 退出 | 所有 goroutine 立即被终止(未 WaitGroup 时) |
| 闭包捕获 | for i := 0; i < 3; i++ { go func() { ... }() } 需 i := i 避免经典坑 |
Go 不需要显式 Thread 类;对照 java-dev ch10 的 ExecutorService,go 关键字即创建轻量协程。
9.3 channel 基础:类型安全的管道
// 无缓冲 channel:发送方阻塞直到接收方就绪(同步握手)
ch := make(chan int)
// 有缓冲 channel:缓冲区满才阻塞发送方
buf := make(chan string, 100)
go func() {
buf <- "log-line-1"
buf <- "log-line-2"
close(buf) // 发送方关闭;关闭后不能再 send
}()
for line := range buf { // range 自动接收直到 channel 关闭
fmt.Println(line)
}
单向 channel(函数签名收窄权限):
func producer(out chan<- string) { // 只写
out <- `{"path":"/api/ping","duration_ms":5}`
close(out)
}
func consumer(in <-chan string) { // 只读
for line := range in {
_ = line
}
}
| 操作 | 无缓冲 | 有缓冲(cap=C) |
|---|---|---|
ch <- v | 阻塞至有人接收 | 缓冲未满时不阻塞 |
v := <-ch | 阻塞至有人发送 | 缓冲非空时不阻塞 |
close(ch) | 关闭后 receive 得零值+false | 同上 |
| 向 closed 发送 | panic | panic |
toolkit-go 模型:
// internal/model/entry.go
type LogEntry struct {
Time time.Time `json:"time"`
Method string `json:"method"`
Path string `json:"path"`
Status int `json:"status"`
DurationMs int `json:"duration_ms"`
RequestID string `json:"request_id"`
}
9.4 sync.WaitGroup:等待一组 goroutine 结束
import (
"sync"
)
func parseLinesConcurrent(lines []string, workers int, threshold int) []LogEntry {
jobs := make(chan string, len(lines))
results := make(chan LogEntry, len(lines))
var wg sync.WaitGroup
worker := func() {
defer wg.Done()
for line := range jobs {
entry, err := parseLogLine(line)
if err != nil {
continue // 坏行跳过,ch11 会统计 skipped
}
if entry.DurationMs >= threshold {
results <- entry
}
}
}
wg.Add(workers)
for i := 0; i < workers; i++ {
go worker()
}
for _, ln := range lines {
jobs <- ln
}
close(jobs) // 关闭 jobs → worker 的 for-range 结束
wg.Wait() // 等所有 worker 退出
close(results) // **必须在 wg.Wait 之后** close results
var slow []LogEntry
for e := range results {
slow = append(slow, e)
}
return slow
}
| 步骤 | 顺序重要性 |
|---|---|
wg.Add(n) | 在 go worker() 之前 或之内,不能少计 |
defer wg.Done() | 每个 worker 退出时 -1 |
close(jobs) | 通知 worker 没有新任务 |
wg.Wait() | 确保不会再有人写 results |
close(results) | 让主 goroutine range 结束 |
常见顺序错误:在 wg.Wait() 之前 close(results),可能仍有 worker 在写 → panic 或丢数据。
9.5 select 多路复用
select 类似 channel 版的 switch,哪个 case 先就绪执行哪个;可设 default 非阻塞。
func receiveWithTimeout(results <-chan LogEntry, d time.Duration) ([]LogEntry, error) {
var out []LogEntry
timeout := time.After(d)
for {
select {
case e, ok := <-results:
if !ok {
return out, nil // channel 已关闭
}
out = append(out, e)
case <-timeout:
return out, fmt.Errorf("collect results timeout after %v", d)
}
}
}
| select 行为 | 说明 |
|---|---|
| 多个 case 就绪 | 随机 选一个(避免饥饿) |
| 无 case 就绪 + 无 default | 阻塞 |
| 有 default | 立即执行 default(轮询模式) |
<-ctx.Done() | ch10 context 取消的标准写法 |
9.6 sync.Mutex:保护共享 map 计数
Go map 不支持并发写;并发读写也会触发 fatal error: concurrent map writes。