Gin 的 log.Println 不能直接往 Kafka 写日志,因为它是同步阻塞调用,而 Kafka 生产者存在网络延迟、重试和缓冲开销,会导致 HTTP 请求线程卡住、QPS 断崖下跌;正确做法是通过内存队列(如 buffered channel)+ 独立 goroutine 异步批量发送,并配合 sync.Pool 复用日志对象、select 非阻塞写入、SyncProducer 精确控制失败处理与降级。

为什么 Gin 的 log.Println 不能直接往 Kafka 写日志
因为 log.Println 是同步阻塞调用,而 Kafka 生产者(如 sarama.SyncProducer)本身就有网络延迟、重试、缓冲等开销。一旦 Kafka 集群抖动或积压,Gin 的 HTTP 请求线程就会卡住,QPS 断崖下跌——这不是日志收集,是主动拒绝服务。
真正可行的路径只有一条:把日志写入内存队列(如 chan 或 buffered channel),由独立 goroutine 异步消费并批量发往 Kafka。
- 别在
gin.HandlerFunc里调producer.SendMessage() - 别用
sarama.AsyncProducer并忽略Errors()和Successes()通道——它会默默丢消息 - 初始化
sarama.Config时必须设Producer.Return.Successes = true,否则你根本不知道发没发成功 - 日志结构建议统一为 JSON:
{"level":"info","ts":"2024-05-20T10:22:33Z","path":"/api/user","status":200,"latency_ms":12.5}
如何用 sync.Pool + chan *LogEntry 控制日志内存占用
高频接口每秒打几百条日志,全靠 make(chan *LogEntry, 1000) 容易撑爆内存。更稳妥的做法是复用日志对象 + 限流缓冲:
json:"ts"
-
logChan建议带缓冲(如make(chan *LogEntry, 5000)),太小容易丢,太大吃内存 - 务必用
select {... default: ...}非阻塞发送,避免中间件被卡住 - 每次用完记得
logEntryPool.Put(entry),否则 GC 压力大
sarama.SyncProducer 批量发送时的关键配置
别迷信 AsyncProducer。在日志场景下,SyncProducer 更可控:你能精确知道哪条失败、何时重试、是否要降级到本地文件。但必须调对参数:
-
Config.Producer.RequiredAcks = sarama.WaitForAll:确保 ISR 全部写入,不丢日志 -
Config.Producer.Flush.Frequency = 100 * time.Millisecond:攒批不要太久,否则延迟高 -
Config.Producer.Flush.Messages = 100:每批最多 100 条,兼顾吞吐和延迟 -
Config.Net.DialTimeout = 5 * time.Second和Config.Net.ReadTimeout = 10 * time.Second:防止 Kafka 不可用时卡死 - 发送失败时,记录错误但不要 panic——日志服务挂了,主业务得继续跑
Gin 中间件如何安全注入 Kafka 日志能力
最简方案不是封装一个 “KafkaLogger” 结构体,而是直接复用 Gin 的 c.Next() 流程,在 defer 里打点:
func KafkaLogMiddleware(producer sarama.SyncProducer, topic string) gin.HandlerFunc {
return func(c *gin.Context) {
start := time.Now()
c.Next()
<pre class="brush:php;toolbar:false;"> entry := logEntryPool.Get().(*LogEntry)
entry.Timestamp = time.Now()
entry.Level = "info"
entry.Path = c.Request.URL.Path
entry.Status = c.Writer.Status()
entry.Latency = float64(time.Since(start).Milliseconds())
msg := &sarama.ProducerMessage{
Topic: topic,
Value: sarama.StringEncoder(fmt.Sprintf("%+v", entry)),
}
_, _, err := producer.SendMessage(msg)
if err != nil {
// 降级:写本地文件 or 打到 stderr,别往上抛
log.Printf("kafka log send failed: %v", err)
}
logEntryPool.Put(entry)
}}
- 别在中间件里 new producer 实例——全局复用一个
- 别用
c.MustGet("xxx")取可能不存在的 key,先c.Get("xxx")判空 - 如果用了
gin.Recovery(),确保 Kafka 日志中间件在它之前注册,否则 panic 日志收不到
异步队列 + 同步生产者 + 对象池 + 降级兜底,这套组合在百万 QPS 网关日志场景里跑得稳。真正难的不是代码,是压测时发现 Kafka broker 的 request.queue.size 溢出,或者磁盘 IO 成瓶颈——这些得看 kafka-topics.sh --describe 和 top -p $(pgrep -f kafka.Kafka)。











