目录

Kafka 全面指南:从基础到实战

目录

1. Kafka 简介

Kafka 是由 LinkedIn 开发并开源的分布式流处理平台,现归属 Apache 基金会。它的设计目标是超高吞吐、持久化存储和水平扩展,既可作为消息队列,也可用于流式计算。

1.1 Kafka 的定位

传统消息队列(RabbitMQ):   消息投递、异步解耦、复杂路由
Kafka:                      高吞吐流式数据管道 + 消息持久化 + 流处理平台

Kafka 把消息称为 Record(记录),把队列称为 Topic(主题),消费后消息不删除,通过**偏移量(Offset)**追踪消费位置,消息可以被多个消费组独立重复消费。

1.2 与其他 MQ 的对比

特性 Kafka RabbitMQ RocketMQ
吞吐量 百万级/秒 万级/秒 十万级/秒
消息延迟 毫秒级(批量) 微秒级 毫秒级
消息顺序 分区内有序 队列内有序 分区内有序
消息堆积 极强(磁盘存储) 一般
消费模式 Pull(拉) Push(推) Push+Pull
消息回溯 支持(重置 Offset) 不支持 支持
延迟消息 不支持 插件支持 原生支持
事务消息 支持(0.11+) 不支持 支持
主要用途 日志收集、流处理、大数据 业务解耦、可靠投递 金融、电商

1.3 Kafka 核心应用场景

  • 日志收集:收集各服务日志,统一输出到 ES / HDFS
  • 消息系统:用户行为事件、订单流水等高并发消息
  • 流处理:结合 Flink / Kafka Streams 实时计算
  • 变更数据捕获(CDC):MySQL Binlog → Kafka → 数仓同步
  • 事件溯源:以 Kafka 作为事件日志,支持消息回放

2. 核心架构

2.1 整体架构图

Producers                  Kafka Cluster                    Consumers
                    ┌──────────────────────────┐
  [App-1] ────────► │  Broker-1                │ ──────► [Consumer Group A]
  [App-2] ────────► │  Broker-2    ZooKeeper   │ ──────► [Consumer Group B]
  [App-3] ────────► │  Broker-3  (KRaft alt.)  │ ──────► [Consumer Group C]
                    └──────────────────────────┘
                              │
                    Topic: order-events
                    ├── Partition 0  [Leader: Broker-1, Follower: Broker-2]
                    ├── Partition 1  [Leader: Broker-2, Follower: Broker-3]
                    └── Partition 2  [Leader: Broker-3, Follower: Broker-1]

2.2 核心术语

术语 说明
Broker Kafka 服务节点,一个 Kafka 集群由多个 Broker 组成
Topic 消息分类的逻辑概念,生产者向 Topic 写,消费者从 Topic 读
Partition(分区) Topic 的物理分片,每个 Partition 是一个有序的日志文件
Offset(偏移量) 消息在 Partition 中的唯一递增编号,消费者用它追踪消费进度
Producer(生产者) 向 Topic 发送消息
Consumer(消费者) 从 Topic 拉取消息
Consumer Group(消费组) 多个 Consumer 组成一组,共同消费同一 Topic,每条消息只被组内一个 Consumer 消费
Replica(副本) Partition 的数据副本,分 Leader 和 Follower
Leader Replica 负责所有读写请求的副本
Follower Replica 同步 Leader 数据,Leader 故障时参与选举
ISR(In-Sync Replicas) 与 Leader 保持同步的副本集合,只有 ISR 中的副本才能参与 Leader 选举
Controller 集群中的一个特殊 Broker,负责 Partition Leader 选举和集群元数据管理
ZooKeeper / KRaft 存储集群元数据;Kafka 2.8+ 引入 KRaft 模式逐步取代 ZooKeeper

2.3 分区(Partition)详解

分区是 Kafka 并发能力的核心。每个分区是一个只追加写入的日志文件(append-only log)

Partition 0:
┌─────┬─────┬─────┬─────┬─────┬─────┐
│  0  │  1  │  2  │  3  │  4  │  5  │  ← Offset
│ msg │ msg │ msg │ msg │ msg │ msg │
└─────┴─────┴─────┴─────┴─────┴─────┘
                              ↑
                        LEO (Log End Offset)
  • 写入:只追加到末尾,顺序写,性能极高
  • 读取:消费者记录 offset,下次从该 offset 读取
  • 删除:不立即删除,按保留策略(时间/大小)清理

分区数量的影响

分区越多 优点 缺点
并行度 更多消费者并行消费 Controller 选举开销大
吞吐量 更高写入并发 文件句柄消耗增多
延迟 Rebalance 时间增长

经验值:分区数 = max(目标吞吐量 / 单分区吞吐量, 消费者数量),通常设为消费者数量的整数倍。

2.4 副本机制(Replication)

Topic: orders, Partition: 0, Replication Factor: 3

Broker-1 [Leader]    Broker-2 [Follower]   Broker-3 [Follower]
┌──────────────┐     ┌──────────────┐      ┌──────────────┐
│ offset 0~100 │ ──► │ offset 0~100 │  ──► │ offset 0~99  │
└──────────────┘     └──────────────┘      └──────────────┘
      ↑                                           ↑
  ISR 成员                                   落后,可能被踢出 ISR

LEO(Log End Offset):每个副本自身最新的 offset+1 HW(High Watermark):所有 ISR 副本都已复制的最小 LEO,消费者只能读取 HW 以下的消息

Broker-1 LEO=101
Broker-2 LEO=101    →  HW = min(101, 101, 99) = 99(消费者最多读到 offset 98)
Broker-3 LEO=99

3. 生产者详解

3.1 消息发送流程

Producer
   │
   ▼
[ProducerRecord]  → 序列化 → 分区器(Partitioner) → RecordAccumulator(批次缓冲)
                                                          │
                                                     Sender 线程
                                                          │
                                                    Kafka Broker
  1. ProducerRecord 封装消息(topic, partition, key, value, headers, timestamp)
  2. 序列化:将 key/value 序列化为字节数组
  3. 分区器:决定消息发到哪个 Partition
  4. RecordAccumulator:按 Partition 缓冲成批次(Batch),减少网络请求
  5. Sender 线程:后台异步将批次发送到对应 Broker

3.2 分区策略

// 分区选择逻辑(伪代码)
func selectPartition(record ProducerRecord) int {
    if record.Partition >= 0 {
        return record.Partition      // 1. 明确指定分区
    }
    if record.Key != nil {
        return hash(record.Key) % numPartitions  // 2. 按 Key 哈希(同Key→同分区)
    }
    return stickyPartition()         // 3. 粘性分区(默认,批次填满再切换)
}

策略对比

策略 场景 特点
指定 Partition 精确控制分布 灵活,需了解分区数
Key 哈希 同 Key 保证顺序(如同一用户) 热点 Key 会导致分区不均
粘性分区(无 Key) 通用场景 Kafka 2.4+ 默认,吞吐更高
轮询(旧版默认) 均匀分布 产生过多小批次

3.3 消息可靠性:acks 参数

acks=0:不等待任何确认,最高吞吐,可能丢消息
acks=1:等待 Leader 写入成功(默认),Leader 崩溃可能丢消息
acks=-1/all:等待所有 ISR 副本写入,最安全,延迟最高
// Go 配置示例(sarama)
config := sarama.NewConfig()
config.Producer.RequiredAcks = sarama.WaitForAll  // acks=-1
config.Producer.Retry.Max = 5
config.Producer.Return.Successes = true

acks=-1 还需配合

  • min.insync.replicas=2(至少 2 个 ISR 副本确认)
  • unclean.leader.election.enable=false(禁止非 ISR 副本成为 Leader)

3.4 幂等生产者(Idempotent Producer)

Kafka 0.11+ 支持,防止网络重试导致的重复消息:

config.Producer.Idempotent = true               // 开启幂等
config.Producer.RequiredAcks = sarama.WaitForAll // 幂等必须配合 acks=-1
config.Net.MaxOpenRequests = 1                   // 幂等必须配合 max.in.flight=1(单分区)

原理:Broker 为每个 Producer 分配 PID,每条消息携带 (PID, SequenceNumber),Broker 检测重复并去重。仅保证单 Partition 内幂等

3.5 事务生产者

跨多个 Partition/Topic 的原子写入:

config.Producer.Transaction.ID = "my-transactional-id" // 开启事务

producer, _ := sarama.NewSyncProducer(brokers, config)

producer.BeginTxn()

producer.SendMessages([]*sarama.ProducerMessage{
    {Topic: "orders", Value: sarama.StringEncoder(`{"id":1}`)},
    {Topic: "inventory", Value: sarama.StringEncoder(`{"item":1,"count":-1}`)},
})

if err != nil {
    producer.AbortTxn()
} else {
    producer.CommitTxn()
}

4. 消费者详解

4.1 消费者组(Consumer Group)

Topic: orders(4个分区)

Consumer Group A(3个消费者):
  Consumer-1 → Partition 0, Partition 1
  Consumer-2 → Partition 2
  Consumer-3 → Partition 3

Consumer Group B(1个消费者):
  Consumer-1 → Partition 0, 1, 2, 3(全部)

核心规则

  • 同一消费组内,每个 Partition 只分配给一个 Consumer
  • 不同消费组之间相互独立,可以消费同一 Topic 的全量数据
  • Consumer 数量 > Partition 数量时,多余的 Consumer 空闲

4.2 Offset 管理

Kafka 0.9+ 将 Offset 存储在内部 Topic __consumer_offsets,不再依赖 ZooKeeper。

消费位置追踪:
┌─────┬─────┬─────┬─────┬─────┬─────┐
│  0  │  1  │  2  │  3  │  4  │  5  │  Partition
└─────┴─────┴─────┴─────┴─────┴─────┘
                    ↑           ↑
              committed         LEO
              offset=3        offset=5

下次拉取从 offset=3 开始

提交方式

方式 说明 风险
自动提交(enable.auto.commit=true 定时自动提交,默认每5s 可能丢消息或重复消费
同步手动提交 CommitOffsets() 阻塞等待 降低吞吐量
异步手动提交 回调形式,不阻塞 提交失败需重试逻辑

4.3 消费语义

语义 实现方式 说明
At Most Once 先提交 Offset,再处理 可能丢消息
At Least Once 先处理,再提交 Offset(默认) 可能重复消费,需幂等
Exactly Once 事务 + 幂等消费 最复杂,性能最低

生产推荐:At Least Once + 消费者幂等(Redis/DB 去重)

4.4 Rebalance(重平衡)

当消费组内 Consumer 数量或 Topic 分区数变化时,触发 Rebalance,重新分配分区。

触发条件

  • Consumer 加入或离开消费组
  • Consumer 心跳超时(session.timeout.ms
  • Topic 分区数变化
  • 订阅的 Topic 数量变化

Rebalance 过程(JoinGroup + SyncGroup)

所有 Consumer → 发送 JoinGroup 请求到 Group Coordinator
                      │
              选出 Group Leader(第一个加入的 Consumer)
                      │
             Leader 执行分区分配策略
                      │
              SyncGroup:Leader 上报分配方案
                      │
              所有 Consumer 收到各自的分区分配

Rebalance 期间所有 Consumer 停止消费,这是 Kafka 延迟的重要来源之一。

减少 Rebalance 的配置

// 增大心跳超时,避免因 GC/网络抖动误判 Consumer 下线
config.Consumer.Group.Session.Timeout = 30 * time.Second  // 默认10s
config.Consumer.Group.Heartbeat.Interval = 3 * time.Second // 建议 = Timeout/3

// 增大 poll 最大间隔,避免消费慢被踢出
config.Consumer.MaxProcessingTime = 5 * time.Minute // 对应 max.poll.interval.ms

分区分配策略

策略 说明
RangeAssignor(默认) 按分区范围分配,可能不均匀
RoundRobinAssignor 轮询分配,更均匀
StickyAssignor 尽量保持上次分配,减少迁移,Kafka 推荐
CooperativeStickyAssignor 增量式 Rebalance,Rebalance 期间不停止全部消费

4.5 Consumer Lag(消费延迟)

Consumer Lag = LEO - Committed Offset

Lag > 0 说明消息有积压
# 查看消费组 Lag
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
    --group my-group --describe

# 输出示例
GROUP       TOPIC       PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
my-group    orders      0          1000            1050            50
my-group    orders      1          980             1020            40

5. 存储机制

5.1 日志分段(Log Segment)

每个 Partition 对应磁盘上的一个目录,目录内有多个日志分段文件:

/data/kafka/orders-0/
├── 00000000000000000000.log        # 数据文件
├── 00000000000000000000.index      # 稀疏索引(offset → 文件位置)
├── 00000000000000000000.timeindex  # 时间索引(timestamp → offset)
├── 00000000000001048576.log        # 下一个分段(前一个满了之后)
├── 00000000000001048576.index
└── 00000000000001048576.timeindex

文件名 = 该分段第一条消息的 offset(20位补零)。

5.2 消息查找流程

查找 offset=368776 的消息:

1. 二分查找找到对应的 .index 文件(offset 在哪个分段)
2. 从 .index 稀疏索引找到最近的物理位置
3. 顺序扫描 .log 文件找到精确 offset

5.3 零拷贝(Zero Copy)

Kafka 消费者读取数据时使用 sendfile() 系统调用,避免数据从内核态复制到用户态再复制回内核态:

传统方式(4次拷贝):
磁盘 → 内核Buffer → 用户Buffer → Socket Buffer → 网卡

零拷贝(2次拷贝):
磁盘 → 内核Buffer ──────────────────────────────► 网卡
                  (sendfile: 内核直接到网卡,无需用户态)

这是 Kafka 高吞吐的重要原因之一。

5.4 消息保留策略

# server.properties

# 基于时间(默认 7 天)
log.retention.hours=168

# 基于大小(超过则删除最旧的 segment)
log.retention.bytes=1073741824   # 1GB per partition

# segment 文件大小(1GB 时切分)
log.segment.bytes=1073741824

# 清理策略:delete(删除)或 compact(压缩,保留每个 Key 的最新值)
log.cleanup.policy=delete

Log Compaction(日志压缩):对每个 Key 只保留最新的 Value,适合状态存储场景(如用户配置、数据库快照)。


6. 完整 Go 实战案例

6.1 依赖安装

# 推荐使用 IBM/sarama(原 Shopify/sarama)
go get github.com/IBM/sarama

# 或者使用 confluent-kafka-go(基于 librdkafka,性能更高)
go get github.com/confluentinc/confluent-kafka-go/kafka

6.2 项目结构

kafka-demo/
├── go.mod
├── pkg/
│   └── kafka/
│       ├── config.go      # 配置
│       ├── producer.go    # 生产者
│       └── consumer.go    # 消费者
├── examples/
│   ├── simple/            # 简单收发
│   ├── order_events/      # 订单事件
│   └── cdc/               # 变更数据捕获
└── main.go

6.3 配置封装

// pkg/kafka/config.go
package kafka

import (
    "time"

    "github.com/IBM/sarama"
)

type Config struct {
    Brokers  []string
    Version  string
    ClientID string
}

func NewProducerConfig(cfg Config) *sarama.Config {
    c := sarama.NewConfig()

    // 版本
    version, _ := sarama.ParseKafkaVersion(cfg.Version) // e.g. "3.6.0"
    c.Version = version
    c.ClientID = cfg.ClientID

    // 可靠性:等待所有 ISR 确认
    c.Producer.RequiredAcks = sarama.WaitForAll
    c.Producer.Retry.Max = 5
    c.Producer.Retry.Backoff = 100 * time.Millisecond

    // 幂等生产者(防重复)
    c.Producer.Idempotent = true
    c.Net.MaxOpenRequests = 1

    // 批量压缩(提升吞吐)
    c.Producer.Compression = sarama.CompressionSnappy
    c.Producer.Flush.Bytes = 1024 * 1024  // 1MB 触发发送
    c.Producer.Flush.Frequency = 500 * time.Millisecond // 或 500ms 触发

    c.Producer.Return.Successes = true
    c.Producer.Return.Errors = true

    return c
}

func NewConsumerConfig(cfg Config) *sarama.Config {
    c := sarama.NewConfig()

    version, _ := sarama.ParseKafkaVersion(cfg.Version)
    c.Version = version
    c.ClientID = cfg.ClientID

    // 从最早的 offset 开始消费(新 group 第一次消费)
    c.Consumer.Offsets.Initial = sarama.OffsetNewest // 或 sarama.OffsetOldest

    // 手动提交 offset
    c.Consumer.Offsets.AutoCommit.Enable = false

    // Rebalance 策略:粘性,减少分区迁移
    c.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{
        sarama.NewBalanceStrategySticky(),
    }

    // 心跳与会话超时
    c.Consumer.Group.Session.Timeout = 30 * time.Second
    c.Consumer.Group.Heartbeat.Interval = 3 * time.Second

    return c
}

6.4 同步生产者

// pkg/kafka/producer.go
package kafka

import (
    "fmt"
    "time"

    "github.com/IBM/sarama"
)

type SyncProducer struct {
    producer sarama.SyncProducer
}

func NewSyncProducer(brokers []string, cfg *sarama.Config) (*SyncProducer, error) {
    p, err := sarama.NewSyncProducer(brokers, cfg)
    if err != nil {
        return nil, fmt.Errorf("new sync producer: %w", err)
    }
    return &SyncProducer{producer: p}, nil
}

// Send 发送消息,返回 partition 和 offset
func (p *SyncProducer) Send(topic, key string, value []byte) (int32, int64, error) {
    msg := &sarama.ProducerMessage{
        Topic:     topic,
        Key:       sarama.StringEncoder(key),
        Value:     sarama.ByteEncoder(value),
        Timestamp: time.Now(),
    }
    partition, offset, err := p.producer.SendMessage(msg)
    if err != nil {
        return 0, 0, fmt.Errorf("send message: %w", err)
    }
    return partition, offset, nil
}

// SendBatch 批量发送
func (p *SyncProducer) SendBatch(msgs []*sarama.ProducerMessage) error {
    if err := p.producer.SendMessages(msgs); err != nil {
        if errs, ok := err.(sarama.ProducerErrors); ok {
            for _, e := range errs {
                fmt.Printf("failed to send: %v, err: %v\n", e.Msg, e.Err)
            }
        }
        return err
    }
    return nil
}

func (p *SyncProducer) Close() error {
    return p.producer.Close()
}

6.5 异步生产者(高吞吐)

// pkg/kafka/async_producer.go
package kafka

import (
    "log"
    "sync"

    "github.com/IBM/sarama"
)

type AsyncProducer struct {
    producer sarama.AsyncProducer
    wg       sync.WaitGroup
}

func NewAsyncProducer(brokers []string, cfg *sarama.Config) (*AsyncProducer, error) {
    p, err := sarama.NewAsyncProducer(brokers, cfg)
    if err != nil {
        return nil, err
    }
    ap := &AsyncProducer{producer: p}
    ap.handleResults()
    return ap, nil
}

func (p *AsyncProducer) handleResults() {
    p.wg.Add(2)

    // 处理成功
    go func() {
        defer p.wg.Done()
        for msg := range p.producer.Successes() {
            log.Printf("sent: topic=%s partition=%d offset=%d",
                msg.Topic, msg.Partition, msg.Offset)
        }
    }()

    // 处理失败(可加入重试/告警逻辑)
    go func() {
        defer p.wg.Done()
        for err := range p.producer.Errors() {
            log.Printf("failed: %v", err)
            // TODO: 写入失败消息到本地文件或告警
        }
    }()
}

// Input 发送消息(非阻塞)
func (p *AsyncProducer) Input(msg *sarama.ProducerMessage) {
    p.producer.Input() <- msg
}

func (p *AsyncProducer) Close() error {
    p.producer.AsyncClose()
    p.wg.Wait()
    return nil
}

6.6 消费者组(完整实现)

// pkg/kafka/consumer.go
package kafka

import (
    "context"
    "fmt"
    "log"
    "sync"

    "github.com/IBM/sarama"
)

// HandlerFunc 消息处理函数,返回 error 时消息不提交(等待重试)
type HandlerFunc func(ctx context.Context, msg *sarama.ConsumerMessage) error

type ConsumerGroup struct {
    group   sarama.ConsumerGroup
    topics  []string
    handler HandlerFunc
}

func NewConsumerGroup(brokers []string, groupID string, topics []string,
    cfg *sarama.Config, handler HandlerFunc) (*ConsumerGroup, error) {

    g, err := sarama.NewConsumerGroup(brokers, groupID, cfg)
    if err != nil {
        return nil, fmt.Errorf("new consumer group: %w", err)
    }
    return &ConsumerGroup{group: g, topics: topics, handler: handler}, nil
}

// Start 启动消费(阻塞,直到 ctx 取消)
func (c *ConsumerGroup) Start(ctx context.Context) error {
    h := &consumerGroupHandler{handler: c.handler}

    for {
        // 每次 Rebalance 后重新进入消费循环
        if err := c.group.Consume(ctx, c.topics, h); err != nil {
            return fmt.Errorf("consume: %w", err)
        }
        if ctx.Err() != nil {
            return ctx.Err()
        }
        log.Println("rebalance completed, resuming...")
    }
}

func (c *ConsumerGroup) Close() error {
    return c.group.Close()
}

// consumerGroupHandler 实现 sarama.ConsumerGroupHandler 接口
type consumerGroupHandler struct {
    handler HandlerFunc
    wg      sync.WaitGroup
}

func (h *consumerGroupHandler) Setup(sess sarama.ConsumerGroupSession) error {
    log.Printf("consumer group setup: member=%s generation=%d",
        sess.MemberID(), sess.GenerationID())
    return nil
}

func (h *consumerGroupHandler) Cleanup(sess sarama.ConsumerGroupSession) error {
    log.Printf("consumer group cleanup, waiting for in-flight messages...")
    h.wg.Wait()
    return nil
}

func (h *consumerGroupHandler) ConsumeClaim(
    sess sarama.ConsumerGroupSession,
    claim sarama.ConsumerGroupClaim,
) error {
    for msg := range claim.Messages() {
        h.wg.Add(1)
        // 串行处理(保证分区内有序)
        func() {
            defer h.wg.Done()
            ctx := sess.Context()
            if err := h.handler(ctx, msg); err != nil {
                log.Printf("handler error: %v, topic=%s partition=%d offset=%d",
                    err, msg.Topic, msg.Partition, msg.Offset)
                // 不提交 offset,等待重试
                return
            }
            // 标记消息,下次 CommitOffsets 时提交
            sess.MarkMessage(msg, "")
            sess.CommitOffsets()
        }()
    }
    return nil
}

6.7 订单事件处理案例

场景:电商系统订单创建后,库存服务、积分服务、通知服务各自消费订单事件。

// examples/order_events/main.go
package main

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "os"
    "os/signal"
    "syscall"
    "time"

    "github.com/IBM/sarama"
    "github.com/your/project/pkg/kafka"
)

const (
    brokerList = "localhost:9092"
    topic      = "order.created"
)

type OrderCreatedEvent struct {
    OrderID   string    `json:"order_id"`
    UserID    int64     `json:"user_id"`
    Items     []Item    `json:"items"`
    Amount    float64   `json:"amount"`
    CreatedAt time.Time `json:"created_at"`
}

type Item struct {
    SkuID    string `json:"sku_id"`
    Quantity int    `json:"quantity"`
    Price    float64 `json:"price"`
}

// ========== 生产者:订单服务 ==========

func createOrder(producer *kafka.SyncProducer, order OrderCreatedEvent) error {
    body, err := json.Marshal(order)
    if err != nil {
        return err
    }
    // 以 OrderID 为 Key,保证同一订单的消息进同一分区(顺序保证)
    partition, offset, err := producer.Send(topic, order.OrderID, body)
    if err != nil {
        return fmt.Errorf("publish order event: %w", err)
    }
    log.Printf("order event sent: order_id=%s partition=%d offset=%d",
        order.OrderID, partition, offset)
    return nil
}

// ========== 消费者:库存服务 ==========

func inventoryHandler(ctx context.Context, msg *sarama.ConsumerMessage) error {
    var event OrderCreatedEvent
    if err := json.Unmarshal(msg.Value, &event); err != nil {
        return err
    }
    log.Printf("[库存服务] 扣减库存: order_id=%s", event.OrderID)
    for _, item := range event.Items {
        log.Printf("  sku=%s qty=%d", item.SkuID, item.Quantity)
        // deductInventory(item.SkuID, item.Quantity)
    }
    return nil
}

// ========== 消费者:积分服务 ==========

func rewardHandler(ctx context.Context, msg *sarama.ConsumerMessage) error {
    var event OrderCreatedEvent
    if err := json.Unmarshal(msg.Value, &event); err != nil {
        return err
    }
    points := int(event.Amount)
    log.Printf("[积分服务] 用户 %d 获得 %d 积分", event.UserID, points)
    // addUserPoints(event.UserID, points)
    return nil
}

func main() {
    brokers := []string{brokerList}
    cfg := kafka.Config{Brokers: brokers, Version: "3.6.0", ClientID: "order-demo"}

    // 启动生产者
    go func() {
        prodCfg := kafka.NewProducerConfig(cfg)
        producer, err := kafka.NewSyncProducer(brokers, prodCfg)
        if err != nil {
            log.Fatal(err)
        }
        defer producer.Close()

        for i := 1; i <= 10; i++ {
            order := OrderCreatedEvent{
                OrderID:   fmt.Sprintf("ORD-%06d", i),
                UserID:    int64(1000 + i),
                Amount:    float64(i) * 99.9,
                CreatedAt: time.Now(),
                Items: []Item{
                    {SkuID: fmt.Sprintf("SKU-%04d", i), Quantity: 1, Price: float64(i) * 99.9},
                },
            }
            createOrder(producer, order)
            time.Sleep(200 * time.Millisecond)
        }
    }()

    ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
    defer stop()

    consCfg := kafka.NewConsumerConfig(cfg)

    // 库存服务消费者组
    inventoryGroup, _ := kafka.NewConsumerGroup(
        brokers, "inventory-service", []string{topic}, consCfg, inventoryHandler)
    defer inventoryGroup.Close()

    // 积分服务消费者组(独立消费,互不影响)
    rewardGroup, _ := kafka.NewConsumerGroup(
        brokers, "reward-service", []string{topic}, consCfg, rewardHandler)
    defer rewardGroup.Close()

    go inventoryGroup.Start(ctx)
    go rewardGroup.Start(ctx)

    <-ctx.Done()
    log.Println("shutting down...")
}

6.8 Offset 手动管理与重置

// 重置到指定时间点重新消费(数据修复场景)
func resetOffsetToTime(brokers []string, groupID, topic string, t time.Time) error {
    cfg := sarama.NewConfig()
    cfg.Version = sarama.V3_6_0_0

    admin, err := sarama.NewClusterAdmin(brokers, cfg)
    if err != nil {
        return err
    }
    defer admin.Close()

    client, err := sarama.NewClient(brokers, cfg)
    if err != nil {
        return err
    }
    defer client.Close()

    partitions, _ := client.Partitions(topic)

    offsets := make(map[string]map[int32]int64)
    offsets[topic] = make(map[int32]int64)

    for _, partition := range partitions {
        // 获取指定时间的 offset
        offset, err := client.GetOffset(topic, partition, t.UnixMilli())
        if err != nil {
            return err
        }
        offsets[topic][partition] = offset
        log.Printf("partition=%d reset to offset=%d (time=%s)", partition, offset, t)
    }

    return admin.ResetConsumerGroupOffsets(groupID, offsets)
}

6.9 实时日志收集案例

场景:多个服务将日志写入 Kafka,日志消费者统一写入 Elasticsearch。

package main

import (
    "context"
    "encoding/json"
    "log"
    "time"

    "github.com/IBM/sarama"
)

type LogEntry struct {
    Service   string    `json:"service"`
    Level     string    `json:"level"`
    Message   string    `json:"message"`
    TraceID   string    `json:"trace_id"`
    Timestamp time.Time `json:"timestamp"`
}

// 生产者:各服务打日志
func writeLog(producer sarama.AsyncProducer, entry LogEntry) {
    body, _ := json.Marshal(entry)
    producer.Input() <- &sarama.ProducerMessage{
        Topic: "app.logs",
        Key:   sarama.StringEncoder(entry.Service),
        Value: sarama.ByteEncoder(body),
    }
}

// 消费者:批量写入 ES
type ESBatcher struct {
    batch   []*LogEntry
    maxSize int
    timeout time.Duration
    ticker  *time.Ticker
}

func NewESBatcher(maxSize int, timeout time.Duration) *ESBatcher {
    return &ESBatcher{
        batch:   make([]*LogEntry, 0, maxSize),
        maxSize: maxSize,
        timeout: timeout,
        ticker:  time.NewTicker(timeout),
    }
}

func (b *ESBatcher) Process(ctx context.Context, msg *sarama.ConsumerMessage) error {
    var entry LogEntry
    if err := json.Unmarshal(msg.Value, &entry); err != nil {
        return err
    }
    b.batch = append(b.batch, &entry)

    select {
    case <-b.ticker.C:
        return b.flush()
    default:
        if len(b.batch) >= b.maxSize {
            return b.flush()
        }
    }
    return nil
}

func (b *ESBatcher) flush() error {
    if len(b.batch) == 0 {
        return nil
    }
    log.Printf("flushing %d logs to ES", len(b.batch))
    // bulkInsertToES(b.batch)
    b.batch = b.batch[:0]
    return nil
}

6.10 Exactly Once 语义实现

// 使用 Kafka 事务实现跨 Topic 的 Exactly Once
// 场景:从 source-topic 消费,处理后写入 dest-topic,保证原子性

func processWithExactlyOnce(
    client sarama.Client,
    srcTopic, dstTopic string,
    groupID string,
) error {
    // 事务生产者
    txnCfg := sarama.NewConfig()
    txnCfg.Producer.Transaction.ID = "eos-processor-1"
    txnCfg.Producer.RequiredAcks = sarama.WaitForAll
    txnCfg.Producer.Idempotent = true
    txnCfg.Net.MaxOpenRequests = 1

    producer, err := sarama.NewSyncProducer(client.Brokers(), txnCfg)
    if err != nil {
        return err
    }
    defer producer.Close()

    // 普通消费者(offset 通过事务提交,不使用 consumer group commit)
    consumer, _ := sarama.NewConsumerFromClient(client)
    pc, _ := consumer.ConsumePartition(srcTopic, 0, sarama.OffsetNewest)

    for msg := range pc.Messages() {
        // 开始事务
        producer.(sarama.SyncProducer)

        // 处理消息
        result := process(msg.Value)

        // 在事务中提交 offset 和写入目标 Topic
        // kafka 事务确保这两个操作原子执行
        offsets := map[string][]*sarama.PartitionOffsetMetadata{
            srcTopic: {{Partition: msg.Partition, Offset: msg.Offset + 1}},
        }
        _ = offsets // 实际使用 producer.AddOffsetsToTxn(offsets, groupID)

        producer.SendMessages([]*sarama.ProducerMessage{
            {Topic: dstTopic, Value: sarama.ByteEncoder(result)},
        })
    }
    return nil
}

func process(data []byte) []byte {
    // 业务处理逻辑
    return data
}

7. Kafka 集群部署

7.1 KRaft 模式(推荐,无 ZooKeeper)

Kafka 2.8 引入、3.3+ 生产可用的 KRaft 模式,消除了对 ZooKeeper 的依赖:

旧架构(ZooKeeper 模式):
Kafka Brokers ←──► ZooKeeper Cluster(元数据存储)

新架构(KRaft 模式):
Kafka Brokers(内置 Raft Controller,自管元数据)

KRaft 优势

  • 启动更快,无需等待 ZooKeeper
  • 架构更简单,运维复杂度降低
  • 支持更多 Partition(ZooKeeper 模式下 ~200K,KRaft 下 ~数百万)

7.2 Docker Compose 三节点集群(KRaft 模式)

# docker-compose.yml
version: '3.8'

x-kafka-common: &kafka-common
  image: apache/kafka:3.7.0
  environment: &kafka-env
    KAFKA_PROCESS_ROLES: broker,controller
    KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka1:9093,2@kafka2:9093,3@kafka3:9093
    KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
    KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
    KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
    KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
    KAFKA_DEFAULT_REPLICATION_FACTOR: "3"
    KAFKA_MIN_INSYNC_REPLICAS: "2"
    KAFKA_NUM_PARTITIONS: "6"
    KAFKA_LOG_RETENTION_HOURS: "168"
    KAFKA_LOG_SEGMENT_BYTES: "1073741824"
    CLUSTER_ID: "MkU3OEVBNTcwNTJENDM2Qk"  # 需固定,使用 kafka-storage random-uuid 生成

services:
  kafka1:
    <<: *kafka-common
    hostname: kafka1
    environment:
      <<: *kafka-env
      KAFKA_NODE_ID: "1"
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
    ports:
      - "9092:9092"
    volumes:
      - kafka1_data:/var/lib/kafka/data

  kafka2:
    <<: *kafka-common
    hostname: kafka2
    environment:
      <<: *kafka-env
      KAFKA_NODE_ID: "2"
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9093
    ports:
      - "9093:9092"
    volumes:
      - kafka2_data:/var/lib/kafka/data

  kafka3:
    <<: *kafka-common
    hostname: kafka3
    environment:
      <<: *kafka-env
      KAFKA_NODE_ID: "3"
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9094
    ports:
      - "9094:9092"
    volumes:
      - kafka3_data:/var/lib/kafka/data

  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka1:9092,kafka2:9092,kafka3:9092

volumes:
  kafka1_data:
  kafka2_data:
  kafka3_data:

7.3 Topic 管理

# 创建 Topic(6个分区,3个副本)
kafka-topics.sh --bootstrap-server localhost:9092 \
    --create --topic orders \
    --partitions 6 \
    --replication-factor 3 \
    --config retention.ms=604800000 \
    --config min.insync.replicas=2

# 查看 Topic 详情
kafka-topics.sh --bootstrap-server localhost:9092 \
    --describe --topic orders

# 修改分区数(只能增加)
kafka-topics.sh --bootstrap-server localhost:9092 \
    --alter --topic orders --partitions 12

# 删除 Topic
kafka-topics.sh --bootstrap-server localhost:9092 \
    --delete --topic orders

# 列出所有 Topic
kafka-topics.sh --bootstrap-server localhost:9092 --list

7.4 生产消费调试

# 生产消息(命令行)
kafka-console-producer.sh --bootstrap-server localhost:9092 \
    --topic orders \
    --property "key.separator=:" \
    --property "parse.key=true"
# 输入: order-001:{"amount":100}

# 消费消息(从头开始)
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
    --topic orders \
    --from-beginning \
    --property print.key=true \
    --property print.partition=true \
    --property print.offset=true

# 查看消费组 Lag
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
    --group inventory-service --describe

# 重置消费 offset 到最早
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
    --group inventory-service \
    --topic orders \
    --reset-offsets --to-earliest \
    --execute

8. 监控与调优

8.1 关键指标

指标 说明 告警阈值
kafka_consumer_group_lag 消费延迟(积压) > 10000
kafka_server_bytesInPerSec 写入吞吐 接近网卡上限
kafka_server_bytesOutPerSec 读取吞吐 接近网卡上限
kafka_server_underReplicatedPartitions 未完全复制的分区数 > 0 立即告警
kafka_controller_activeControllerCount 集群 Controller 数量 不等于1时告警
kafka_server_requestHandlerAvgIdlePercent 请求处理线程空闲率 < 20% 时告警
kafka_log_logFlushRateAndTimeMs 日志刷盘延迟 持续 > 1000ms

8.2 JVM 调优(Broker)

# kafka-server-start.sh
export KAFKA_HEAP_OPTS="-Xms6g -Xmx6g"
export KAFKA_JVM_PERFORMANCE_OPTS="-XX:+UseG1GC \
    -XX:MaxGCPauseMillis=20 \
    -XX:InitiatingHeapOccupancyPercent=35 \
    -XX:G1HeapRegionSize=16M \
    -XX:MinMetaspaceFreeRatio=50 \
    -XX:MaxMetaspaceFreeRatio=80"

8.3 OS 调优

# 增大文件句柄数(每个分区需要多个文件句柄)
echo "* soft nofile 1000000" >> /etc/security/limits.conf
echo "* hard nofile 1000000" >> /etc/security/limits.conf

# 调整内核参数
sysctl -w vm.swappiness=1          # 几乎不用 swap
sysctl -w net.core.rmem_max=134217728
sysctl -w net.core.wmem_max=134217728
sysctl -w net.ipv4.tcp_rmem="4096 131072 134217728"
sysctl -w net.ipv4.tcp_wmem="4096 131072 134217728"

8.4 Producer 调优参数

// 追求吞吐量
config.Producer.Flush.Bytes = 5 * 1024 * 1024    // 5MB 批次
config.Producer.Flush.Frequency = time.Second      // 或1秒发一次
config.Producer.Compression = sarama.CompressionLZ4 // LZ4 压缩,速度最快

// 追求低延迟
config.Producer.Flush.Bytes = 1                    // 立即发送
config.Producer.Flush.Frequency = time.Millisecond
config.Producer.Compression = sarama.CompressionNone

8.5 Consumer 调优参数

// 增大每次 fetch 的数据量,减少 fetch 次数
config.Consumer.Fetch.Min = 1024 * 1024     // 最小 fetch 1MB
config.Consumer.Fetch.Default = 10 * 1024 * 1024 // 默认 fetch 10MB
config.Consumer.MaxWaitTime = 500 * time.Millisecond // 最长等待 500ms

9. 高频面试题

Q1: Kafka 为什么吞吐量这么高?

五大核心原因

① 顺序写磁盘 Kafka 只追加写日志文件(Append Only),磁盘顺序写性能接近内存随机写(约 600MB/s vs 100MB/s)。

② 零拷贝(Zero Copy) 消费者读取数据使用 sendfile() 系统调用,数据从页缓存直接到网卡,减少 2 次数据拷贝和 2 次系统调用。

③ 批量发送与压缩 Producer 将多条消息打包成批次(Batch)发送,Consumer 批量拉取(Fetch),配合 Snappy/LZ4/ZSTD 压缩大幅减少网络传输量。

④ 页缓存(Page Cache) Kafka 利用操作系统的页缓存而非 JVM 堆内存,避免 GC 影响,读热数据直接从内存返回。

⑤ 分区并行 多 Partition 可并行写入和读取,水平扩展吞吐量。

Q2: Kafka 如何保证消息不丢失?

三端保障

① 生产者端

acks=-1(等待所有 ISR 确认)
+ min.insync.replicas=2(至少2个ISR副本)
+ retries > 0(失败重试)
+ 幂等生产者(防重复)

② Broker 端

replication.factor >= 3(3副本)
+ unclean.leader.election.enable=false(禁止非ISR副本当选Leader,防止丢数据)
+ min.insync.replicas=2

③ 消费者端

手动提交 offset(autoCommit=false)
+ 业务处理成功后再 CommitOffsets

Q3: Kafka 消息重复消费如何处理?

原因:消费者处理成功但还没提交 offset 时崩溃,重启后重新从上次提交的 offset 消费。

解决方案(幂等消费):

func handleMessage(ctx context.Context, msg *sarama.ConsumerMessage) error {
    // 从消息中提取业务唯一ID(或使用 topic+partition+offset 作为唯一ID)
    var payload struct {
        OrderID string `json:"order_id"`
    }
    json.Unmarshal(msg.Value, &payload)

    // 方案1:Redis 去重
    key := fmt.Sprintf("kafka:processed:%s:%d:%d",
        msg.Topic, msg.Partition, msg.Offset)
    ok, _ := redis.SetNX(ctx, key, 1, 24*time.Hour)
    if !ok {
        return nil // 已处理,直接跳过
    }

    // 方案2:数据库唯一索引
    // INSERT IGNORE INTO processed_messages (topic, partition, offset) VALUES (?,?,?)

    return processOrder(payload.OrderID)
}

Q4: Kafka 如何保证消息有序?

分区内有序,分区间无序

场景 方案
同一业务实体严格有序(如同一订单) 使用 OrderID 作为 Key,相同 Key 路由到同一分区
全局有序(吞吐低) 只用 1 个分区 + 1 个消费者
业务层面排序 消息带序号,消费者暂存乱序消息等待补齐

注意:即使 Key 相同,修改分区数后相同 Key 可能路由到不同分区,顺序被破坏。

Q5: Kafka Rebalance 有什么问题,如何优化?

问题:Rebalance 期间所有 Consumer 停止消费(Stop The World),大规模场景下延迟可达数秒到数分钟。

优化方向

手段 说明
增大 session.timeout.ms 避免因短暂 GC/网络抖动误判 Consumer 下线
增大 max.poll.interval.ms 给消费逻辑更多时间,避免因处理慢被踢出组
使用 StickyAssignor 尽量复用上次分配,减少分区迁移数量
使用 CooperativeStickyAssignor 增量式 Rebalance,未迁移的分区继续消费
减少 Consumer 动态变化 使用静态成员 ID(group.instance.id),重启不触发 Rebalance
// 静态成员(固定 instance ID,重启不触发 Rebalance)
config.Consumer.Group.InstanceID = "consumer-instance-1"

Q6: Kafka 的 ISR 机制是什么?

ISR(In-Sync Replicas):与 Leader 保持同步的副本集合。

加入/移出 ISR 的条件

  • Follower 落后 Leader 的消息数超过 replica.lag.time.max.ms(默认30s)内未同步完,被踢出 ISR
  • Follower 追上 Leader 后,重新加入 ISR

为什么需要 ISR

  • acks=-1 只需要 ISR 内所有副本确认,不需要等待落后的 Follower
  • Leader 宕机时,只有 ISR 内的 Follower 才能参与选举,保证数据不丢失
ISR = {Broker-1(Leader), Broker-2}    Broker-3 落后被踢出 ISR

acks=-1 + min.insync.replicas=2:
只要 Broker-1 和 Broker-2 确认,消息即可返回成功
Broker-3 异步追赶即可

Q7: Kafka 的 HW 和 LEO 是什么关系?

LEO(Log End Offset):每个副本自身下一条写入的位置(最新 offset + 1)
HW(High Watermark):所有 ISR 副本 LEO 的最小值

规则:
- 消费者只能读取 HW 以下的消息(防止读到未完全复制的消息)
- Leader 宕机后,新 Leader 的 HW 之前的数据才是安全的

示例:
  Leader  LEO=100
  Follow1 LEO=98
  Follow2 LEO=95
  → HW = 95,消费者最多读到 offset 94

Q8: Kafka 和 RabbitMQ 怎么选?

场景 选择 理由
日志收集 / 大数据管道 Kafka 高吞吐、持久化、可回放
实时流处理 Kafka 与 Flink/Spark 生态集成好
复杂路由(Exchange 规则) RabbitMQ 原生支持多种路由策略
延迟消息 RabbitMQ(插件)或 RocketMQ Kafka 原生不支持
消息可靠投递 + 业务解耦 RabbitMQ 轻量,运维简单
顺序消息 + 事务消息 RocketMQ 原生支持
超高吞吐(百万级/s) Kafka 性能天花板最高

Q9: Kafka 消息积压怎么处理?

排查步骤

# 1. 查看积压量
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
    --group my-group --describe

# 2. 查看 Consumer 处理速度(Consume Rate vs Produce Rate)
# 通过 Prometheus/JMX 看 records-consumed-rate vs records-produced-rate

处理方案

方案 适用场景
扩容消费者(≤ 分区数) 消费能力不足,分区数足够
先扩分区再扩消费者 消费者数量受分区数限制
临时新建 Topic 迁移 堆积量极大,需要快速消化:新建更多分区的 Topic,将消息转移
优化消费逻辑 消费慢是因为 DB 慢查询、外部调用超时
跳过积压消息 非关键消息,重置 offset 到最新
// 并发消费同一分区的消息(牺牲顺序换吞吐)
func (h *handler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
    sem := make(chan struct{}, 20) // 20并发
    var wg sync.WaitGroup
    for msg := range claim.Messages() {
        sem <- struct{}{}
        wg.Add(1)
        go func(m *sarama.ConsumerMessage) {
            defer func() { <-sem; wg.Done() }()
            h.process(m)
            sess.MarkMessage(m, "")
        }(msg)
    }
    wg.Wait()
    return nil
}

Q10: unclean.leader.election.enable 是什么,为什么要关闭?

当 ISR 中所有副本都不可用时,是否允许非 ISR 副本(落后较多的副本)成为 Leader:

unclean.leader.election.enable=true(默认 false):
  优点:可用性高,即使 ISR 全挂也能选出 Leader 继续服务
  缺点:可能丢失未同步的消息(数据不一致)

unclean.leader.election.enable=false(推荐):
  优点:不丢数据,一致性强
  缺点:ISR 全挂时 Partition 不可用(需等 ISR 副本恢复)

生产环境建议关闭,配合足够的副本数(3+)和监控,通过运维恢复节点而非允许脏选举。


10. 最佳实践总结

10.1 Topic 设计

✅ 分区数 = 目标消费者数 × N(留有扩展余地)
✅ Replication Factor = 3(生产环境)
✅ min.insync.replicas = 2(防止脑裂丢数据)
✅ retention.ms 根据业务设置(通常 7-30 天)
✅ 同一业务实体的消息用相同的 Key(保证顺序)
❌ 不要频繁修改分区数(破坏 Key 路由的有序性)

10.2 生产者最佳实践

✅ acks=-1 + retries + 幂等生产者(重要业务)
✅ 根据吞吐/延迟权衡 batch.size 和 linger.ms
✅ 启用压缩(Snappy/LZ4,大消息用 ZSTD)
✅ 监控 record-error-rate,发现发送失败
✅ 异步发送时一定要处理 Errors channel
❌ 不要在主业务流程中使用同步发送的超长超时

10.3 消费者最佳实践

✅ 关闭自动提交(autoCommit=false),手动控制 offset
✅ 消费逻辑实现幂等(Redis/DB 去重)
✅ 合理配置 session.timeout 和 max.poll.interval
✅ 使用 StickyAssignor 或 CooperativeStickyAssignor
✅ 监控 Consumer Lag,设置积压告警
✅ 优雅关闭:先停止消费,处理完 in-flight 消息,再提交 offset
❌ 不要在 ConsumeClaim 中做无限阻塞操作
❌ 不要在消费者内同步调用外部慢服务(超时会触发 Rebalance)

10.4 常见陷阱

❌ acks=1 + Follower 还没同步 → Leader 宕机 → 消息丢失
✅ 改为 acks=-1 + min.insync.replicas=2

❌ 消费者处理慢超过 max.poll.interval.ms → 被踢出消费组 → Rebalance 循环
✅ 增大 max.poll.interval.ms 或优化消费逻辑

❌ 期望 Kafka 支持延迟消息(不支持)
✅ 改用 RocketMQ 或在 Kafka 上层实现延迟队列(时间轮 + 状态Topic)

❌ 用 Kafka 做 RPC(请求-响应模式),延迟高且复杂
✅ RPC 场景改用 gRPC 或 HTTP

❌ 修改了分区数但没更新 Key 路由逻辑
✅ 修改分区数后,相同 Key 的路由分区发生变化,需评估影响

小结

知识点 核心要点
高吞吐原因 顺序写 + 零拷贝 + 批量 + 页缓存 + 分区并行
分区 并发单元,分区内有序;Key 哈希保证同 Key 同分区
副本 & ISR Leader 读写,Follower 同步;ISR 副本参与选举,保证不丢数据
HW & LEO 消费者只读 HW 以下消息;HW = ISR 中最小 LEO
消息可靠 acks=-1 + min.insync.replicas=2 + 手动 offset + 幂等
Rebalance Consumer 变化触发;StickyAssignor + 静态 ID 减少影响
消息积压 扩容 Consumer(≤ 分区数)+ 优化消费逻辑
选型 高吞吐/流处理/日志选 Kafka;复杂路由/延迟消息选 RabbitMQ/RocketMQ