
本文介绍一种基于 goroutine 和 channel 的生产级并发方案,用于同时处理多个超大日志文件:每个文件内行级并行处理 + 独立输出协程,避免竞态、阻塞和内存溢出。
本文介绍一种基于 goroutine 和 channel 的生产级并发方案,用于同时处理多个超大日志文件:每个文件内行级并行处理 + 独立输出协程,避免竞态、阻塞和内存溢出。
在 Go 中并发处理多个超大日志文件(如各含 400 万行),关键不在于“尽可能多开 goroutine”,而在于合理分层解耦、控制资源、避免共享状态竞争。原始伪代码存在多个严重问题:wg2.Wait() 在单文件处理中阻塞整个 goroutine,导致行级并发失效;defer wg1.Done() 放在 open(newfile) 后但未检查错误,可能 panic;go processRow(line) 启动的 goroutine 没有同步机制,且结果无法可靠收集;更严重的是——所有 goroutine 共享同一个 fScanner 实例,因 bufio.Scanner 非并发安全,将引发数据错乱或 panic。
✅ 正确做法是采用 “生产者-消费者”通道模型,按职责分层:
- 输入层:每个文件由独立 *os.File 和 bufio.Scanner 流式读取(无内存压力);
- 计算层:为每行启动 goroutine 处理,结果通过带缓冲 channel(如 chan string, 10)异步投递;
- 输出层:固定数量(如 10 个)写入协程从 channel 持续消费,用 bufio.Writer 批量刷盘,显著降低 I/O 开销;
- 同步层:使用 sync.WaitGroup 精确等待所有行处理完成,再关闭 channel,触发写入协程自然退出。
以下是可直接运行的核心实现(已修复原答案中的几处隐患):
func processRow(r string) string {
// 示例:简单日志清洗(实际替换为你的业务逻辑)
return strings.TrimSpace(r) + " [PROCESSED]"
}
func writeRow(outFile *os.File, ch <p>⚠️ <strong>关键注意事项</strong>:</p>
- 永远不要在 goroutine 中直接使用循环变量(如 scanner.Text() 的返回值),必须显式传参复制(见 go func(l string));
- 输入文件应在 processFile 内部打开并关闭,避免主 goroutine 提前 Close() 导致读取失败;
- channel 缓冲区不宜过大(如 100000):虽提升吞吐,但会显著增加内存占用(每行字符串约数十~数百字节 × 缓冲数);
- 写入协程数量非越多越好:磁盘 I/O 是瓶颈,通常 4–16 个足够;可通过压测调整 numWriters;
- 务必检查 scanner.Err():防止因文件损坏或权限问题静默跳过错误;
- 若处理逻辑涉及共享状态(如全局计数器),必须加锁或改用 sync/atomic。
该方案在真实场景中可稳定处理 TB 级日志:内存占用恒定(仅缓冲 channel + 单行副本),CPU 利用率接近核心数上限,I/O 效率最大化。记住 Go 并发的黄金法则:“不要通过共享内存来通信,而应通过通信来共享内存”——channel 就是你最值得信赖的协程纽带。











