阶段一 · 基础与核心模型

本地开发环境与运维基础

一句话总结

用 docker-compose 一键拉起 3-Broker KRaft 集群 + kafka-ui,所有实验可本地复现.
掌握 CLI 基本操作(创建 Topic、生产消费、查看 Group lag),以及生产环境最该关注的几个监控指标.

docker-compose:3-Broker KRaft 集群

Kafka 3.3+ 支持 KRaft 模式(不依赖 ZooKeeper),本地开发推荐直接用 KRaft,更轻量.

# docker-compose.yml
services:
kafka-0:
image: bitnami/kafka:3.7
ports:
- "9092:9092"
environment:
- KAFKA_CFG_NODE_ID=0
- KAFKA_CFG_PROCESS_ROLES=broker,controller
- KAFKA_CFG_LISTENERS=PLAINTEXT://:19092,CONTROLLER://:9093,EXTERNAL://:9092
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka-0:19092,EXTERNAL://localhost:9092
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka-0:9093,1@kafka-1:9093,2@kafka-2:9093
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_INTER_BROKER_LISTENER_NAME=PLAINTEXT
- KAFKA_KRAFT_CLUSTER_ID=local-kraft-cluster-001
- KAFKA_CFG_DEFAULT_REPLICATION_FACTOR=3
- KAFKA_CFG_MIN_INSYNC_REPLICAS=2
- KAFKA_CFG_NUM_PARTITIONS=6
volumes:
- kafka-0-data:/bitnami/kafka

kafka-1:
image: bitnami/kafka:3.7
ports:
- "9093:9092"
environment:
- KAFKA_CFG_NODE_ID=1
- KAFKA_CFG_PROCESS_ROLES=broker,controller
- KAFKA_CFG_LISTENERS=PLAINTEXT://:19092,CONTROLLER://:9093,EXTERNAL://:9092
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka-1:19092,EXTERNAL://localhost:9093
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka-0:9093,1@kafka-1:9093,2@kafka-2:9093
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_INTER_BROKER_LISTENER_NAME=PLAINTEXT
- KAFKA_KRAFT_CLUSTER_ID=local-kraft-cluster-001
- KAFKA_CFG_DEFAULT_REPLICATION_FACTOR=3
- KAFKA_CFG_MIN_INSYNC_REPLICAS=2
- KAFKA_CFG_NUM_PARTITIONS=6
volumes:
- kafka-1-data:/bitnami/kafka

kafka-2:
image: bitnami/kafka:3.7
ports:
- "9094:9092"
environment:
- KAFKA_CFG_NODE_ID=2
- KAFKA_CFG_PROCESS_ROLES=broker,controller
- KAFKA_CFG_LISTENERS=PLAINTEXT://:19092,CONTROLLER://:9093,EXTERNAL://:9092
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka-2:19092,EXTERNAL://localhost:9094
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka-0:9093,1@kafka-1:9093,2@kafka-2:9093
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_INTER_BROKER_LISTENER_NAME=PLAINTEXT
- KAFKA_KRAFT_CLUSTER_ID=local-kraft-cluster-001
- KAFKA_CFG_DEFAULT_REPLICATION_FACTOR=3
- KAFKA_CFG_MIN_INSYNC_REPLICAS=2
- KAFKA_CFG_NUM_PARTITIONS=6
volumes:
- kafka-2-data:/bitnami/kafka

kafka-ui:
image: provectuslabs/kafka-ui:latest
ports:
- "8080:8080"
environment:
- KAFKA_CLUSTERS_0_NAME=local
- KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=kafka-0:19092,kafka-1:19092,kafka-2:19092
depends_on:
- kafka-0
- kafka-1
- kafka-2

volumes:
kafka-0-data:
kafka-1-data:
kafka-2-data:

启动:

docker compose up -d

# 验证集群状态
docker compose ps
# 浏览器打开 http://localhost:8080 查看 kafka-ui

为什么用 bitnami/kafka

bitnami 镜像开箱支持 KRaft,环境变量即可完成全部配置,无需手动格式化日志目录.
另一个选择是 apache/kafka(官方镜像,3.7+ 提供),配置方式略有不同.


连接方式

Go 代码连接本地集群:

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092", "localhost:9093", "localhost:9094"),
// ...
)

三个端口分别映射到三个 Broker 的 EXTERNAL listener.SeedBrokers 只需要填一个也能工作(客户端会通过 metadata 发现全部 Broker),但填多个增加初始连接的容错性.


CLI 工具:kafka-topics / kafka-console-*

进入任一 Kafka 容器执行 CLI 命令:

# 进入容器
docker compose exec kafka-0 bash

Topic 管理

# 创建 Topic
kafka-topics.sh --create \
--topic order-events \
--partitions 6 \
--replication-factor 3 \
--bootstrap-server localhost:19092

# 查看所有 Topic
kafka-topics.sh --list --bootstrap-server localhost:19092

# 查看 Topic 详情(分区分布、副本、ISR)
kafka-topics.sh --describe --topic order-events --bootstrap-server localhost:19092

# 修改分区数(只能增不能减)
kafka-topics.sh --alter --topic order-events --partitions 12 --bootstrap-server localhost:19092

# 删除 Topic
kafka-topics.sh --delete --topic order-events --bootstrap-server localhost:19092

命令行生产/消费(调试利器)

# 生产消息(交互式输入,Ctrl+C 退出)
kafka-console-producer.sh \
--topic order-events \
--property "key.separator=:" \
--property "parse.key=true" \
--bootstrap-server localhost:19092
# 输入: user-123:{"action":"purchase"}

# 消费消息(从头开始读)
kafka-console-consumer.sh \
--topic order-events \
--from-beginning \
--property "print.key=true" \
--property "key.separator= | " \
--bootstrap-server localhost:19092

# 消费指定 Partition 的指定 offset 范围
kafka-console-consumer.sh \
--topic order-events \
--partition 0 \
--offset 100 \
--max-messages 10 \
--bootstrap-server localhost:19092

Consumer Group 管理

# 查看所有 Consumer Group
kafka-consumer-groups.sh --list --bootstrap-server localhost:19092

# 查看 Group 详情(每个 Partition 的 lag)
kafka-consumer-groups.sh --describe --group order-processor --bootstrap-server localhost:19092

输出示例:

GROUP            TOPIC          PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG   CONSUMER-ID          HOST
order-processor order-events 0 1520 1525 5 consumer-1-xxx /172.18.0.4
order-processor order-events 1 980 980 0 consumer-1-xxx /172.18.0.4
order-processor order-events 2 2100 2350 250 consumer-2-xxx /172.18.0.5
列 含义
CURRENT-OFFSET Consumer 已提交的 offset
LOG-END-OFFSET Partition 最新消息的 offset
LAG 积压量 = LOG-END-OFFSET - CURRENT-OFFSET

LAG 是最重要的运维指标

LAG 持续增长意味着 Consumer 跟不上生产速率.如果 LAG 增长到接近 retention 时间对应的消息量,老消息会被删除,导致数据丢失.
生产环境必须对 LAG 设告警.

重置 offset(数据回放)

# 重置到最早(重新消费全部消息)
kafka-consumer-groups.sh --reset-offsets \
--group order-processor \
--topic order-events \
--to-earliest \
--execute \
--bootstrap-server localhost:19092

# 重置到指定时间点
kafka-consumer-groups.sh --reset-offsets \
--group order-processor \
--topic order-events \
--to-datetime "2025-06-01T00:00:00.000" \
--execute \
--bootstrap-server localhost:19092

# 重置到最新(跳过所有积压)
kafka-consumer-groups.sh --reset-offsets \
--group order-processor \
--topic order-events \
--to-latest \
--execute \
--bootstrap-server localhost:19092

重置 offset 前必须先停消费者

Group 内有活跃 Consumer 时无法重置 offset.先停掉所有 Consumer 实例,再执行 reset.


现代 CLI 替代:kcat(原 kafkacat)

kcat 比官方 CLI 更轻量,适合本地开发快速调试:

# macOS 安装
brew install kcat

# 生产
echo '{"action":"test"}' | kcat -P -b localhost:9092 -t order-events -k "user-123"

# 消费(从尾部开始,显示 metadata)
kcat -C -b localhost:9092 -t order-events -o end -f '%T %k | %s\n'

# 查看 Topic 的 metadata(分区、副本、ISR)
kcat -L -b localhost:9092 -t order-events

kafka-ui:可视化管理

打开 http://localhost:8080 后,kafka-ui 提供:

功能 用途
Topics 列表 查看分区数、副本、消息量、配置
Topic 消息浏览 按 Partition/offset/时间查看消息内容
Consumer Groups 查看每个 Group 的 lag,定位慢消费
Brokers 查看集群节点状态、配置
生产消息 直接在 UI 里发测试消息(免 CLI)

开发阶段用 kafka-ui 比 CLI 直观得多,尤其是查看消息内容和 Group lag.


生产环境监控指标

本地开发时不需要完整监控,但了解关键指标有助于理解系统行为.生产环境需要关注的 Top 指标:

Broker 端

指标 含义 告警阈值建议
UnderReplicatedPartitions 副本未同步的 Partition 数 > 0 持续 5min
IsrShrinksPerSec ISR 缩减速率 突增
ActiveControllerCount 集群中 Controller 数量 ≠ 1
OfflinePartitionsCount 不可用 Partition 数 > 0
BytesInPerSec / BytesOutPerSec 流量 接近网卡/磁盘上限
RequestHandlerAvgIdlePercent 请求处理线程空闲率 < 30%

Producer 端

指标 含义 关注点
record-send-rate 每秒发送消息数 基线波动
record-error-rate 每秒发送失败数 > 0
request-latency-avg 平均请求延迟 突增
buffer-available-bytes 缓冲区剩余空间 接近 0 表示背压
batch-size-avg 平均 batch 大小 远小于 batch.size 说明 linger 太短

Consumer 端

指标 含义 关注点
Consumer Lag 积压消息数 最核心指标,持续增长必须告警
records-consumed-rate 每秒消费消息数 和 produce rate 对比
commit-latency-avg offset 提交延迟 突增可能表示 Coordinator 压力大
rebalance-rate-per-hour 每小时 rebalance 次数 频繁 rebalance 说明有不稳定的 Consumer

franz-go 中的 Metrics Hook

franz-go 提供 kgo.WithHooks 接口,可以接入 Prometheus:

import "github.com/twmb/franz-go/plugin/kprom"

metrics := kprom.NewMetrics("kafka")

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.WithHooks(metrics),
// ...
)

// metrics.Handler() 是标准的 http.Handler,暴露给 Prometheus 抓取
http.Handle("/metrics", metrics.Handler())

暴露的指标包括:连接数、请求延迟、produce/fetch 字节数、错误计数等,开箱可用.


日常开发 Cheatsheet

场景 命令
启动集群 docker compose up -d
停止集群(保留数据) docker compose stop
销毁集群(清空数据) docker compose down -v
查看 Broker 日志 docker compose logs -f kafka-0
快速发测试消息 echo 'test' | kcat -P -b localhost:9092 -t my-topic
看消息内容 kcat -C -b localhost:9092 -t my-topic -o beginning -c 10
看 Group lag kafka-consumer-groups.sh --describe --group xxx --bootstrap-server localhost:19092
重置 offset 先停 Consumer,再 --reset-offsets --to-earliest --execute
杀一个 Broker 测试容错 docker compose stop kafka-1

小结

要点 说明
KRaft 模式 3.3+ 无需 ZooKeeper,本地开发更轻量
docker-compose 一键起 3 Broker + kafka-ui,复制即用
CLI 三板斧 kafka-topics.sh 管 Topic,kafka-consumer-groups.sh 管 Group,kcat 快速调试
LAG 是命脉 Consumer lag 持续增长 = 即将丢数据,必须告警
Metrics Hook franz-go + kprom 开箱暴露 Prometheus 指标

至此,阶段一(基础与核心模型)7 篇全部完成.你已经掌握了 Kafka 的核心概念、生产消费模型、Go 客户端用法、可靠性保证和本地环境搭建.下一阶段我们深入 Kafka 的内部架构:存储引擎为什么快、副本机制怎么同步、Controller 怎么选举.