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 之间做出正确选型
  • 完成基础生产调优:批量参数、分区数与副本数的权衡

前置知识

架构概览

组件角色说明
Broker消息代理存储消息的服务器节点,集群由多个 Broker 组成
Topic消息主题逻辑分类,每个 Topic 可有多个 Partition
Partition分区并行处理单元,消息按 key 哈希分配到分区
Replica副本每个 Partition 的多份拷贝,分布在不同 Broker
Producer生产者发送消息的客户端
Consumer Group消费者组组内各消费者消费不同分区,实现负载均衡
ZooKeeper / KRaft元数据管理KRaft(Kafka 3.3+)替代 ZooKeeper 做选举
KRaft 模式 Kafka 3.3 起支持 KRaft(Kafka Raft),无需 ZooKeeper。新部署推荐直接使用 KRaft,减少运维复杂度。

数据流全景

一条消息从生产到消费的完整路径如下(文字版架构图):

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.sizelinger.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 端批量发送时会为每条消息调用它,且发送失败重试时分区必须保持不变,否则消息会乱序。

热点分区 key 哈希可能造成数据倾斜:一个超大 key(如热点商品)把流量全部打到单个分区,单分区吞吐上限约 20-50MB/s。解决:key 加随机后缀打散,消费端按后缀聚合;或用自定义分区器把热点 key 轮询到多个分区。

4. 消息语义与可靠性

语义acks 配置重试行为适用
At-Most-Onceacks=0不重试允许丢失(日志采集)
At-Least-Onceacks=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 可靠性的核心,三个术语必须分清:

术语全称含义
LEOLog End Offset每个副本日志的下一条写入位置(最后一条消息的 offset + 1)
HWHigh WatermarkISR 中所有副本都复制到的最低 LEO;消费者只能读到 HW 之前的消息
ISRIn-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 生产环境设置 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 持续增长说明消费能力不足,需要增加消费者实例或优化处理速度。

分区再平衡 消费者数量不要超过分区数——多余的消费者会空闲。建议分区数 = 消费者数 = Broker 数 的整数倍。

再平衡(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.msmax.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

KafkaRabbitMQRedis 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.threads38网络线程数
num.io.threads816磁盘 I/O 线程数
log.retention.hours168(7天)按需消息保留时间
log.segment.bytes1GB512MB段文件大小
batch.size1638432768生产者批量大小
linger.ms05-10批量发送等待时间
compression.typenonelz4生产者压缩算法,CPU 换带宽
replication.factor13副本数,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=earliestlatest,确认 offset 未被 log.retention 清除
消息消费重复同一条消息被处理多次检查 enable.auto.commit 是否为 false 并手动提交 offset;确认 enable.idempotence=true
Producer 发送超时send() 长时间阻塞或 TimeoutException检查 batch.sizelinger.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.hourslog.retention.bytes,避免磁盘占满
统一序列化格式使用 Avro/Protobuf 等 Schema 管理消息格式,配合 Schema Registry 做版本兼容
Topic 命名与保留策略规范化命名 <域>-<业务>-<事件>;日志型 delete 保留、事件型按业务周期、状态型 compact,变更走单
LAG 分级告警LAG > 1 万告警、> 10 万拉群、持续 30 分钟触发扩容流程;先排查再扩容
Broker 部署规范化多块独立磁盘 + log.dirs 分散 IO、JVM 堆 4-8G、sysctl 与 limits 按检查清单逐项验证

练习题

  1. 在本地搭建 3 节点 KRaft 集群,创建 6 分区 3 副本的 Topic test-orders,用命令行生产者发送 100 条消息,然后用消费者组消费并观察 LAG 变化。
  2. 编写一个 Java 生产者,分别使用 acks=0acks=allacks=all + enable.idempotence=true 三种配置,记录每种配置下发送 1000 条消息的耗时和成功数,对比分析可靠性与性能的权衡。
  3. 模拟消费者组再平衡:启动 2 个消费者消费 test-orders,然后关闭其中一个,观察另一个消费者的分区分配变化以及 LAG 恢复过程。
  4. kafka-producer-perf-test.sh 对同一 Topic 分别测试 acks=1acks=all 的吞吐和延迟,量化可靠性带来的吞吐代价。
  5. 模拟一次消费堆积事故:让消费者故意 sleep 3 秒处理每条消息,观察 LAG 增长与 max.poll.interval.ms 触发踢出组的日志,按第 5 节流程排查并给出处置方案。
  6. 为你的业务设计 Topic 规划表:命名、分区数、副本数、保留策略(时间/容量/compact),并说明每类 Topic 的 LAG 告警阈值。
  7. 对一台 Broker 做上线前体检:对照"Broker 端生产环境检查清单"逐项验证 limits、sysctl、log.dirs 与 JVM 堆配置,输出检查结果表。

学习检查点

学完本章后,请检验自己是否掌握以下内容:

检查项自测问题验证方法
概念理解能用自己的话解释 Kafka 的分区、副本、消费者组概念尝试向他人讲解
命令操作能不查文档完成 Kafka 集群部署、Topic 创建和消息收发在终端实际执行
原理掌握能说出 Kafka 的零拷贝和顺序写入的性能优化原理画出流程图
故障排查能独立排查 Kafka 消费者组延迟或 Broker 宕机的问题模拟故障并修复
最佳实践能说明为什么需要合理设置 Kafka 的分区数和副本因子对比不同方案

本章总结

Kafka 的高吞吐来自三个设计决策:Partition 并行、顺序写磁盘、消费者组水平扩展。它适合"大量事件流 + 重放 + 多消费者"的场景,但代价是运维复杂度与消息延迟。选型时记住:需要可靠投递和灵活路由选 RabbitMQ,需要百万级吞吐和流处理选 Kafka,轻量缓存级队列选 Redis Stream——没有万能的消息队列,只有匹配场景的架构。

延伸阅读

↑ 回到顶部