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 基本上就能稳定运行了。
核心要点:
- 生产者用异步 + 批量 + 压缩提升吞吐
- 消费者用手动提交 + 批量处理保证可靠性
- 分区数要提前规划,消费者实例数不超过分区数
- 死信队列和重试机制是生产环境的标配