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
ProducerRecord封装消息(topic, partition, key, value, headers, timestamp)- 序列化:将 key/value 序列化为字节数组
- 分区器:决定消息发到哪个 Partition
- RecordAccumulator:按 Partition 缓冲成批次(Batch),减少网络请求
- 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 |
xingliuhua