编程 Go-Kafka 生产消费实战:异步生产者 + 批量消费 + 重试死信,百万级消息处理

2026-07-21 08:15:14 +0800 CST views 11

Go-Kafka 生产消费:百万级消息处理实战

来源:微信公众号

Kafka 的高吞吐量不是吹的,单机每秒几十万消息很轻松。但这背后是有代价的:配置复杂、运维门槛高、消费端的 offset 管理也容易踩坑。

消费者重启后消息重复消费了几十万条——这是很多初学者踩过的坑,根因是 offset 提交策略的问题。

Kafka 的核心设计是顺序写磁盘 + 零拷贝,所以吞吐量能做到这么高。消息按 topic 分区存储,每个分区内部是有序的,但分区之间不保证顺序。理解这点很重要,否则设计出来的系统可能会有乱序问题。

生产者配置

segmentio/kafka-go 是 Go 里比较干净的 Kafka 客户端:

type Producer struct {
    writer *kafka.Writer
    topic  string
}

func NewProducer(brokers []string, topic string) *Producer {
    writer := &kafka.Writer{
        Addr:         kafka.TCP(brokers...),
        Topic:        topic,
        Balancer:     &kafka.LeastBytes{},
        BatchSize:    100,
        BatchTimeout: 10 * time.Millisecond,
        RequiredAcks: kafka.RequireAll,
        Async:        false,
    }
    return &Producer{writer: writer, topic: topic}
}

生产环境高吞吐配置:

writer := &kafka.Writer{
    Addr:         kafka.TCP(brokers...),
    Topic:        topic,
    Balancer:     &kafka.Hash{},      // 按 key 哈希分区,保证同 key 消息有序
    BatchSize:    500,                // 每批 500 条
    BatchTimeout: 50 * time.Millisecond,
    RequiredAcks: kafka.RequireOne,   // 等 leader 确认
    Async:        true,               // 异步发送
    Compression:  kafka.Snappy,       // 压缩
}

Async: true 是提升吞吐的关键——消息先丢到本地 buffer,后台批量发送。但要注意,异步发送如果 broker 挂了,可能会丢消息。如果业务不能丢消息,就用同步发送,或者异步发送 + 错误回调处理。

异步生产者封装

用 channel + worker 做更灵活的异步封装:

type AsyncProducer struct {
    writer    *kafka.Writer
    msgChan   chan kafka.Message
    errChan   chan error
    wg        sync.WaitGroup
    ctx       context.Context
    cancel    context.CancelFunc
    batchSize int
}

func NewAsyncProducer(brokers []string, topic string, batchSize int, workers int) *AsyncProducer {
    ctx, cancel := context.WithCancel(context.Background())
    writer := &kafka.Writer{
        Addr:         kafka.TCP(brokers...),
        Topic:        topic,
        Balancer:     &kafka.Hash{},
        BatchSize:    batchSize,
        BatchTimeout: 50 * time.Millisecond,
        RequiredAcks: kafka.RequireOne,
        Async:        true,
        Compression:  kafka.Snappy,
    }
    p := &AsyncProducer{
        writer:    writer,
        msgChan:   make(chan kafka.Message, batchSize*workers),
        errChan:   make(chan error, 100),
        ctx:       ctx,
        cancel:    cancel,
        batchSize: batchSize,
    }
    for i := 0; i < workers; i++ {
        p.wg.Add(1)
        go p.worker()
    }
    go p.errorHandler()
    return p
}

发送方完全不阻塞,消息堆积在 channel 里,后台 worker 定时批量刷到 Kafka。channel 满了之后会阻塞发送方,起到背压效果。

消费者实现

Kafka 消费者最核心的问题是 offset 管理

type Consumer struct {
    reader *kafka.Reader
}

func NewConsumer(brokers []string, topic, groupID string) *Consumer {
    return &Consumer{
        reader: kafka.NewReader(kafka.ReaderConfig{
            Brokers:        brokers,
            Topic:          topic,
            GroupID:        groupID,
            MinBytes:       10e3,
            MaxBytes:       10e6,
            MaxWait:        time.Second,
            CommitInterval: time.Second, // 自动提交间隔
            StartOffset:    kafka.LastOffset,
        }),
    }
}

自动提交(CommitInterval)省心但可能丢消息或重复消费:如果消费者处理了消息但还没自动提交就挂了,重启后会重新消费。手动提交更可控,但代码复杂度更高。

批量消费

如果每条消息都写一次数据库,吞吐量肯定上不去。批量消费 + 批量写入是标准做法:

type BatchConsumer struct {
    reader    *kafka.Reader
    batchSize int
    timeout   time.Duration
}

func (c *BatchConsumer) ConsumeBatch(ctx context.Context, handler func(msgs []kafka.Message) error) error {
    batch := make([]kafka.Message, 0, c.batchSize)
    timer := time.NewTimer(c.timeout)
    defer timer.Stop()

    for {
        select {
        case <-ctx.Done():
            if len(batch) > 0 {
                c.processBatch(ctx, batch, handler)
            }
            return ctx.Err()
        case <-timer.C:
            if len(batch) > 0 {
                c.processBatch(ctx, batch, handler)
                batch = batch[:0]
            }
            timer.Reset(c.timeout)
        default:
            // 读取消息...
        }
    }
}

批量消费有两个触发条件:数量达到 batchSize 或者时间超过 timeout。既能保证吞吐量,又不会因为消息太少而一直等。

并发消费

Kafka 的并发模型和分区强相关:一个分区只能被一个消费者实例消费,所以消费者实例数不能超过分区数,否则多出来的实例会空闲。

kafka-go 的 reader 在 consumer group 模式下会自动做负载均衡,多个 reader 用同一个 GroupID 就会分配到不同分区。

消息重试与死信

Kafka 本身没有死信队列的概念,需要自己实现:

type RetryConsumer struct {
    reader      *kafka.Reader
    retryWriter *kafka.Writer
    dlqWriter   *kafka.Writer
    maxRetries  int
}

func (c *RetryConsumer) Consume(ctx context.Context, handler func(msg kafka.Message) error) error {
    for {
        msg, err := c.reader.FetchMessage(ctx)
        if err != nil {
            return err
        }
        retryCount := c.getRetryCount(msg)
        if err := handler(msg); err != nil {
            if retryCount < c.maxRetries {
                c.sendToRetry(ctx, msg, retryCount+1)
            } else {
                c.sendToDLQ(ctx, msg)
            }
        }
        c.reader.CommitMessages(ctx, msg)
    }
}

重试 topic 可以配单独的 retention time,比如只保留一天。也可以给重试 topic 加一个 delay consumer,过几分钟再消费,避免立刻重试又失败。

常见踩坑

问题解决方案
分区数不够根据并发需求提前规划,后期扩容分区会影响消息顺序
key 分布不均匀不加 key 让消息轮询到各分区,或对 key 做二次哈希
offset 提交太频繁批量 commit,每 100 条或每 5 秒一次
消费者处理太慢导致 rebalance减少 MaxWait、增大 SessionTimeout,或异步化处理逻辑
消费延迟监控consumer lag 是最重要的监控指标,lag 持续增长需扩容
消息顺序与并发矛盾严格有序只能单分区单消费者;按 key 分区保证同 key 有序

总结

Kafka 用好了确实是吞吐量利器,但相比 RabbitMQ 它的运维复杂度和开发门槛都更高。分区规划、offset 管理、消费者重平衡——这三个问题搞定了,Kafka 基本上就能稳定运行了。

核心要点:

  • 生产者用异步 + 批量 + 压缩提升吞吐
  • 消费者用手动提交 + 批量处理保证可靠性
  • 分区数要提前规划,消费者实例数不超过分区数
  • 死信队列和重试机制是生产环境的标配

推荐文章

PHP 压缩包脚本功能说明
2024-11-19 03:35:29 +0800 CST
防止 macOS 生成 .DS_Store 文件
2024-11-19 07:39:27 +0800 CST
Vue中的`key`属性有什么作用?
2024-11-17 11:49:45 +0800 CST
什么是Vue实例(Vue Instance)?
2024-11-19 06:04:20 +0800 CST
程序员茄子在线接单