不能在HTTP handler中直接调用sarama.AsyncProducer.Input(),因其非阻塞写入内存队列、不等待Kafka响应,导致日志丢失、请求阻塞、OOM;必须通过logrus.Hook+ring buffer解耦采集与发送,并严格配置RequiredAcks、Retry.Max、Flush.Bytes和MaxOpenRequests。

直接用 sarama.AsyncProducer 往 Kafka 里塞日志,不加缓冲、不调参数、不处理错误,99% 的服务会在上线后第 3 天开始丢日志或 OOM。这不是配置问题,是架构级踩坑。
为什么不能在 HTTP handler 里直接调 producer.Input()
常见错误现象:panic: send on closed channel、HTTP 请求超时、Kafka client 内存持续上涨。根本原因是 sarama.AsyncProducer.Input() 只往内部无缓冲 channel 塞消息,不保证送达,也不阻塞等待——塞不进就 panic 或静默丢弃(取决于 Config.Producer.Return.Errors 是否开启)。
- Web 请求生命周期是毫秒级,而 Kafka 批量攒批默认要等 500ms(
Config.Producer.Flush.Frequency),天然不匹配 - 每个请求 new 一个 producer 开销大;复用又容易因重连、分区变更触发
InvalidTopicError或UnknownTopicOrPartitionError - handler 返回时,消息可能还在 producer 内部 buffer 里没发出去,也没人监听
Errors()channel
必须配的 4 个 sarama.Config 参数
这些不是“可选优化”,是保底可用的前提。默认配置在日志场景下几乎必出问题:乱序、重复、丢失、OOM。
-
Config.Producer.RequiredAcks = sarama.WaitForAll:否则acks=0时 broker 接收即返,网络丢包就丢日志 -
Config.Producer.Retry.Max = 3:默认为 0,分区 leader 切换时直接失败,不重试 -
Config.Producer.Flush.Bytes = 1024 * 1024(1MB):避免小包频繁发;但别设太大,否则 buffer 积压导致延迟升高 -
Config.Net.MaxOpenRequests = 1:防止并发请求过多触发 Kafka 的NOT_COORDINATOR错误
logrus Hook + ring buffer 是最稳的日志中转模式
核心是把日志写入本地无锁 ring buffer 或带限速的 chan,再由独立 goroutine 批量消费、序列化、投 Kafka。彻底解耦 handler 和 Kafka 网络路径。
立即学习“go语言免费学习笔记(深入)”;
-
logrus.Hook.Fire()方法只做select { case logChan ,不碰 Kafka client -
logChan建议设为带缓冲的 channel(如make(chan *logrus.Entry, 10000)),太小易阻塞 handler,太大吃内存 - 消费 goroutine 用
time.Ticker触发 flush,或按条数(如每 200 条)/大小(如累计 1MB)触发 batch send,别依赖 Kafka 自动 flush - 序列化统一用
json.Marshal(),不用fmt.Sprintf—— 后者无法保证字段顺序,且易引入空格或换行破坏 JSON 格式
结构体零拷贝序列化的关键点
Go 默认 json.Marshal 会分配新 []byte,高频日志下 GC 压力陡增。关键不是换序列化库,而是复用缓冲区 + 避免指针逃逸。
- 日志结构体字段全部用值类型(
int64、string),禁用*string或嵌套 struct 指针 - 用
sync.Pool管理bytes.Buffer,每次写入前buf.Reset(),写完直接buf.Bytes()传给sarama.ProducerMessage.Value - 时间戳存
int64(UnixMilli),别存time.Time—— 后者含指针,强制逃逸到堆 - 禁止将
*logrus.Entry直接传给 producer,它含未导出字段,json包会忽略,导致关键字段丢失
真正难的不是写代码,是控制好 sync.Pool 的复用粒度、chan 容量与 sarama.AsyncProducer 错误通道的消费节奏——这三个地方一松,日志就开始漏,而且漏得悄无声息。



















