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

Kafka生产者到消费者消息投递延迟30秒问题排查求助

Kafka消息投递延迟排查(Go + Sarama + AWS MSK)

我们有多款基于GoLang开发的微服务,通过Kafka消息总线进行消息交互。其中一款微服务向一个分区数为3、副本因子为2的Kafka Topic写入消息,采用AWS MSK作为Kafka Broker,使用Shopify的Sarama客户端连接Broker。

当向微服务施加负载,生产者生成500条单条大小约1KB的消息时,消息投递出现30秒的延迟。我们期望消息在生产后能即时投递,Kafka完全能满足该场景需求,需要协助排查延迟原因。

生产者代码

package kf

import (
    "fmt"
    "github.com/Shopify/sarama"
    "github.com/segmentio/kafka-go"
    "net"
    "strconv"
)

type Producer struct {
    flowEventProducer sarama.SyncProducer
    topic             string
}

func InitProducer(brokers []string, topic string) *Producer {
    CreateKafkaTopic(brokers[0], topic)
    p := &Producer{}
    prod, err := newFlowWriter(brokers)
    if err != nil {
        panic("failed to connect to producer")
    }
    p.flowEventProducer = prod
    p.topic = topic
    return p
}

func CreateKafkaTopic(kafkaURL, topic string) {
    conn, err := kafka.Dial("tcp", kafkaURL)
    if err != nil {
        panic(err.Error())
    }
    controller, err := conn.Controller()
    if err != nil {
        panic(err.Error())
    }

    var controllerConn *kafka.Conn
    controllerConn, err = kafka.Dial("tcp", net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port)))
    if err != nil {
        panic(err.Error())
    }

    defer controllerConn.Close()
    topicConfigs := []kafka.TopicConfig{
        {
            Topic:             topic,
            NumPartitions:     3,
            ReplicationFactor: 2,
        },
    }
    err = controllerConn.CreateTopics(topicConfigs...)
    if err != nil {
        panic(err.Error())
    }
    defer conn.Close()
}

func newFlowWriter(brokers []string) (sarama.SyncProducer, error) {
    config := sarama.NewConfig()
    version := "2.6.2"
    kafkaVer, err := sarama.ParseKafkaVersion(version)
    if err != nil {
        panic("failed to parse kafka version, producer will not run")
    }
    config.Producer.Partitioner = sarama.NewHashPartitioner
    config.Net.MaxOpenRequests = 10
    config.Producer.RequiredAcks = sarama.WaitForLocal
    config.Producer.Return.Successes = true
    config.Version = kafkaVer
    producer, err := sarama.NewSyncProducer(brokers, config)

    return producer, err
}

func (p *Producer) WriteMessage(uuid string, data []byte) error {
    msg := &sarama.ProducerMessage{
        Topic: p.topic,
        Key:   sarama.ByteEncoder(uuid),
        Value: sarama.ByteEncoder(data),
    }

    part, off, err := p.flowEventProducer.SendMessage(msg)
    if err != nil {
        return err
    } else {
        fmt.Printf("message wriiten on part:%d and offset: %d", part, off)
    }
    return nil
}

消费者代码

package kf

import (
    "context"
    "encoding/json"
    "fmt"
    "github.com/Shopify/sarama"
)

type Consumer struct {
    flowEventReader sarama.ConsumerGroup
    topic           string
    brokerUrls      []string
}

type data struct {
    Name     string `json:"name"`
    Employee string `json:"employee"`
}

func InitConsumer(brokers []string, topic string) *Consumer {
    c := &Consumer{}
    c.topic = topic
    c.brokerUrls = brokers
    var (
        err error
    )
    conf := createSaramaKafkaConf()
    c.flowEventReader, err = sarama.NewConsumerGroup(c.brokerUrls, "myconf", conf)
    if err != nil {
        panic("failed to create consumer group on kafka cluster")
    }

    return c
}

type KafkaConsumerGroupHandler struct {
    Cons *Consumer
}

func (c *Consumer) HandleMessages() {
    // Consume from kafka and process
    for {
        var err = c.flowEventReader.Consume(context.Background(), []string{c.topic}, &KafkaConsumerGroupHandler{Cons: c})
        if err != nil {
            fmt.Println("FAILED")
            continue
        }
    }

}
func (*KafkaConsumerGroupHandler) Setup(_ sarama.ConsumerGroupSession) error   { return nil }
func (*KafkaConsumerGroupHandler) Cleanup(_ sarama.ConsumerGroupSession) error { return nil }
func (l *KafkaConsumerGroupHandler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
    for msg := range claim.Messages() {
        l.Cons.logMessage(msg)
        sess.MarkMessage(msg, "")
    }
    return nil
}

func (c *Consumer) logMessage(msg *sarama.ConsumerMessage) {
    d := &data{}
    err := json.Unmarshal(msg.Value, d)
    if err != nil {
        fmt.Println(err)
    }
    fmt.Printf("messages: key: %s and val:%+v", string(msg.Key), d)
}

func createSaramaKafkaConf() *sarama.Config {
    conf := sarama.NewConfig()
    version := "2.6.2"
    kafkaVer, err := sarama.ParseKafkaVersion(version)
    if err != nil {
        panic("failed to parse kafka version, executor will not run")
    }
    conf.Version = kafkaVer
    conf.Consumer.Offsets.Initial = sarama.OffsetOldest
    conf.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{sarama.BalanceStrategyRoundRobin}

    return conf
}

延迟排查方向

生产者端问题

  • 同步发送无批量配置:当前使用SyncProducer的SendMessage单条发送,没有显式配置批量参数。Sarama默认的Producer.Flush.Frequency(默认500ms)可能导致消息攒批等待,建议显式设置Producer.Flush.Frequency为10ms或更小,强制立即刷新批量。同时检查Producer.Flush.Bytes是否设置过大,导致攒够字节数才发送。
  • 分区哈希倾斜:使用NewHashPartitioner,如果所有消息的uuid哈希后集中到某1-2个分区,会导致单个broker节点负载过高,消息写入排队。可以通过Kafka监控查看各分区的消息堆积、写入延迟指标。
  • 网络与连接问题:生产者初始化时仅通过单个broker创建topic,后续SyncProducer使用broker列表,但如果存在TCP连接握手延迟、DNS解析慢或网络抖动,会增加单条消息的发送耗时。建议在WriteMessage中添加发送耗时日志,统计每条消息从调用到返回的时间。

AWS MSK集群问题

  • broker资源瓶颈:检查AWS MSK监控的CPU使用率、磁盘IOPS、网络吞吐量指标,如果某节点CPU超过70%或磁盘IO饱和,会直接导致消息写入延迟。
  • 分区leader分布不均:如果3个分区的leader都集中在同一个broker节点,该节点会成为性能瓶颈。可以通过kafka-topics.sh --describe --topic <topic-name> --bootstrap-server <msk-broker>查看分区的leader分布。
  • 跨AZ同步延迟:如果MSK集群跨可用区部署,副本同步的跨AZ网络延迟会影响消息可见性(虽然RequiredAcks=WaitForLocal只等待leader写入,但消费者若从follower拉取会有延迟)。确认topic副本的AZ分布,以及AZ间的网络延迟是否正常。

消费者端问题

  • 消费阻塞:当前消费者的logMessage中使用同步fmt.Printf输出,高负载下会阻塞消费线程,导致消息堆积,给人“投递延迟”的错觉。建议替换为异步日志组件,或去掉不必要的控制台输出。
  • 频繁重平衡:消费者HandleMessages循环中,Consume出错会立即重试,可能导致频繁重连触发消费者组重平衡,中断消费流程。检查日志中是否有重平衡相关的警告,添加重平衡事件的日志记录。
  • 拉取配置不合理:Sarama默认的Consumer.Fetch.Min(默认1)通常没问题,但如果设置过大,会导致消费者等待足够多的消息才拉取。确认Consumer.Fetch.Max是否限制了拉取大小,导致单批次拉取消息过少,增加拉取次数。

内容的提问来源于stack exchange,提问作者Ashutosh Pandey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 15:20:33