You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 08:05:10