Go + Kafka 客户端实战:Sarama、Segmentio 与 Consumer Group

Go 开发 Kafka 应用:Sarama 与 Segmentio kafka-go 客户端、Consumer Group、事务生产者、Go 微服务集成

Go 语言凭借原生并发模型和轻量协程调度,是处理高吞吐消息流的理想选择。Kafka 生态中有两个主流 Go 客户端:IBM 维护的 sarama 功能全面、兼容性好;Segment 出品的 kafka-go API 简洁、依赖轻量。本文深入对比两库用法,从基础 Producer 和 Consumer 到 Consumer Group、手动提交、事务生产者,再到微服务集成,帮你做出正确选择。

1. Segmentio kafka-go:Writer/Reader

kafka-go 是 Segment 开源的纯 Go 实现的 Kafka 客户端,不依赖 librdkafka C 库,零 CGO 依赖,这意味着编译产出的二进制文件更小,交叉编译也不需要额外的工具链。它的 API 设计非常直观,核心类型只有两个:kafka.Writerkafka.Reader

1.1 基础生产者 Writer

kafka.Writer 支持同步和异步发送模式。在同步模式下,WriteMessages 是阻塞调用,直到 broker 确认或超时;异步模式通过配置 Async: true 开启,写操作立即返回,由内部 goroutine 负责批量发送。对于高吞吐场景,建议开启异步并配合合理的批量参数。

package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "github.com/segmentio/kafka-go"
)

func main() {
    writer := kafka.NewWriter(kafka.WriterConfig{
        Brokers:      []string{"localhost:9092", "localhost:9093"},
        Topic:        "test-topic",
        Balancer:     &kafka.LeastBytes{},
        BatchSize:    100,
        BatchTimeout: 10 * time.Millisecond,
        RequiredAcks: -1, // 等价于 acks=all
        Async:        false,
        Completion: func(messages []kafka.Message, err error) {
            if err != nil {
                log.Printf("batch write error: %v", err)
                return
            }
            log.Printf("batch written, count=%d", len(messages))
        },
    })
    defer writer.Close()

    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()

    msgs := []kafka.Message{
        {Key: []byte("key-1"), Value: []byte(`{"name":"alice","age":30}`)},
        {Key: []byte("key-2"), Value: []byte(`{"name":"bob","age":25}`)},
    }

    if err := writer.WriteMessages(ctx, msgs...); err != nil {
        log.Fatalf("write messages failed: %v", err)
    }
    fmt.Println("messages written successfully")
}

上述代码展示了 kafka-go 核心用法。kafka.Balancer 决定消息如何分配到分区,内置实现包括 LeastBytesRoundRobinCRC32Balancer 等。设置了 Key 时默认会基于 Key 哈希,保证相同 Key 的消息总是在同一分区。RequiredAcks 取值为 0(不等待)、1(leader 确认)、-1(所有 ISR 确认),生产环境建议设 -1BatchSizeBatchTimeout 控制批量触发,Completion 回调每批完成后调用,适合做指标统计。

1.2 基础消费者 Reader

kafka.Reader 的设计类似于一个迭代器,每次 ReadMessage 返回一条消息。它封装了自动重连、偏移量自动提交、消费组管理等底层细节,让消费端的开发变得非常简洁。

package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "github.com/segmentio/kafka-go"
)

func main() {
    reader := kafka.NewReader(kafka.ReaderConfig{
        Brokers:     []string{"localhost:9092"},
        Topic:       "test-topic",
        GroupID:     "my-consumer-group",
        StartOffset: kafka.FirstOffset,
        MinBytes:    10e3,
        MaxBytes:    10e6,
        MaxWait:     500 * time.Millisecond,
    })
    defer reader.Close()

    ctx := context.Background()
    for {
        msg, err := reader.ReadMessage(ctx)
        if err != nil {
            log.Fatalf("read message failed: %v", err)
        }

        fmt.Printf("partition=%d, offset=%d, key=%s, value=%s\n",
            msg.Partition, msg.Offset, msg.Key, msg.Value)

        // kafka-go 默认不自动提交偏移量
        // 需要在确认消息处理成功后手动调用 reader.CommitMessages(ctx, msg)
    }
}

需要注意,kafka.Reader 默认不自动提交偏移量。生产环境最佳实践是手动提交:先处理业务逻辑,再提交偏移量,实现 at-least-once 语义。要求 exactly-once 时需在外部存储中记录已处理消息 ID。MinBytesMaxBytes 控制 fetch 请求大小,MaxWait 是 broker 等待累积数据的超时时长,合理设置可平衡吞吐与延迟。

kafka-go 的使用场景:初创项目、对包体积敏感的服务(如 AWS Lambda Function)、跨平台编译需求强烈(不依赖 CGO)的容器化部署。它的缺点是缺少一些高级功能,如完整的消费者组生命周期回调、事务 API,以及相对有限的监控指标暴露。

2. Sarama:SyncProducer / AsyncProducer / PartitionConsumer

sarama 是 ShopStyle 创建、目前由 IBM 维护的 Go Kafka 客户端,功能覆盖全面,几乎实现了 Kafka 协议的所有核心功能。如果你的团队需要用到事务、消费者组的完整生命周期控制、SASL/SCRAM 安全认证等高级特性,sarama 是更好的选择。

2.1 同步生产者 SyncProducer

SyncProducer 是阻塞式的生产者,每条消息发送后都会等待 broker 的确认。它适合对延迟不敏感但对可靠性要求很高的场景,比如订单状态通知、支付回调等不能丢失消息的业务。

package main

import (
    "fmt"
    "log"

    "github.com/IBM/sarama"
)

func main() {
    config := sarama.NewConfig()
    config.Producer.RequiredAcks = sarama.WaitForAll // acks=all
    config.Producer.Retry.Max = 5
    config.Producer.Return.Successes = true
    config.Producer.Return.Errors = true
    config.Producer.Partitioner = sarama.NewHashPartitioner() // 基于 Key 哈希分区

    producer, err := sarama.NewSyncProducer(
        []string{"localhost:9092"}, config)
    if err != nil {
        log.Fatalf("create producer failed: %v", err)
    }
    defer producer.Close()

    msg := &sarama.ProducerMessage{
        Topic: "test-topic",
        Key:   sarama.StringEncoder("order-12345"),
        Value: sarama.ByteEncoder([]byte(`{"status":"paid","amount":199.00}`)),
        Headers: []sarama.RecordHeader{
            {Key: []byte("trace-id"), Value: []byte("abc-123-def")},
        },
    }

    partition, offset, err := producer.SendMessage(msg)
    if err != nil {
        log.Fatalf("send failed: %v", err)
    }
    fmt.Printf("message sent to partition=%d, offset=%d\n", partition, offset)
}

saramaEncoder 接口(如 StringEncoderByteEncoderJSONEncoder)用于消息序列化,开发也可以实现自定义的 Encoder 来对接 protobuf、avro 等序列化协议。RecordHeader 用于传递 Kafka 消息的头部信息,在链路追踪和元数据传递中非常实用。

2.2 异步生产者 AsyncProducer

AsyncProducer 是高性能场景的首选。它内部使用两个 channel 分别返回成功和失败的通知,注册的 ProducerMessage 会在发送完成后被回调填充 partition 和 offset 信息。

package main

import (
    "fmt"
    "log"

    "github.com/IBM/sarama"
)

func main() {
    config := sarama.NewConfig()
    config.Producer.Return.Successes = true
    config.Producer.Return.Errors = true
    config.Producer.Flush.Frequency = 10      // 最多缓存 10 条或 100ms 后发送
    config.Producer.Flush.MaxMessages = 100
    config.Producer.RequiredAcks = sarama.WaitForAll

    producer, err := sarama.NewAsyncProducer(
        []string{"localhost:9092"}, config)
    if err != nil {
        log.Fatalf("create async producer failed: %v", err)
    }
    defer producer.AsyncClose()

    // 必须开启 goroutine 消费成功和失败 channel,否则会被阻塞
    go func() {
        for msg := range producer.Successes() {
            fmt.Printf("success: partition=%d offset=%d\n",
                msg.Partition, msg.Offset)
        }
    }()
    go func() {
        for err := range producer.Errors() {
            log.Printf("send failed: %s, err=%v", err.Msg.Topic, err.Err)
        }
    }()

    for i := 0; i < 1000; i++ {
        producer.Input() <- &sarama.ProducerMessage{
            Topic: "test-topic",
            Key:   sarama.StringEncoder(fmt.Sprintf("key-%d", i)),
            Value: sarama.ByteEncoder([]byte(fmt.Sprintf("value-%d", i))),
        }
    }

    fmt.Println("all messages queued")
}

AsyncProducer 最关键的一点是必须启动 goroutine 消费 Successes()Errors() 两个 channel。如果不消费,内部缓冲区填满后写入操作会阻塞甚至死锁。Flush.FrequencyFlush.MaxMessages 控制批量发送的时机,和 Java 客户端的 linger.msbatch.size 对应。

2.3 分区消费者 PartitionConsumer

saramaPartitionConsumer 提供了底层消费原语,让你可以直接消费指定分区的消息,绕过消费者组协议。它适合一些特殊场景,比如数据迁移、指定时间点的回溯消费、消费端完全自定义负载均衡逻辑等。

package main

import (
    "fmt"
    "log"

    "github.com/IBM/sarama"
)

func main() {
    config := sarama.NewConfig()
    config.Consumer.Return.Errors = true
    config.Consumer.Offsets.Initial = sarama.OffsetOldest // 从头消费

    client, err := sarama.NewClient([]string{"localhost:9092"}, config)
    if err != nil {
        log.Fatalf("create client failed: %v", err)
    }
    defer client.Close()

    consumer, err := sarama.NewConsumerFromClient(client)
    if err != nil {
        log.Fatalf("create consumer failed: %v", err)
    }

    partitionConsumer, err := consumer.ConsumePartition(
        "test-topic", 0, sarama.OffsetNewest)
    if err != nil {
        log.Fatalf("consume partition failed: %v", err)
    }
    defer partitionConsumer.Close()

    for msg := range partitionConsumer.Messages() {
        fmt.Printf("partition=%d, offset=%d, value=%s\n",
            msg.Partition, msg.Offset, msg.Value)
    }
}

PartitionConsumer 有两个需要特别注意的地方。一是它完全绕过消费组,不会提交偏移量到 __consumer_offsets,停机重启后需要自行管理从哪里开始消费;二是它一次只能消费一个分区,如果 Topic 有多个分区,需要为每个分区独立创建 PartitionConsumer,并由调用方负责并发调度。这种完全的透明性和控制力是一把双刃剑。

3. Sarama Consumer Group:Setup / Cleanup / ConsumeClaim

当使用 Sarama 消费多分区 Topic 并需要高可用和水平扩展时,必须使用 Consumer Group。Sarama 通过 ConsumerGroupHandler 接口暴露消费者组的生命周期回调,这三个方法构成了消费端的核心骨架。

package main

import (
    "context"
    "fmt"
    "log"
    "os"
    "os/signal"
    "sync"
    "syscall"

    "github.com/IBM/sarama"
)

type consumerGroupHandler struct{}

func (h consumerGroupHandler) Setup(sarama.ConsumerGroupSession) error {
    log.Println("consumer group setup: partition assignment completed")
    return nil
}

func (h consumerGroupHandler) Cleanup(sarama.ConsumerGroupSession) error {
    log.Println("consumer group cleanup: releasing resources")
    return nil
}

func (h consumerGroupHandler) ConsumeClaim(
    session sarama.ConsumerGroupSession,
    claim sarama.ConsumerGroupClaim) error {

    for msg := range claim.Messages() {
        log.Printf("received: topic=%s partition=%d offset=%d value=%s",
            msg.Topic, msg.Partition, msg.Offset, msg.Value)

        if err := processMessage(msg); err != nil {
            log.Printf("process failed: %v, skipping commit", err)
            continue
        }

        // 业务处理成功后手动提交偏移量
        session.MarkMessage(msg, "")
    }
    return nil
}

func processMessage(msg *sarama.ConsumerMessage) error {
    fmt.Printf("processing: %s\n", msg.Value)
    return nil
}

func main() {
    config := sarama.NewConfig()
    config.Consumer.Group.Rebalance.Strategy =
        sarama.NewBalanceStrategyRoundRobin()
    config.Consumer.Offsets.Initial = sarama.OffsetOldest
    config.Consumer.Offsets.AutoCommit.Enable = false // 关闭自动提交
    config.Version = sarama.V2_8_0_0

    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    group, err := sarama.NewConsumerGroup(
        []string{"localhost:9092"}, "my-consumer-group", config)
    if err != nil {
        log.Fatalf("create consumer group failed: %v", err)
    }

    topics := []string{"test-topic"}

    var wg sync.WaitGroup
    wg.Add(1)
    go func() {
        defer wg.Done()
        for {
            if err := group.Consume(ctx, topics, consumerGroupHandler{}); err != nil {
                log.Printf("consume error: %v", err)
            }
            if ctx.Err() != nil {
                return
            }
        }
    }()

    sig := make(chan os.Signal, 1)
    signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
    <-sig

    log.Println("shutdown signal received, closing consumer group...")
    cancel()
    wg.Wait()
    group.Close()
}

这个示例覆盖了消费者组的完整生命周期管理。Setup 在分区分配完成后调用,是初始化资源(如数据库连接池预热、指标埋点)的好时机;Cleanup 在分区释放前调用,用于清理资源;ConsumeClaim 则是实际的消息消费主循环,claim.Messages() 返回的是一个 channel,循环退出时表示该分区的消费已结束(Rebalance 期间会被重新分配)。

关键点在于 config.Consumer.Offsets.AutoCommit.Enable = false,关闭了 sarama 的自动偏移量提交机制,转而在 processMessage 成功后再调用 session.MarkMessage。这样可以确保消息被成功处理后,对应的偏移量才会被标记,兼顾可靠性和灵活性。如果 ConsumeClaim 返回 error,Consumer Group 会退出当前消费并尝试重连,外层 for 循环中的 group.Consume 会被重新调用。

3.1 Rebalance 监听:处理分区分配变更

Rebalance 是消费者组的固有行为,当消费者实例加入或退出时触发。在 Setup 方法中,可以通过 session.Claims() 获取当前实例被分配到的 Topic 和分区列表,动态调整消费并发度或其他资源分配。

func (h consumerGroupHandler) Setup(session sarama.ConsumerGroupSession) error {
    claims := session.Claims()
    for topic, partitions := range claims {
        log.Printf("assigned: topic=%s partitions=%v
", topic, partitions)
    }
    return nil
}

如果 Rebalance 过于频繁(如大量使用 Spot 实例或持续发布的微服务),可以考虑在 sarama 中启用 Group.InstanceId(静态成员),但这要求 Kafka 集群版本在 2.3 以上。另一个常见的优化是调整 Session.TimeoutHeartbeat.Interval,使消费者有足够的时间完成消息处理而不会因心跳超时被误判为死亡。

4. Sarama 事务生产者

sarama 同样支持 Kafka 的事务机制,实现 exactly-once semantics。事务机制的核心是通过 transactional.id 将多个消息发送绑定为一个原子操作,要么全部成功,要么全部回滚,同时结合 sendOffsetsToTransaction 将消费偏移量和生产输出绑定到同一个事务中。

package main

import (
    "fmt"
    "log"

    "github.com/IBM/sarama"
)

func main() {
    config := sarama.NewConfig()
    config.Producer.RequiredAcks = sarama.WaitForAll
    config.Producer.Return.Successes = true
    config.Producer.Idempotent = true // 开启幂等性
    config.Producer.Transaction.ID = "my-txn-producer-1"
    config.Net.MaxOpenRequests = 1    // 事务模式下必须设为 1

    producer, err := sarama.NewSyncProducer(
        []string{"localhost:9092"}, config)
    if err != nil {
        log.Fatalf("create transactional producer failed: %v", err)
    }
    defer producer.Close()

    // 初始化事务
    if err := producer.InitTransactions(); err != nil {
        log.Fatalf("init transactions failed: %v", err)
    }

    // 开始事务
    if err := producer.BeginTxn(); err != nil {
        log.Fatalf("begin transaction failed: %v", err)
    }

    msg1 := &sarama.ProducerMessage{
        Topic: "output-topic-1",
        Key:   sarama.StringEncoder("event-001"),
        Value: sarama.ByteEncoder([]byte(`{"type":"order_created"}`)),
    }
    msg2 := &sarama.ProducerMessage{
        Topic: "output-topic-2",
        Key:   sarama.StringEncoder("event-001"),
        Value: sarama.ByteEncoder([]byte(`{"type":"notification"}`)),
    }

    if _, _, err := producer.SendMessage(msg1); err != nil {
        producer.AbortTxn()
        log.Fatalf("send message 1 failed: %v", err)
    }
    if _, _, err := producer.SendMessage(msg2); err != nil {
        producer.AbortTxn()
        log.Fatalf("send message 2 failed: %v", err)
    }

    // 提交事务
    if err := producer.CommitTxn(); err != nil {
        log.Fatalf("commit transaction failed: %v", err)
    }
    fmt.Println("transaction committed successfully")
}

事务生产者有几个严格的约束条件。Net.MaxOpenRequests 必须设为 1,因为事务模式下最多只允许一个未完成的请求。Producer.Idempotent 必须开启。transactional.id 必须是应用范围内全局唯一的标识符,如果应用有多实例部署,通常结合主机名或 Pod 名称生成,例如 my-app-producer-${HOSTNAME}

在端到端 exactly-once 场景中,还需要 Consumer 设置 IsolationLevel: sarama.ReadCommitted,使其只读取已提交的事务消息。这样即使事务中途回滚,未提交的消息对消费者不可见,实现了从消费到生产的全链路事务保证。

config.Consumer.IsolationLevel = sarama.ReadCommitted

事务生产者虽然提供了最强的可靠性保证,但性能开销也显著高于普通生产者。事务协调器的存在增加了延迟,每个事务的 begin、commit、abort 都需要额外的 RPC。因此,事务只适合订单支付、账户转账等对一致性要求极高的业务场景,而不是日志采集或埋点上报这类高频冗余场景。

5. 与 Go 微服务集成:优雅关闭与 Context 传递

在基于 Go 的微服务架构中,Kafka Consumer 通常作为后台常驻 goroutine 运行,与 HTTP/gRPC 服务共享进程生命周期。这种模式下,最关键的工程实践是优雅关闭(Graceful Shutdown):当进程收到 SIGTERM 或 SIGINT 时,先停止接收新的 HTTP 请求,等正在处理的消息消费完成后再安全退出。

package main

import (
    "context"
    "fmt"
    "log"
    "net/http"
    "os"
    "os/signal"
    "sync"
    "syscall"
    "time"

    "github.com/IBM/sarama"
)

type KafkaWorker struct {
    consumerGroup sarama.ConsumerGroup
    handler       sarama.ConsumerGroupHandler
    topics        []string
}

func (w *KafkaWorker) Run(ctx context.Context) error {
    for {
        if err := w.consumerGroup.Consume(ctx, w.topics, w.handler); err != nil {
            log.Printf("consume error: %v", err)
        }
        if ctx.Err() != nil {
            return ctx.Err()
        }
    }
}

func (w *KafkaWorker) Close() error {
    return w.consumerGroup.Close()
}

type serverHandler struct {
    worker *KafkaWorker
}

func (s *serverHandler) health(w http.ResponseWriter, r *http.Request) {
    w.WriteHeader(http.StatusOK)
    w.Write([]byte("ok"))
}

func main() {
    config := sarama.NewConfig()
    config.Consumer.Offsets.AutoCommit.Enable = false
    config.Consumer.Group.Rebalance.Strategy = sarama.NewBalanceStrategyRange()

    group, err := sarama.NewConsumerGroup(
        []string{"localhost:9092"}, "microservice-group", config)
    if err != nil {
        log.Fatalf("create consumer group failed: %v", err)
    }

    worker := &KafkaWorker{
        consumerGroup: group,
        handler:       consumerGroupHandler{},
        topics:        []string{"order-events", "inventory-events"},
    }

    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    var wg sync.WaitGroup
    wg.Add(1)
    go func() {
        defer wg.Done()
        if err := worker.Run(ctx); err != nil {
            log.Printf("worker stopped: %v", err)
        }
    }()

    srv := &http.Server{
        Addr:    ":8080",
        Handler: http.HandlerFunc(serverHandler{worker: worker}.health),
    }

    go func() {
        if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
            log.Fatalf("http server error: %v", err)
        }
    }()

    sig := make(chan os.Signal, 1)
    signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
    <-sig

    log.Println("starting graceful shutdown...")

    cancel() // 通知 Kafka worker 停止
    wg.Wait() // 等待消费轮次完成

    shutdownCtx, shutdownCancel := context.WithTimeout(
        context.Background(), 30*time.Second)
    defer shutdownCancel()
    srv.Shutdown(shutdownCtx)

    worker.Close()
    log.Println("shutdown completed")
}

在这个架构中,http.Server 提供 HTTP 健康检查端点,Kubernetes 或负载均衡器可以通过该端点判断服务是否就绪。Kafka Consumer 作为独立 goroutine 在后台运行,context.WithCancel 作为协程间的协同信号。当收到 SIGTERM 后,cancel() 被调用,Consumer Group 会退出当前的 Consume 调用,WaitGroup 保证未处理完的消息有机会完成。最后 HTTP Server 关闭,整个进程安全退出。

5.1 Context 与 Trace 传递

在微服务体系中,消息的 Producer 和消费处理链通常分布在不同的服务中。将 trace-id 或 span 信息注入 Kafka 消息的 Headers 中,是实现分布式链路追踪的常见手段。sarama.RecordHeaderkafka-go.Message.Headers 都提供了标准的 header 支持,可以与 OpenTelemetry 或 Jaeger 集成。

// 发送端注入 trace context
headers := []sarama.RecordHeader{
    {Key: []byte("trace-id"), Value: []byte(traceID)},
    {Key: []byte("span-id"), Value: []byte(spanID)},
}
producer.Input() <- &sarama.ProducerMessage{
    Topic:   "order-events",
    Headers: headers,
    Value:   sarama.ByteEncoder(payload),
}

// 消费端提取 trace context
for _, h := range msg.Headers {
    if string(h.Key) == "trace-id" {
        ctx = context.WithValue(ctx, "trace-id", string(h.Value))
    }
}

sarama 的 Consumer Group Handler 和 kafka-go 的 Reader 都在同一个 goroutine 中顺序处理消息,天然保证了单个分区内的消息处理顺序性。如果业务允许用多个 goroutine 并发处理不同分区的消息,可以在 ConsumeClaim 中将每条消息投递到一个 goroutine 队列中,但这会失去 offset 提交和消息处理之间的原子保证,需要额外设计幂等和重试机制。

6. 订单事件消费者实战

下面是一个贴近生产环境的完整示例:一个订单服务的 Kafka 消费者,监听 order-events Topic,处理订单创建、取消和支付完成三种事件,并将处理结果写入 MySQL 数据库。

package main

import (
    "context"
    "database/sql"
    "encoding/json"
    "fmt"
    "log"
    "os"
    "os/signal"
    "strings"
    "sync"
    "syscall"
    "time"

    "github.com/IBM/sarama"
    _ "github.com/go-sql-driver/mysql"
)

type OrderEvent struct {
    OrderID   string    `json:"order_id"`
    EventType string    `json:"event_type"`
    Amount    float64   `json:"amount"`
    UserID    string    `json:"user_id"`
    Timestamp time.Time `json:"timestamp"`
}

type OrderProcessor struct {
    db *sql.DB
}

func (p *OrderProcessor) Process(msg *sarama.ConsumerMessage) error {
    var event OrderEvent
    if err := json.Unmarshal(msg.Value, &event); err != nil {
        return fmt.Errorf("unmarshal failed: %w", err)
    }

    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()

    switch strings.ToLower(event.EventType) {
    case "order_created":
        return p.handleOrderCreated(ctx, event)
    case "order_paid":
        return p.handleOrderPaid(ctx, event)
    case "order_cancelled":
        return p.handleOrderCancelled(ctx, event)
    default:
        log.Printf("unknown event type: %s", event.EventType)
        return nil
    }
}

func (p *OrderProcessor) handleOrderCreated(
    ctx context.Context, event OrderEvent) error {
    _, err := p.db.ExecContext(ctx,
        "INSERT INTO orders (order_id, user_id, amount, status, created_at) VALUES (?, ?, ?, ?, ?)",
        event.OrderID, event.UserID, event.Amount, "pending", event.Timestamp)
    if err != nil && strings.Contains(err.Error(), "Duplicate entry") {
        log.Printf("duplicate order event: %s", event.OrderID)
        return nil
    }
    return err
}

func (p *OrderProcessor) handleOrderPaid(
    ctx context.Context, event OrderEvent) error {
    _, err := p.db.ExecContext(ctx,
        "UPDATE orders SET status = ?, paid_at = ? WHERE order_id = ?",
        "paid", time.Now(), event.OrderID)
    return err
}

func (p *OrderProcessor) handleOrderCancelled(
    ctx context.Context, event OrderEvent) error {
    _, err := p.db.ExecContext(ctx,
        "UPDATE orders SET status = ? WHERE order_id = ?",
        "cancelled", event.OrderID)
    return err
}

type orderConsumerHandler struct {
    processor *OrderProcessor
}

func (h orderConsumerHandler) Setup(sarama.ConsumerGroupSession) error {
    log.Println("order consumer setup")
    return nil
}

func (h orderConsumerHandler) Cleanup(sarama.ConsumerGroupSession) error {
    log.Println("order consumer cleanup")
    return nil
}

func (h orderConsumerHandler) ConsumeClaim(
    session sarama.ConsumerGroupSession,
    claim sarama.ConsumerGroupClaim) error {

    for msg := range claim.Messages() {
        log.Printf("received order event: topic=%s partition=%d offset=%d",
            msg.Topic, msg.Partition, msg.Offset)

        if err := h.processor.Process(msg); err != nil {
            log.Printf("process failed: %v, topic=%s offset=%d",
                err, msg.Topic, msg.Offset)
            // 生产环境:发送到 DLQ(死信队列)或报警
            continue
        }

        session.MarkMessage(msg, "")
        log.Printf("committed offset=%d", msg.Offset)
    }
    return nil
}

func main() {
    dsn := os.Getenv("DB_DSN")
    if dsn == "" {
        dsn = "user:password@tcp(localhost:3306)/orderdb?parseTime=true"
    }
    db, err := sql.Open("mysql", dsn)
    if err != nil {
        log.Fatalf("open db failed: %v", err)
    }
    db.SetMaxOpenConns(20)
    db.SetMaxIdleConns(10)
    db.SetConnMaxLifetime(30 * time.Minute)
    defer db.Close()

    config := sarama.NewConfig()
    config.Consumer.Group.Rebalance.Strategy = sarama.NewBalanceStrategyRange()
    config.Consumer.Offsets.Initial = sarama.OffsetOldest
    config.Consumer.Offsets.AutoCommit.Enable = false
    config.Version = sarama.V2_8_0_0

    brokersEnv := os.Getenv("KAFKA_BROKERS")
    if brokersEnv == "" {
        brokersEnv = "localhost:9092"
    }
    brokers := strings.Split(brokersEnv, ",")

    group, err := sarama.NewConsumerGroup(brokers, "order-service-group", config)
    if err != nil {
        log.Fatalf("create consumer group failed: %v", err)
    }
    defer group.Close()

    processor := &OrderProcessor{db: db}
    handler := orderConsumerHandler{processor: processor}

    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    var wg sync.WaitGroup
    wg.Add(1)
    go func() {
        defer wg.Done()
        for {
            if err := group.Consume(ctx, []string{"order-events"}, handler); err != nil {
                log.Printf("consume error: %v", err)
            }
            if ctx.Err() != nil {
                return
            }
            time.Sleep(3 * time.Second) // Reconnect backoff
        }
    }()

    sig := make(chan os.Signal, 1)
    signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
    <-sig

    log.Println("shutting down order consumer...")
    cancel()
    wg.Wait()
    log.Println("order consumer stopped")
}

这个实战示例体现了多个生产级实践。首先是数据库连接的合理管理:SetMaxOpenConns 避免连接数过多压垮 MySQL 实例,SetConnMaxLifetime 保证连接定期轮换。其次是错误处理策略:当消息处理失败时,不是简单的 panic 或无限重试(可能阻塞后续消息),而是记录日志后继续消费,真实系统会将失败消息投递到死信队列(Dead Letter Queue)供后续人工介入处理。

最重要的是幂等性保证:handleOrderCreated 中捕获了唯一索引冲突(Duplicate entry),这确保了即使同一订单创建事件因重试被投递多次,数据库也只保留一条记录。配合手动 offset 提交(session.MarkMessage),这个消费者实现了业务语义上的 exactly-once 处理。

7. kafka-go vs sarama:选型决策

在 Go 项目的 Kafka 客户端选型上,两个库各有侧重。kafka-go 适合对包体积、交叉编译、部署流程有洁癖的现代化项目,它的 API 简洁、文档清晰,学习曲线平缓。如果你的 Kafka 使用场景只涉及基本的生产消费,且不需要事务、SASL/GSSAPI 等高级认证机制,kafka-go 通常能以更少的代码量完成任务。

sarama 则适合功能需求复杂、运维深度要求高的企业级项目。它完整支持 Kafka 协议的高级特性,包括事务、压缩、SASL/SCRAM、SSL/TLS、消费者组静态成员等,社区活跃度高,遇到边缘场景时更容易找到解决方案。sarama 在并发和性能调优上也有更多的可配置参数,对于需要根据业务负载精细调优的团队来说是更稳健的选择。

维度kafka-gosarama
CGO 依赖
依赖大小很轻(~500KB)中等(~2MB)
API 复杂度低,参数少中高,参数丰富
Consumer Group内置,支持有限完整,生命周期控制
事务支持完整
SASL/SCRAM部分完整
LZ4/snappy 压缩
社区活跃度中等高(IBM 维护)
版本兼容性较好对旧版本更完善

8. 总结

本文从 Segmentio 的 kafka-go 入手,介绍了 Writer 和 Reader 的基础用法以及批量和压缩配置;随后深入 sarama 的 SyncProducer、AsyncProducer 和 PartitionConsumer,展示了完整的生产者 API 覆盖;重点讲解了 sarama Consumer Group 的核心生命周期回调 Setup/Cleanup/ConsumeClaim,以及手动偏移量提交的最佳实践;通过事务生产者的完整代码示例,说明了 Go 中实现 exactly-once 语义的技术路径;结合 Context 和优雅关闭机制,给出了 Kafka Consumer 与 Go 微服务集成的推荐架构;最后以订单事件处理为例,展示了从消息消费、业务处理到数据库写入的完整生产级消费链路。

Kafka 客户端的选择没有银弹。在快速迭代的小型项目中,kafka-go 的轻量和简洁让你可以更快地完成工作;在需要高可靠、强一致性的大型企业系统中,sarama 的完整特性和成熟度提供了更充分的安全垫。无论选择哪个库,掌握手动偏移量提交、幂等消费、优雅关闭和错误隔离这四项实践,都是构建健壮的 Go Kafka 应用的基石。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 事件驱动架构:Event Sourcing、CQRS 与 Saga 模式
  2. Kafka 运维监控与故障恢复:JMX 指标、Lag 监控与分区重分配
  3. Kafka 详解:分布式日志系统、ISR 与一致性保证