go 不支持生成器函数,应避免用 channel 模拟;日志清洗需基于 bufio.scanner 流式读取、字节流变换(transformfunc)、缓冲写入与错误处理保障可靠性。

Go 里没有生成器函数,别被 Python 思维带偏
Go 语言原生不支持 yield 或类似 Python 的生成器函数。试图用闭包+channel 模拟“生成器”,容易掉进 goroutine 泄漏、channel 阻塞、资源未关闭的坑里。真实场景中,日志清洗不是靠“生成一个无限流”,而是靠“按批拉取 + 流式处理 + 可中断”。
用 bufio.Scanner + io.Reader 实现真正的流式读取
海量日志文件(GB 级)不能一次性加载进内存。核心是让解析逻辑与 I/O 解耦,靠 bufio.Scanner 控制每次读取行为,避免 ReadAll 或 Lines() 这类全量加载操作。
常见错误:直接对大文件调用 strings.Split(string(data), "\n") —— 内存爆掉,OOM 直接 kill。
-
bufio.Scanner默认单行上限 64KB,超长日志行会报scanner.ErrTooLong,需提前用Scanner.Buffer扩容 - 不要用
Scanner.Text()后再做正则匹配;对每行先粗筛(比如用bytes.HasPrefix判断是否含"ERROR"),再进复杂解析 - 把日志解析逻辑封装成独立函数,输入
[]byte,输出结构体指针或 error,便于单元测试和复用
转换逻辑用 func([]byte) ([]byte, error) 接口统一收口
清洗规则(如脱敏手机号、补全时间戳、过滤 DEBUG 日志)本质是字节流变换。定义统一签名能让 pipeline 组合更灵活,也方便加中间件式处理(比如统计、采样、限速)。
示例:
type TransformFunc func([]byte) ([]byte, error)
var cleanIP = func(line []byte) ([]byte, error) {
return bytes.ReplaceAll(line, []byte("192.168.1.100"), []byte("xxx.xxx.xxx.xxx")), nil
}
var dropDebug = func(line []byte) ([]byte, error) {
if bytes.Contains(line, []byte("DEBUG")) {
return nil, nil // 返回 nil 表示丢弃
}
return line, nil
}
注意:nil 返回值代表“跳过该行”,不是错误;错误才应中断流程。多个 transform 串起来时,用简单 for 循环比 channel 管道更可控、更容易调试。
写入阶段必须显式控制 buffer 和 flush 频率
日志清洗后写入新文件或发往 Kafka,如果每行都 WriteString + Flush,I/O 开销爆炸。但 buffer 太大又可能丢失最后几条(进程 crash 或 SIGKILL 时未 flush)。
- 用
bufio.NewWriterSize(w, 64*1024)设置 64KB 缓冲区,平衡吞吐与安全性 - 每写入 N 行(比如 1000 行)主动
Flush(),避免延迟过高 - 务必在 defer 中调用
w.Flush(),否则程序正常退出时最后一块 buffer 会丢失 - 若目标是 Kafka,优先用官方
segmentio/kafka-go的Writer,它自带 batch 和重试,别自己封装裸 socket
真正难的不是怎么“流式”,而是怎么在 OOM、磁盘满、网络抖动、信号中断这些边界下保证数据不丢、不错、可重入。这些细节藏在 defer、if err != nil 和 close 检查里,而不是语法糖里。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











