基于Uber-go/zap与kafka-go的多日志实例竞态问题排查
问题根源与解决方法
先拆解下你遇到的两个问题核心原因:
1. 双异步实例时消费者收不到消息
- Kafka配置拼写错误:你的Docker Compose里有两个关键环境变量名写错了:
KAFKA_LISTENERS-INTERNAL应该是KAFKA_LISTENERS,KAFKA_ADVERTISED应该是KAFKA_ADVERTISED_LISTENERS。这直接导致Kafka无法正确配置外部监听地址,Go生产者虽然能连到localhost:9092,但消息根本没法正确路由到集群里;异步模式下kafka-go不会主动抛出错误,消息直接丢在缓冲区里,程序退出后就没了。 - zap的Sync时机完全错误:
getLogger里的defer logger.Sync()会在Logger实例创建完成后立刻执行——这时候你还没打任何日志,完全是无用操作。而main函数结束时你又没手动调用两个Logger的Sync(),异步Kafka Writer的缓冲区还没来得及发消息,程序就终止了,消费者自然收不到内容。
2. 一同步一异步时消息错乱
- 还是Kafka配置的锅:错误的监听配置导致消息发送逻辑异常,同步模式下虽然会阻塞等待,但Kafka本身的路由问题可能引发消息截断、路由错误;再加上你没有处理Kafka发送的错误,这些异常就表现成了消息串位、片段错乱。
- 另外,你给JSONEncoder用了带颜色的
CapitalColorLevelEncoder,ANSI转义码会破坏JSON结构,也会导致消费者看到的内容错乱。
一步步修复方案:
第一步:修复Kafka Docker Compose配置
把错误的环境变量名修正,这是最核心的一步:
version: '3.8' services: zookeeper: image: confluentinc/cp-zookeeper networks: - kafka-net container_name: zookeeper environment: ZOOKEEPER_CLIENT_PORT: 2181 ports: - 2181:2181 kafka: image: confluentinc/cp-kafka networks: - kafka-net container_name: kafka environment: KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 ALLOW_PLAINTEXT_LISTENER: "yes" # 修正这两个变量名 KAFKA_LISTENERS: INTERNAL://kafka:29092,EXTERNAL://localhost:9092 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,EXTERNAL://localhost:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL ports: - 9092:9092 - 29092:29092 depends_on: - zookeeper restart: on-failure networks: kafka-net: driver: bridge
第二步:修复zap Logger的Sync逻辑
- 删除
getLogger里的defer logger.Sync(),这个操作完全没用。 - 在main函数结束前手动调用两个Logger的
Sync(),确保异步消息全部发送:
func main() { loggerA := logger.Init("test-service", "localhost:9092", "topica", false, false) loggerB := logger.Init("test-service", "localhost:9092", "topicb", false, true) // 新增defer,在程序退出前刷缓冲区 defer func() { _ = loggerA.Sync() _ = loggerB.Sync() }() ctx := context.Background() ctx2 := context.WithValue(ctx, logger.UID, "abc123") loggerA.CInfo(ctx2, "topic-a log 1") loggerB.CInfo(ctx2, "topic-b log 1") loggerA.CInfo(ctx2, "topic-a log 2") loggerB.CInfo(ctx2, "topic-b log 2") loggerA.CInfo(ctx2, "topic-a log 3") loggerB.CInfo(ctx2, "topic-b log 3") }
第三步:给KafkaProducer添加错误处理和资源管理
- 修改
KafkaProducer结构体,增加上下文和关闭方法,避免静默失败:
type KafkaProducer struct { Producer producerInterface Topic string ctx context.Context cancel context.CancelFunc } func NewKafkaProducer(c *KConfig) *KafkaProducer { ctx, cancel := context.WithCancel(context.Background()) return &KafkaProducer{ Producer: kafka.NewWriter(kafka.WriterConfig{ Brokers: []string{c.Broker}, Topic: c.Topic, Balancer: &kafka.Hash{}, Async: c.Async, RequiredAcks: -1, // -1 = 等待所有副本确认 }), Topic: c.Topic, ctx: ctx, cancel: cancel, } } func (kp *KafkaProducer) Write(msg []byte) (int, error) { err := kp.Producer.WriteMessages(kp.ctx, kafka.Message{ Key: []byte(""), Value: msg, }) if err != nil { // 这里可以换成你自己的错误日志方式,至少能定位问题 fmt.Printf("[Kafka Error] 发送到topic %s失败: %v\n", kp.Topic, err) } return len(msg), err } // 添加关闭方法,释放Kafka连接 func (kp *KafkaProducer) Close() error { kp.cancel() if w, ok := kp.Producer.(*kafka.Writer); ok { return w.Close() } return nil }
- 修改
Logger结构体,保存KafkaProducer实例,方便后续关闭:
type Logger struct { *zap.Logger Ns string Kp *KafkaProducer // 新增字段 }
- 在
Init函数里赋值:
func Init(namespace, broker, topic string, debug, async bool) *Logger { var kp *KafkaProducer = nil if broker != "" && topic != "" { kp = NewKafkaProducer(&KConfig{ Broker: broker, Topic: topic, Async: async, }) } logger := getLogger(debug, kp) return &Logger{logger, namespace, kp} }
- 最后在main的defer里加上关闭操作:
defer func() { _ = loggerA.Sync() _ = loggerB.Sync() _ = loggerA.Kp.Close() _ = loggerB.Kp.Close() }()
第四步:拆分Encoder配置,避免JSON带颜色
把控制台和Kafka的Encoder分开配置,控制台保留颜色,Kafka用纯文本Level编码:
func getLogger(debug bool, kp *KafkaProducer) *zap.Logger { var cores []zapcore.Core // 控制台用带颜色的Encoder consoleEncoder := zapcore.NewConsoleEncoder(zapcore.EncoderConfig{ TimeKey: "timeStamp", LevelKey: "level", NameKey: "logger", CallerKey: "caller", FunctionKey: zapcore.OmitKey, MessageKey: "msg", StacktraceKey: "stacktrace", LineEnding: zapcore.DefaultLineEnding, EncodeLevel: zapcore.CapitalColorLevelEncoder, EncodeTime: zapcore.ISO8601TimeEncoder, EncodeDuration: zapcore.SecondsDurationEncoder, }) // Kafka用不带颜色的JSON Encoder kafkaEncoder := zapcore.NewJSONEncoder(zapcore.EncoderConfig{ TimeKey: "timeStamp", LevelKey: "level", NameKey: "logger", CallerKey: "caller", FunctionKey: zapcore.OmitKey, MessageKey: "msg", StacktraceKey: "stacktrace", LineEnding: zapcore.DefaultLineEnding, EncodeLevel: zapcore.CapitalLevelEncoder, EncodeTime: zapcore.ISO8601TimeEncoder, EncodeDuration: zapcore.SecondsDurationEncoder, }) cores = append(cores, zapcore.NewCore(consoleEncoder, zapcore.Lock(os.Stdout), getPriority(debug)), zapcore.NewCore(consoleEncoder, zapcore.Lock(os.Stderr), zap.ErrorLevel), ) if kp != nil { cores = append(cores, zapcore.NewCore(kafkaEncoder, zapcore.Lock(zapcore.AddSync(kp)), kafkaPriority)) } logger := zap.New(zapcore.NewTee(cores...)) return logger }
改完以上内容后,记得重启Kafka容器让配置生效,应该就能解决消息丢失和错乱的问题了。
内容的提问来源于stack exchange,提问作者Mark Robson
相关产品推荐
相关产品推荐

