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
相关产品推荐
相关产品推荐

