6.5 Apache Kafka 消息队列深入——架构、安装与生产调优
预计阅读时间:22 分钟
📖 目录
6.4:消息队列入门 介绍了 RabbitMQ 和 Redis Stream 的基本用法。Apache Kafka 是高吞吐、分布式事件流平台,适合日志聚合、流处理、事件溯源等场景。Kafka 单集群可达百万级 QPS,是大数据和微服务架构的消息骨干。
学习目标
- 理解 Kafka 的核心架构:Broker、Topic、Partition、Offset 的分工与设计动机
- 掌握 KRaft 模式下的集群安装、Topic 创建与数据生产消费
- 区分三种消息语义(At-Most-Once / At-Least-Once / Exactly-Once)及其实现代价
- 理解 Consumer Group 的分区分配与再平衡机制,能排查消费堆积
- 结合业务场景在 Kafka / RabbitMQ / Redis Stream 之间做出正确选型
- 完成基础生产调优:批量参数、分区数与副本数的权衡
前置知识
- 6.4:消息队列入门 消息队列入门——RabbitMQ 与 Redis Stream 的基础概念与术语
- 4.5:网络故障排查 网络故障排查——理解端口、监听与防火墙(Kafka 默认 9092)
- 3.12:系统监控与告警 系统监控与告警——消费堆积排查依赖监控数据
- 2.1:Shell 脚本入门 Shell 脚本入门——运行安装与运维脚本需要的基础
架构概览
| 组件 | 角色 | 说明 |
|---|---|---|
| Broker | 消息代理 | 存储消息的服务器节点,集群由多个 Broker 组成 |
| Topic | 消息主题 | 逻辑分类,每个 Topic 可有多个 Partition |
| Partition | 分区 | 并行处理单元,消息按 key 哈希分配到分区 |
| Replica | 副本 | 每个 Partition 的多份拷贝,分布在不同 Broker |
| Producer | 生产者 | 发送消息的客户端 |
| Consumer Group | 消费者组 | 组内各消费者消费不同分区,实现负载均衡 |
| ZooKeeper / KRaft | 元数据管理 | KRaft(Kafka 3.3+)替代 ZooKeeper 做选举 |
数据流全景
一条消息从生产到消费的完整路径如下(文字版架构图):
Producer ──┐
│ batch.size + linger.ms 攒批
v
┌─────────┐ partition 0 ──► Replica(Leader) ──► ISR(Followers) ──► 消费者读 HW
│ Broker │ partition 1 ──► Replica(Leader) ──► ISR(Followers)
│ Topic │ partition 2 ──► ...
└─────────┘
│
│ Consumer Group 按分区订阅
v
Consumer A ── partition 0, 2
Consumer B ── partition 1, 3
每个消费者独立维护 offset
- Producer → Broker:消息先进入客户端缓冲区,达到
batch.size或linger.ms才批量发送,这是 Kafka 高吞吐的第一层来源。 - Broker → Partition:Broker 根据分区策略决定消息写入哪个分区;每个分区是追加写(顺序 IO)的日志文件,落盘后按
replication.factor同步到 ISR 副本。 - Partition → Consumer:Consumer Group 内的消费者与分区按分配策略绑定,一个分区同一时刻只被组内一个消费者消费,消费进度(offset)持久化在
__consumer_offsets内部 Topic。
这条链路决定了 Kafka 的两个特性:分区内有序(同一分区追加写、顺序读)和组内水平扩展(加消费者只触发分区再分配,不影响已落盘数据)。
1. 安装与集群部署
# 下载 Kafka(KRaft 模式,无 ZooKeeper)
wget https://downloads.apache.org/kafka/4.0.0/kafka_2.13-4.0.0.tgz
tar xzf kafka_2.13-4.0.0.tgz
cd kafka_2.13-4.0.0
# 生成集群 UUID
KAFKA_CLUSTER_ID=$(bin/kafka-storage.sh random-uuid)
# 格式化存储目录
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID \
-c config/kraft/server.properties
# 启动 Broker
bin/kafka-server-start.sh config/kraft/server.properties
# 验证
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
多节点集群
# 每个 Broker 配置不同 node.id 和 listeners
# node1
broker.id=1
listeners=PLAINTEXT://node1:9092
controller.quorum.voters=1@node1:9093,2@node2:9093,3@node3:9093
# node2
broker.id=2
listeners=PLAINTEXT://node2:9092
controller.quorum.voters=1@node1:9093,2@node2:9093,3@node3:9093
# node3
broker.id=3
listeners=PLAINTEXT://node3:9092
controller.quorum.voters=1@node1:9093,2@node2:9093,3@node3:9093
# 3 节点集群可容忍 1 节点故障
2. Topic 管理
# 创建 Topic(6 分区,3 副本)
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic orders \
--partitions 6 \
--replication-factor 3
# 查看 Topic 列表
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
# 查看 Topic 详情
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--describe --topic orders
# 修改分区数(只能增加,不能减少)
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--alter --topic orders --partitions 12
Topic 命名与生命周期管理
Topic 名称是团队间的事实标准,命名规范直接影响运维效率;保留策略按消息性质分三类配置:
# 命名规范:<域>-<业务>-<事件>(如 order-created / user-registered)
# 便于按前缀分权、按模式做监控告警聚合、防止同名冲突
# 类型一:日志型(访问日志、审计)——delete + 按时间/容量保留
bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter \
--entity-type topics --entity-name access-log \
--add-config cleanup.policy=delete,retention.ms=43200000,retention.bytes=536870912000
# 类型二:事件型(订单、支付)——delete + 按业务重放窗口保留(如 7 天)
# 类型三:状态型(用户最新状态、配置)——compact,按 key 保留最新值
bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter \
--entity-type topics --entity-name user-profile \
--add-config cleanup.policy=compact
# 生命周期管理:Topic 创建/扩分区/改保留策略全部走变更单,
# 避免"随手建 Topic"导致无人清理、磁盘被占满
分区数量规划:先按峰值吞吐估算(见第 7 节公式),再考虑未来 6-12 个月的业务增长——分区只能增加不能减少,但单 Broker 分区数超过 4000 会拖垮元数据管理。
3. 生产者与消费者
命令行快速测试
# 生产者(控制台输入消息)
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders
# 消费者(从最早消息开始读)
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders --from-beginning --group test-consumer
Java 生产者
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("acks", "all"); // 等待所有副本确认
props.put("retries", 3); // 自动重试 3 次
props.put("linger.ms", 5); // 批量发送等待 5ms
props.put("batch.size", 16384); // 批量大小 16KB
Producer producer = new KafkaProducer<>(props);
ProducerRecord record =
new ProducerRecord<>("orders", "order-123", "{\"total\":99.9}");
producer.send(record).get(); // 同步发送
producer.close();
Java 消费者
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "order-service");
props.put("auto.offset.reset", "earliest"); // 无 offset 时从最早开始
props.put("enable.auto.commit", false); // 手动提交 offset
KafkaConsumer consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("orders"));
while (true) {
ConsumerRecords records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord record : records) {
processOrder(record.key(), record.value());
}
consumer.commitSync(); // 手动提交
}
分区策略
Producer 端决定消息去哪个分区,有三种方式。默认行为是:有 key 走哈希、无 key 走轮询——多数场景不需要自定义:
| 策略 | 行为 | 适用场景 |
|---|---|---|
| key 哈希 | 对 key 做 murmur2 哈希后取模分区数,相同 key 永远进同一分区 | 同一订单/用户的消息需要严格有序、按 key 聚合 |
| 轮询(无 key) | 消息依次均匀分配到各分区 | 无顺序要求、追求最大吞吐的通用日志 |
| 自定义分区器 | 实现 Partitioner 接口,按业务规则分配 | 按地域/租户隔离、热点 key 打散 |
| 粘性分区(2.4+) | 同一批次消息先攒满一个分区再换下一个 | 小消息高频场景,减少请求数并改善批次效率 |
// 自定义分区器:VIP 用户(key 前缀 v-)固定进 partition 0,其余按哈希
public class VipPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
int numPartitions = partitions.size();
if (keyBytes == null) {
// 无 key:轮询(用线程安全的原子计数器)
return counter.getAndIncrement() % numPartitions;
}
if (key.toString().startsWith("v-")) {
return 0; // VIP 流量单独进 0 号分区
}
return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
}
@Override public void configure(Map<String, ?> configs) {}
@Override public void close() {}
}
// 使用:props.put("partitioner.class", "com.example.VipPartitioner");
自定义分区器要保证 partition() 是纯函数(无副作用、可重复调用),因为 Producer 端批量发送时会为每条消息调用它,且发送失败重试时分区必须保持不变,否则消息会乱序。
4. 消息语义与可靠性
| 语义 | acks 配置 | 重试行为 | 适用 |
|---|---|---|---|
| At-Most-Once | acks=0 | 不重试 | 允许丢失(日志采集) |
| At-Least-Once | acks=all | 重试至成功 | 大多数业务场景 |
| Exactly-Once | 幂等 + 事务 | 幂等生产者 | 金融、订单 |
// 幂等生产者配置(Kafka 0.11+)
props.put("enable.idempotence", true); // 启用幂等,自动 acks=all
// 幂等 + 事务(跨 Topic 原子写入)
props.put("transactional.id", "order-transaction");
producer.initTransactions();
producer.beginTransaction();
producer.send(record1);
producer.send(record2);
producer.commitTransaction();
Exactly-Once 的实现细节
幂等生产者(0.11+)的原理:每个 Producer 实例初始化时分配 producerId(PID),每发一条消息分配自增 sequence 序号,Broker 端按 (PID, 分区) 校验序号——序号不连续说明有重复或乱序,直接拒绝。它只能保证单会话内不重复,跨会话或跨 Topic 的原子性要靠事务:
// 完整事务示例:订单创建 + 扣库存 + 发消息,要么全成要么全回滚
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("transactional.id", "order-txn-001"); // 必须固定且唯一
props.put("enable.idempotence", "true"); // 事务要求幂等
props.put("acks", "all");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions(); // 注册事务,获取 PID
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("orders", "order-1", "{\"total\":99.9}"));
producer.send(new ProducerRecord<>("inventory", "sku-88", "-1"));
producer.sendOffsetsToTransaction(
consumerOffsets, "order-service"); // 可选:消费位点也纳入事务
producer.commitTransaction();
} catch (ProducerFencedException e) {
// 事务 ID 冲突(旧实例未关闭),需重建 Producer
producer.close();
} catch (KafkaException e) {
producer.abortTransaction(); // 业务失败则回滚
}
- transactional.id 的作用:Broker 用它绑定 PID。同一 transactional.id 的新 Producer 启动时会"封死"旧 Producer(抛出 ProducerFencedException),防止僵尸实例写入。
- 事务状态存储:事务日志写入内部 Topic
__transaction_state(默认 50 分区 3 副本),transaction.state.log.replication.factor必须 ≥ 副本数。 - 代价:开启事务会多写一次事务日志,吞吐下降约 10-20%,且
read_committed消费者会过滤未提交消息,端到端延迟略增。能用幂等解决就别上事务。
ISR 与 HW/LEO 机制
副本同步是 Kafka 可靠性的核心,三个术语必须分清:
| 术语 | 全称 | 含义 |
|---|---|---|
| LEO | Log End Offset | 每个副本日志的下一条写入位置(最后一条消息的 offset + 1) |
| HW | High Watermark | ISR 中所有副本都复制到的最低 LEO;消费者只能读到 HW 之前的消息 |
| ISR | In-Sync Replicas | 与 Leader 保持同步的副本集合,由 replica.lag.time.max.ms(默认 30s)判定 |
# 写入与确认流程(replication.factor=3,acks=all)
# 1. Producer 发消息到 Leader,LEO: 100 -> 101
# 2. Follower 1 拉取复制,LEO: 100 -> 101
# 3. Follower 2 拉取复制,LEO: 100 -> 101
# 4. 三个副本都追上后,HW 从 100 -> 101,向 Producer 返回成功
# 5. 此时消费者才能读到 offset=100 这条消息
# 故障场景:Follower 2 宕机 30s 未同步,被踢出 ISR
# 此时 acks=all 仍能成功(ISR=[Leader, Follower1]),但只剩 2 副本冗余
# 若 min.insync.replicas=2,ISR 降为 1 个时生产者直接报 NotEnoughReplicasException
min.insync.replicas=2(配合 replication.factor=3、acks=all),保证至少 2 个副本确认才返回成功,避免"Leader 单独确认但数据未复制"的丢消息窗口。消费一致性:HW 机制保证消费者永远读不到"可能丢失"的消息——即使 Leader 宕机,从 ISR 中选出新 Leader 后,未达 HW 的消息会被截断,消费者不会读到前后不一致的数据。
5. 消费者组与分区再平衡
消费者组内每个消费者负责不同分区。当消费者加入或离开时触发再平衡(Rebalance),期间消息消费暂停。
# 消费者组状态查看
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-service --describe
# 输出:
# GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# order-service orders 0 1500 1520 20
# order-service orders 1 1480 1480 0
LAG(堆积量)是监控关键指标——LAG 持续增长说明消费能力不足,需要增加消费者实例或优化处理速度。
再平衡(Rebalance)的触发与影响
再平衡是"组内分区所有权重新分配"的过程,任何成员变化都会触发:
- 消费者加入:新实例启动并调用
subscribe(),触发一次分配 - 消费者退出:主动
close()或优雅退出,触发重新分配 - 消费者崩溃:心跳超时(
session.timeout.ms内未发心跳)被判定下线,触发重新分配 - 消费卡死:处理消息耗时超过
max.poll.interval.ms(默认 5 分钟),即使心跳正常也被踢出组 - 订阅变化:调用
subscribe()更换 Topic 列表
影响:再均衡期间全组消费暂停(stop-the-world),且旧分配被收回、新分配未生效前,所有消费者都要重新拉取分区并重置本地状态。频繁再均衡 = 消费抖动 + LAG 瞬时上涨。
| 分配策略 | 算法 | 特点 |
|---|---|---|
| Range(默认) | 按 Topic 逐一分区,连续片段分配 | 分区数不均时易倾斜,一个消费者可能多拿整个 Topic 的分区段 |
| RoundRobin | 所有 Topic 分区混排轮流分配 | 更均匀,但再均衡时全组停顿 |
| Sticky | 尽量保留上次分配,只移动必要分区 | 减少移动量,但同样全组停顿 |
| Cooperative Sticky(2.4+) | 分多轮渐进式再均衡 | 仅受影响消费者暂停,推荐生产使用 |
降低再均衡影响的措施:调大 session.timeout.ms 与 max.poll.interval.ms 容忍偶发慢消费;设置 partition.assignment.strategy=CooperativeStickyAssignor;消费逻辑幂等化,保证再均衡后的重复消费无副作用。
消费者参数与 offset 管理
# 生产消费者配置清单
props.put("max.poll.records", 500); // 单次 poll 返回上限,控制处理时长
props.put("max.poll.interval.ms", 300000); // 两次 poll 最大间隔(默认 5 分钟)
props.put("session.timeout.ms", 10000); // 心跳超时
props.put("heartbeat.interval.ms", 3000); // 心跳频率,建议 session.timeout / 3
props.put("fetch.min.bytes", 1024); // 攒够 1KB 再返回,减少请求次数
props.put("fetch.max.wait.ms", 500); // 最多等 500ms,控制延迟上限
# 手动重置 offset(排查问题后重放消息)
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-service \
--topic orders:0 --reset-offsets --to-earliest --execute
# 整个组重置到 3 天前
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-service --topic orders \
--reset-offsets --to-datetime 2025-06-12T00:00:00.000 --execute
max.poll.records 与 max.poll.interval.ms 的联动:如果单条消息处理耗时 1 秒,一次 poll 500 条就需要 500 秒——超过 max.poll.interval.ms 就会被判定"卡死"踢出组,触发一次无谓的再均衡。要么调小 max.poll.records,要么调大 max.poll.interval.ms,二者必须匹配。
消费堆积(LAG 增长)排查流程
LAG 增长先分清是"生产变快"还是"消费变慢",再逐层定位。按下面的顺序排查,不要上来就加消费者:
# 第一步:确认范围——哪个 Topic/分区的 LAG 在涨
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group order-service --describe
# 所有分区均匀上涨 → 消费端整体变慢;单分区上涨 → 分区热点或分配不均
# 第二步:消费端是不是在处理?
# 进程在消费但慢 → 业务处理耗时增加(查应用慢日志);
# 不在消费(卡死/退出)→ 检查是否触发 max.poll.interval.ms 被踢出组
# 第三步:分区分配是否均衡(Range 策略易倾斜)
# 对比各消费者实例的分区数,明显不均则换 CooperativeStickyAssignor
# 第四步:生产端是否突增?(促销、数据回放)
# 看 BytesInPerSec 是否放大——是则先扩容消费实例,不是则回到第二步
经验阈值:LAG 允许的"业务容忍量"取决于消息时效性——日志同步可容忍小时级,订单通知只能容忍分钟级。告警按 LAG 分级:LAG > 1 万告警、> 10 万拉群、持续 30 分钟触发扩容流程。排查结论要落到文档:每次堆积事故记录时间线、根因、处置动作,沉淀为"堆积处置手册"。
6. Kafka vs RabbitMQ vs Redis Stream
| Kafka | RabbitMQ | Redis Stream | |
|---|---|---|---|
| 消息模型 | Pull(消费者拉取) | Push(Broker 推送) | Pull |
| 吞吐量 | 极高(100K-1M msg/s) | 中等(10K-50K msg/s) | 高(100K+ msg/s) |
| 消息持久化 | 磁盘(长期保留) | 内存 + 持久化队列 | 内存 + AOF |
| 消息回放 | 支持(按 offset 重读) | 不支持(消费即删) | 支持(按 ID) |
| 消息顺序 | 分区内严格有序 | 队列内有序 | Stream 内有序 |
| 事务 | 支持(Exactly-Once) | 不支持 | 不支持 |
| 延迟消息 | 需额外插件 | 原生支持 TTL | 不支持 |
7. 生产调优
| 参数 | 默认值 | 推荐值 | 说明 |
|---|---|---|---|
num.network.threads | 3 | 8 | 网络线程数 |
num.io.threads | 8 | 16 | 磁盘 I/O 线程数 |
log.retention.hours | 168(7天) | 按需 | 消息保留时间 |
log.segment.bytes | 1GB | 512MB | 段文件大小 |
batch.size | 16384 | 32768 | 生产者批量大小 |
linger.ms | 0 | 5-10 | 批量发送等待时间 |
compression.type | none | lz4 | 生产者压缩算法,CPU 换带宽 |
replication.factor | 1 | 3 | 副本数,broker 默认 1,生产建议 3 |
消息压缩与协议优化
压缩是 Kafka 吞吐优化的杠杆:用 CPU 换带宽与磁盘。Producer 端开启后,Broker 原样存储压缩消息,消费者端自动解压,端到端透明。
| 压缩算法 | 特点 | 推荐场景 |
|---|---|---|
| lz4 | 压缩解压快,压缩率中等 | 吞吐优先的默认选择 |
| zstd | 压缩率高,CPU 开销略大 | JSON/文本消息多、带宽紧张的场景 |
| gzip | 压缩率最高,CPU 开销最大 | 批量大、对延迟不敏感的离线同步 |
props.put("compression.type", "lz4"); // Java 生产者一行开启
props.put("linger.ms", 20); // 压缩配合攒批,压缩率更高
# 注意事项:压缩按批次生效——linger.ms=0 或 batch.size 过小时,
# 每个批次只有一两条消息,压缩几乎无收益,反而浪费 CPU
吞吐量验证:producer-perf-test 基准测试
# 官方自带性能工具:验证集群吞吐与延迟
bin/kafka-producer-perf-test.sh \
--topic perf-test --num-records 1000000 \
--record-size 1024 --throughput -1 \
--producer-props bootstrap.servers=localhost:9092 acks=1
# 输出:
# 1000000 records sent, 65482.3 records/sec (63.9 MB/sec), 8.7 ms avg latency
# 解读:单生产者 6.5 万条/秒、64MB/s,平均延迟 8.7ms
# 把 acks 换成 all 再跑一次对比,可量化可靠性带来的吞吐代价
# 消费端基准
bin/kafka-consumer-perf-test.sh \
--bootstrap-server localhost:9092 --topic perf-test \
--messages 1000000
# 输出:data.consumed.in.MB、MB.sec、data.consumed.in.nMsg、nMsg.sec
监控 Kafka 的关键 JMX 指标
| 指标 | 含义 | 告警阈值建议 |
|---|---|---|
| BrokerTopicMetrics - BytesInPerSec | 集群写入吞吐 | 持续高于 80% 磁盘带宽上限 |
| BrokerTopicMetrics - FailedProduceRequestsPerSec | 写入失败请求数 | 大于 0 即告警 |
| RequestMetrics - RequestQueueTimeMs | 请求在队列等待时间 | P99 大于 200ms 说明 IO 线程不足 |
| KafkaRequestHandlerPool - RequestHandlerAvgIdlePercent | 处理线程空闲率 | 低于 30% 说明线程池打满 |
Kafka 通过 JMX 暴露全部指标,配合 3.12:系统监控与告警 的 Prometheus + jmx_exporter 采集,再加上消费组 LAG 告警,就构成完整的 Kafka 监控面。
分区数与副本数的权衡
- 分区数 = 并行度上限:读写并行度、消费者数、Producer 端并发都受分区数限制;单分区吞吐上限约 20-50MB/s,需要更高吞吐就加分区。
- 分区不是越多越好:每个分区约占 1-2 个文件句柄与内存缓冲,Broker 端元数据变更(Topic 增删、分区迁移)要广播到所有 Broker,分区过多时管理开销可观。
- 副本数与可靠性:3 副本可容忍 1 个 Broker 故障,5 副本可容忍 2 个;每多一个副本多一份磁盘与网络复制开销,写入延迟略增。
# 经验公式(3 Broker 集群、日志型场景)
# 分区总数 ≤ Broker 数 × 2000(Kafka 官方建议单 Broker 分区数不超过 4000)
# 吞吐需求 200MB/s ÷ 单分区 30MB/s ≈ 7 个分区,取 12(留余量且为 3 的倍数)
Broker 端生产环境检查清单
Kafka 是"磁盘 IO 敏感"的系统,Broker 所在宿主机的内核参数与磁盘布局直接影响稳定性。上线前逐项确认:
# /etc/security/limits.conf —— 文件句柄(每分区约占用 1-2 个句柄)
kafka soft nofile 100000
kafka hard nofile 100000
kafka soft nproc 32768
kafka hard nproc 32768
# /etc/sysctl.conf —— 网络与内存
net.core.somaxconn = 4096 # accept 队列,防连接溢出
net.ipv4.tcp_max_syn_backlog = 4096
vm.swappiness = 1 # 尽量不换页(页缓存是 Kafka 读路径的关键)
# server.properties —— 磁盘相关
log.dirs=/data1/kafka,/data2/kafka # 多目录分散 IO(各挂一块独立磁盘)
num.recovery.threads.per.data.dir=1 # 崩溃恢复时的并行度
# 磁盘选型:多块独立磁盘优于 RAID5(随机写命中率低);
# 顺序追加写场景下单块 NVMe 通常足够,容量按保留策略估算
# 内存分配:JVM 堆 4-8G 足够,剩余内存全部留给页缓存,
# 不要给堆分配超过 8G——Kafka 大量读写走页缓存而非堆
常见错误
| 问题 | 现象 | 排查与解决 |
|---|---|---|
| NotLeaderForPartitionException | 生产者报 Leader 不可用 | 检查目标 Broker 是否存活,运行 kafka-topics.sh --describe 查看 Leader 分配情况 |
| OffsetOutOfRangeException | 消费者启动时报 offset 越界 | 设置 auto.offset.reset=earliest 或 latest,确认 offset 未被 log.retention 清除 |
| 消息消费重复 | 同一条消息被处理多次 | 检查 enable.auto.commit 是否为 false 并手动提交 offset;确认 enable.idempotence=true |
| Producer 发送超时 | send() 长时间阻塞或 TimeoutException | 检查 batch.size 和 linger.ms 配置,确认网络连通性及 Broker 负载 |
| 消费堆积(LAG 持续增长) | LAG 指标不断增大 | 增加消费者实例(不超过分区数)、调大 max.poll.records、优化消息处理速度 |
| Rebalance 风暴 | 日志反复出现 rebalance,消费频繁中断 | 检查单次 poll 处理时长是否超过 max.poll.interval.ms,升级为 CooperativeSticky 分配策略 |
| 磁盘被消息撑满 | Broker 磁盘使用率飙升至 100%,写入失败 | 检查 log.retention 与 Topic 的 cleanup.policy;磁盘使用率 > 85% 挂告警,日志型 Topic 及时清理 |
| 消息格式不兼容 | 升级后消费者解析失败或丢字段 | 消息格式用 Schema Registry(Avro/Protobuf)管理,字段变更走 schema 演进,避免线上解析崩溃 |
| 单分区吞吐成瓶颈 | Topic 吞吐上不去,但 Broker 整体资源充足 | 单分区吞吐上限约 20-50MB/s;按峰值吞吐重算分区数(只能增加),把热点 key 用后缀打散 |
最佳实践
| 建议 | 说明 |
|---|---|
| 生产者 acks=all + idempotence | 启用幂等生产者,自动 acks=all,保证 At-Least-Once 语义且无重复写入 |
| 消费者手动提交 offset | 关闭 enable.auto.commit,在消息处理成功后调用 commitSync(),避免处理失败仍提交 offset |
| 合理规划分区数 | 分区数决定最大并行度,建议分区数 = 消费者数 = Broker 数的整数倍,避免消费者空闲 |
| 设置合理的保留策略 | 根据业务需求配置 log.retention.hours 和 log.retention.bytes,避免磁盘占满 |
| 统一序列化格式 | 使用 Avro/Protobuf 等 Schema 管理消息格式,配合 Schema Registry 做版本兼容 |
| Topic 命名与保留策略规范化 | 命名 <域>-<业务>-<事件>;日志型 delete 保留、事件型按业务周期、状态型 compact,变更走单 |
| LAG 分级告警 | LAG > 1 万告警、> 10 万拉群、持续 30 分钟触发扩容流程;先排查再扩容 |
| Broker 部署规范化 | 多块独立磁盘 + log.dirs 分散 IO、JVM 堆 4-8G、sysctl 与 limits 按检查清单逐项验证 |
练习题
- 在本地搭建 3 节点 KRaft 集群,创建 6 分区 3 副本的 Topic
test-orders,用命令行生产者发送 100 条消息,然后用消费者组消费并观察 LAG 变化。 - 编写一个 Java 生产者,分别使用
acks=0、acks=all和acks=all + enable.idempotence=true三种配置,记录每种配置下发送 1000 条消息的耗时和成功数,对比分析可靠性与性能的权衡。 - 模拟消费者组再平衡:启动 2 个消费者消费
test-orders,然后关闭其中一个,观察另一个消费者的分区分配变化以及 LAG 恢复过程。 - 用
kafka-producer-perf-test.sh对同一 Topic 分别测试acks=1与acks=all的吞吐和延迟,量化可靠性带来的吞吐代价。 - 模拟一次消费堆积事故:让消费者故意 sleep 3 秒处理每条消息,观察 LAG 增长与 max.poll.interval.ms 触发踢出组的日志,按第 5 节流程排查并给出处置方案。
- 为你的业务设计 Topic 规划表:命名、分区数、副本数、保留策略(时间/容量/compact),并说明每类 Topic 的 LAG 告警阈值。
- 对一台 Broker 做上线前体检:对照"Broker 端生产环境检查清单"逐项验证 limits、sysctl、log.dirs 与 JVM 堆配置,输出检查结果表。
学习检查点
学完本章后,请检验自己是否掌握以下内容:
| 检查项 | 自测问题 | 验证方法 |
|---|---|---|
| 概念理解 | 能用自己的话解释 Kafka 的分区、副本、消费者组概念 | 尝试向他人讲解 |
| 命令操作 | 能不查文档完成 Kafka 集群部署、Topic 创建和消息收发 | 在终端实际执行 |
| 原理掌握 | 能说出 Kafka 的零拷贝和顺序写入的性能优化原理 | 画出流程图 |
| 故障排查 | 能独立排查 Kafka 消费者组延迟或 Broker 宕机的问题 | 模拟故障并修复 |
| 最佳实践 | 能说明为什么需要合理设置 Kafka 的分区数和副本因子 | 对比不同方案 |
本章总结
Kafka 的高吞吐来自三个设计决策:Partition 并行、顺序写磁盘、消费者组水平扩展。它适合"大量事件流 + 重放 + 多消费者"的场景,但代价是运维复杂度与消息延迟。选型时记住:需要可靠投递和灵活路由选 RabbitMQ,需要百万级吞吐和流处理选 Kafka,轻量缓存级队列选 Redis Stream——没有万能的消息队列,只有匹配场景的架构。
延伸阅读
- 6.4:消息队列入门 消息队列入门——RabbitMQ 与 Redis Stream
- 4.7:压力测试实战 压力测试实战——用 wrk/hey 测试消息系统吞吐
- 3.12:系统监控与告警 系统监控与告警(Kafka JMX 指标 + Prometheus)