Iris 中 Kafka 日志必须解耦:禁止在 handler 直接调用 sarama.AsyncProducer,应通过 logrus Hook + 缓冲 channel + 独立 goroutine 批量投递,避免阻塞、丢日志与 OOM。

Iris 框架本身不内置 Kafka 支持,MVC 结构里直接在 Controller 或 Handler 中调用 Kafka 生产者必然阻塞请求、丢日志、OOM——这不是配置问题,是架构错误。
sarama.AsyncProducer 不能出现在 HTTP handler 主流程里
常见错误现象:panic: send on closed channel、HTTP 请求超时、Kafka client 内存持续上涨、日志静默丢失。根本原因是 sarama.AsyncProducer.Input() 只往内部无界/小缓冲 chan 塞消息,不等 broker 确认就返回;而 Iris 的 HTTP handler 生命周期仅毫秒级,Input() 成功 ≠ 消息已发,更不等于已落盘。
- 禁止在
iris.Context处理逻辑中直接调producer.Input() - 禁止每个请求 new 一个
sarama.AsyncProducer:连接开销大,且短生命周期导致Close()来不及执行,残留 goroutine 卡住 broker - 必须复用单例
sarama.AsyncProducer,但需配好重试与错误监听,否则分区变更时会抛UnknownTopicOrPartitionError
用 logrus Hook + ring buffer 做真正解耦
Iris MVC 日志集成的关键不是“怎么连 Kafka”,而是“怎么不让 Kafka 连上 handler”。正确路径是:Iris 写日志 → logrus hook → 本地缓冲 → 独立 goroutine 批量投 Kafka。
- 定义带缓冲的 channel:
logChan := make(chan *logrus.Entry, 10000),太小易阻塞 handler,太大吃内存 - 实现
logrus.Hook,Fire()方法只做select { case logChan ,不碰任何 Kafka client - 启一个常驻 goroutine,用
time.Ticker(如每 200ms)或按条数(如每 200 条)、大小(如累计 1MB)触发 flush,调producer.Input() - 序列化必须用
json.Marshal(),禁用fmt.Sprintf:后者破坏 JSON 格式,且反射逃逸严重
sarama.AsyncProducer 必须改的 4 个配置项
默认配置在 Iris 日志场景下必崩:消息乱序、重复、OOM、连接抖动。这不是可选项,是保底前提。
-
Config.Producer.RequiredAcks = sarama.WaitForAll:否则acks=0,broker 接收即返,网络丢包就丢日志 -
Config.Producer.Retry.Max = 3:默认 0,分区 leader 切换时直接失败,不重试 -
Config.Producer.Flush.Bytes = 1024 * 1024(1MB):避免小包刷屏,但别设为 0 或过小 -
Config.ChannelBufferSize = 1024:默认 256,突发日志时容易满而 panic,必须手动放大
结构体序列化要防指针逃逸和时间漂移
Kafka 消息自带的时间戳是 broker 接收时间,日志分析依赖客户端打点时刻。不能靠 time.Now() 在发送前塞字段,因为 AsyncProducer 的缓冲和重试会让实际发送时间漂移几百毫秒甚至秒级。
- 在日志结构体中预埋
Timestamp int64字段,由 logrus hook 的Fire()入口统一赋值:entry.Time.UnixMilli() - 字段全用值类型:
string、int64,禁用*string或嵌套 struct 指针,避免强制堆分配 - 若用
sync.Pool管理bytes.Buffer,每次写入前buf.Reset(),写完用buf.Bytes()直接传给sarama.ProducerMessage.Value - 禁用
fmt.Sprintf拼接日志字符串——它内部调 reflect,GC 压力陡增;改用strconv.AppendInt等底层函数
最易被忽略的是:Kafka 不是日志系统,它不承诺顺序、不保证不丢、不负责流控。所有“高吞吐”都建立在你主动做缓冲、限速、批量、错误隔离的基础上——Iris 的优雅只是表象,背后全是自己写的 goroutine 控制逻辑。



















