Kafka KRaft 部署、可靠生产与消费组治理工具手册
从“端口通了但消费者不动”开始
应用启动时成功连接 bootstrap.servers,健康检查也是绿色,生产者却持续超时;另一个服务能拉到元数据,却消费不到任何记录。运维看到 broker 进程仍在,开发看到 9092 也能建立 TCP 连接,双方都认为 Kafka 没问题。
Kafka 的客户端连接不是一次静态 TCP 会话。客户端先通过 bootstrap 地址取得集群元数据,再直接连接目标 partition leader;消费者还要加入 group、获得 partition assignment、从已提交 offset 继续 fetch。只要 broker 返回了容器内 hostname、目标 leader 不在 ISR、group 正在 rebalance、旧 offset 已经越过记录,都会出现“能连上但不能工作”。
排障时需要同时回答四个问题:
KRaft metadata quorum 是否有 active controller,元数据是否能推进。目标 partition 的 leader、replicas 和 ISR 是否符合预期。Producer 的 record 是否获得预期 acks,失败结果是否确定。
Consumer group 当前分配、committed offset、log end offset 和 lag 分别是什么。
Kafka 更接近可分区、可保留、可回放的提交日志,而不是带自动死信和逐条删除的传统队列。消息被消费者读取后不会立即从 broker 删除;保留策略独立决定数据何时清理,consumer group 的 offset 只记录消费位置。理解这一点,才能正确解释重复、积压、回放和容量。
Topic、partition、replica 与 offset 如何协作
Topic 是逻辑事件流,partition 是实际的有序追加日志。每个 partition 在任一时刻有一个 leader,producer 与普通 consumer 都通过 leader 读写;其他 replica 从 leader 拉取日志。ISR 是当前与 leader 保持同步、具备参与安全选主条件的副本集合。
Offset 是 partition 内的单调位置,不是全局消息编号。Kafka 只保证同一 partition 内的记录顺序;不同 partition 的记录没有统一先后关系。Producer 若提供 key,默认分区器会让相同 key 稳定落到同一 partition,从而建立“单个订单有序”一类局部顺序。无 key 的流量更均衡,却无法依赖业务键顺序。
Consumer group 把一个 topic 的 partition 分给组内成员。一个 partition 在同一 group 内同一时刻只分给一个成员,因此有效消费并行度上限首先受 partition 数限制。六个消费者读取只有三个 partition 的 topic,最多三个消费者有工作;三个消费者读取十二个 partition,每个成员会处理多个 partition。
Committed offset 表示这个 group 下一次恢复时从哪里继续,不表示业务副作用一定成功。消费者先提交 offset、后写数据库会丢处理;先写数据库、后提交 offset,在两步之间崩溃会重复处理。通用解法是业务幂等键、数据库唯一约束或 inbox;Kafka transaction 能原子提交 Kafka 输出记录与消费 offset,但不能自动把外部数据库事务纳入同一个原子边界。
KRaft 把集群元数据从 ZooKeeper 收回 Kafka
Kafka 4.x 只使用 KRaft。Controller 维护 broker、topic、partition、ACL 和配额等集群元数据,并通过 metadata quorum 复制;broker 保存业务日志并处理客户端请求。每个 server 的 process.roles 可以是 controller、broker 或开发环境中的 broker,controller。
Combined mode 把 controller 与 broker 放在同一进程,部署简单,适合本地实验和低风险环境;关键生产环境应分离角色,让 controller 不受业务 I/O、页缓存和滚动扩容干扰。Kafka 的KRaft 指南建议使用三个或五个 controller:三个可以容忍一个 controller 故障,五个可以容忍两个,元数据可用性始终需要多数派。
KRaft 元数据有自己的 log、high watermark、leader epoch 和快照。Broker 进程存活不代表 controller quorum 健康;检查命令是:
/opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server localhost:9092 \
describe --status关键证据包括 LeaderId、HighWatermark、CurrentVoters 和 follower lag。没有 leader 或多数 controller 不可达时,创建 topic、分区选主和 ACL 变更等控制面操作不能正常推进。
Kafka 4.3 支持动态 controller quorum。新集群可以用 kafka-storage.sh format --standalone 启动首个 voter,再用 add-controller 扩展;旧式静态 quorum 使用 controller.quorum.voters 固定成员。两种方式的格式化元数据不同,不能只修改一行配置来切换。开发 Compose 为了可读性可使用单节点静态 voter,生产自动化应明确选择动态还是静态,并保存 cluster id、directory id 和成员变更记录。
部署形态对应不同的恢复能力
官方二进制包最适合观察 kafka-storage.sh format、配置文件和日志目录,但需要团队管理 Java、服务进程、磁盘和权限。官方 Docker 镜像适合一次性本地验证;Compose 适合把 listener、数据卷、topic 初始化和清理边界固化进项目。
生产自建集群通常把三个 controller 与多个 broker 分离,broker 跨故障域放置,topic 使用多副本,并让客户端连接多个 bootstrap 地址。托管 Kafka 可以代管 broker、磁盘替换和部分升级工作,却不会替应用决定 key、partition、acks、offset、幂等、retention、ACL 和成本上限。
| 形态 | 组成与原理 | 适用信号 | 主要代价与故障模式 |
|---|---|---|---|
| 单节点 combined KRaft | 一个进程同时是 broker/controller,副本因子为 1 | 开发、契约测试、CLI 学习 | 任一进程或磁盘故障即不可用,不能验证 ISR 容错 |
| 分离式 KRaft 集群 | 3/5 controller + 多 broker,业务日志按 partition 复制 | 核心事件流、自建平台 | 磁盘、网络、再分配、升级和容量值守成本高 |
| 托管 Kafka | 供应商管理集群生命周期,应用使用服务端点 | 希望减少基础设施维护 | 按吞吐/存储/流量计费,配额与特性存在供应商边界 |
| 跨集群复制 | 独立集群间异步复制事件 | 灾备、地域隔离、迁移 | RPO 非零、offset/ACL/topic 配置需单独治理 |
生产高可用不是“启动三台 Kafka”。Topic 的 replication factor、min.insync.replicas、producer acks 和故障域分布必须配套;否则三台 broker 上的 topic 仍可能只有一个副本。
在开发机启动官方 Kafka 4.3.1 镜像
Apache Kafka 官方下载页列出了受支持的 4.3.1 发布和 apache/kafka:4.3.1 镜像。镜像必须固定补丁版本,升级前核对 release notes、客户端兼容和数据目录,不使用 latest。
创建 .env:
KAFKA_IMAGE=apache/kafka:4.3.1
KAFKA_HOST_PORT=9092
KAFKA_TOPIC=te.orders.events.v1
KAFKA_GROUP=te.orders.projector.v1创建 compose.yaml:
services:
kafka:
image: ${KAFKA_IMAGE:?set KAFKA_IMAGE}
container_name: te-kafka
hostname: kafka
restart: unless-stopped
ports:
- "127.0.0.1:${KAFKA_HOST_PORT:-9092}:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,HOST:PLAINTEXT
KAFKA_LISTENERS: CONTROLLER://:9093,INTERNAL://:29092,HOST://:9092
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,HOST://127.0.0.1:9092
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
KAFKA_NUM_PARTITIONS: 3
KAFKA_DEFAULT_REPLICATION_FACTOR: 1
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk
volumes:
- kafka-data:/var/lib/kafka/data
healthcheck:
test:
- CMD-SHELL
- /opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092 >/dev/null 2>&1
interval: 10s
timeout: 10s
retries: 12
start_period: 30s
volumes:
kafka-data:三个 listener 各有职责:CONTROLLER 只用于 controller quorum,INTERNAL 把 kafka:29092 告诉同一 Compose 网络里的客户端,HOST 把 127.0.0.1:9092 告诉宿主机客户端。listeners 是进程绑定地址,advertised.listeners 是 broker 返回给客户端的可达地址;两者混淆是容器 Kafka 最常见的故障源。
这份配置把内部 topic 的副本因子降为 1,只为单节点开发环境能启动 transaction 与 consumer group。它不能作为生产基线。生产环境通常至少使用三个 broker,并按可容忍故障数设置内部 topic 和业务 topic 的 replication factor 与 min ISR。
启动并检查:
docker compose up -d kafka
docker compose ps kafka
docker logs --tail 100 te-kafka
docker exec te-kafka /opt/kafka/bin/kafka-broker-api-versions.sh \
--bootstrap-server localhost:9092
docker exec te-kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server localhost:9092 describe --statusAPI versions 成功只证明能访问 bootstrap broker;metadata quorum 输出有 leader 才证明 KRaft 元数据可推进。宿主机应用还要实际 produce/fetch,才能证明 advertised HOST 地址可达。
创建一个能解释的 Topic
显式创建三个 partition、一个副本的开发 topic:
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic te.orders.events.v1 \
--partitions 3 \
--replication-factor 1 \
--config cleanup.policy=delete \
--config retention.ms=86400000 \
--config segment.ms=3600000retention.ms 的一天和 segment.ms 的一小时是开发演示值。Kafka 按 log segment 清理,达到 retention 不代表某条记录在那个毫秒立即物理删除;较短 segment 能让开发实验更快观察清理,但会增加文件滚动。生产值必须根据回放窗口、磁盘预算和下游最长恢复时间决定。
查看 topology:
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic te.orders.events.v1单节点预期每个 partition 都显示 Leader: 1、Replicas: 1、Isr: 1。这是结构证据,不是高可用证据。三 broker 生产 topic 若设置 replication factor 3,健康状态应看到三个 replica 都在 ISR;ISR 缩小时,应结合 replica lag、网络和磁盘判断,不能立刻执行 unclean leader election。
生产创建示例通常类似:
kafka-topics.sh --bootstrap-server kafka-a:9092,kafka-b:9092,kafka-c:9092 \
--create --topic orders.events.v1 \
--partitions 12 --replication-factor 3 \
--config min.insync.replicas=2 \
--config cleanup.policy=delete \
--config retention.ms=604800000分区数 12、保留七天都只是示意。分区只能增加、不能直接缩小;增加后 key 到 partition 的映射可能变化,并行度、文件句柄、controller 元数据、consumer rebalance 和小批量开销都会增加。先由目标吞吐、单 partition 基准、key 基数、消费者并行度和未来增长推导,而不是用固定公式拍数。
正向实验:生产、消费并观察 Offset
先生产三条带 key 的记录。Console producer 的 parse.key=true 用制表符分隔 key 与 value:
printf 'order-1001\t{"eventId":"evt-1","type":"OrderCreated"}\norder-1002\t{"eventId":"evt-2","type":"OrderCreated"}\norder-1001\t{"eventId":"evt-3","type":"OrderPaid"}\n' | \
docker exec -i te-kafka /opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic te.orders.events.v1 \
--property parse.key=true \
--property key.separator=$'\t' \
--producer-property acks=all \
--producer-property enable.idempotence=trueConsole producer 正常退出说明这批 send 没有以异常结束。要观察 partition、offset 和 key,使用独立 group 消费:
docker exec te-kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic te.orders.events.v1 \
--group te.orders.projector.v1 \
--from-beginning \
--timeout-ms 10000 \
--property print.key=true \
--property print.partition=true \
--property print.offset=true预期三条记录可见,order-1001 的两条记录位于同一 partition 且 offset 递增。Console consumer 退出后查看 group:
docker exec te-kafka /opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group te.orders.projector.v1每个 partition 会显示 CURRENT-OFFSET、LOG-END-OFFSET 和 LAG。消费完成时 lag 应回到 0;没有记录的 partition 也可能显示有效 offset。Group 没有 active member 时,consumer id 可为空,但已提交 offset 仍保留。
反向实验:证明旧 Group 不会自动重读
再次使用同一个 group 和 --from-beginning 启动 consumer:
docker exec te-kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic te.orders.events.v1 \
--group te.orders.projector.v1 \
--from-beginning \
--timeout-ms 5000预期不会重新打印前三条记录。--from-beginning 只在 group 对目标 partition 没有已提交 offset 时决定起点;已有 offset 的 group 仍从 committed position 继续。换一个新 group 才能从 earliest 读取:
docker exec te-kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic te.orders.events.v1 \
--group te.orders.audit-replay.v1 \
--from-beginning \
--timeout-ms 10000如果生产 group 需要回放,先停止全部成员,预览 reset 结果,再执行:
docker exec te-kafka /opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group te.orders.projector.v1 \
--topic te.orders.events.v1 \
--reset-offsets --to-earliest
docker exec te-kafka /opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group te.orders.projector.v1 \
--topic te.orders.events.v1 \
--reset-offsets --to-earliest --execute第一条不带 --execute,只输出计划。Offset reset 会造成重复消费或跳过记录,必须记录操作者、topic、partition、旧 offset、新 offset、业务幂等确认和回滚点。共享环境默认不给普通开发账号 Alter group 权限。
反向实验:让 Listener 错误留下明确证据
把 KAFKA_ADVERTISED_LISTENERS 中的 HOST 临时改为宿主机无法解析的地址:
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,HOST://kafka.invalid:9092重建容器后,宿主机客户端可以先连到 127.0.0.1:9092,随后在访问 partition leader 时失败,并在日志中出现无法解析 kafka.invalid 或连接该地址失败。容器内 CLI 通过 INTERNAL listener 仍可能正常,所以“容器内命令成功”不能证明宿主机接入正确。
恢复为 HOST://127.0.0.1:9092 并重建:
docker compose up -d --force-recreate kafka
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 --list复杂环境要分别从宿主机、Compose 网络、CI runner 和真实应用网络执行 metadata + produce + consume。Listener 名称不是 DNS;每个 listener 广播的 hostname 和端口必须对使用它的网络域可达。
Producer 的 acks、幂等与不确定结果
acks=0 不等待 broker 响应,客户端无法知道消息是否到达;acks=1 等 leader 写入,但 leader 在 follower 复制前永久故障仍可能丢失;acks=all 等待当前 ISR 满足 broker 的 min.insync.replicas。Kafka 的Topic 配置明确说明,当 ISR 少于 min.insync.replicas 且 producer 使用 acks=all 时,写入会以副本不足异常失败。
生产者可靠基线:
bootstrap.servers=kafka-a:9092,kafka-b:9092,kafka-c:9092
client.id=orders-api-producer
acks=all
enable.idempotence=true
retries=2147483647
delivery.timeout.ms=120000
request.timeout.ms=30000
max.in.flight.requests.per.connection=5
compression.type=zstd
linger.ms=5
batch.size=65536enable.idempotence=true 让 broker 用 producer id、epoch 和 sequence 去除同一 producer session 内因重试产生的重复,并要求兼容的 acks、retries 和 max in flight 配置。它不能跨 producer 身份永久去重,也不能阻止业务代码主动发送两次相同事件。最终仍要在消息中携带稳定 eventId,下游以该键幂等。
生产 API 成功返回的 metadata 应至少记录 topic、partition 和 offset。超时或连接中断属于不确定结果:record 可能没写入,也可能已写入但响应丢失。不能把所有 timeout 都记成“发送失败”后无脑生成新 eventId;应保留同一业务事件标识重试,让幂等 producer 与幂等 consumer共同收敛重复。
Producer 的批量与压缩是吞吐和延迟取舍。linger.ms 增大能收集更大 batch、提升压缩与吞吐,但增加低流量等待;batch.size 太小会增加请求开销,太大会增加内存;大 record 还受 producer、broker、topic 和 consumer 多处大小上限约束。容量测试应使用真实消息大小分布,而不是只测几十字节字符串。
Transaction 解决的是 Kafka 内部原子写入
设置唯一且稳定的 transactional.id 后,producer 可以原子写多个 partition,并将消费 offset 与输出记录一起提交。消费者使用 isolation.level=read_committed 时不会看到已 abort 的事务记录。典型 consume-transform-produce 流程是:
initTransactions
beginTransaction
poll input records
produce output records
sendOffsetsToTransaction
commitTransaction实例重启后,新的 producer epoch 会 fencing 旧实例,避免两个相同 transactional id 的 writer 同时提交。transactional.id 因此必须按逻辑实例分配并受 ACL 保护,不能让所有副本共享一个随意常量,也不能每次重试生成新 id。
Transaction 不会把 MySQL、HTTP 调用和对象存储写入纳入 Kafka 原子提交。需要“数据库状态 + 事件”一致时,常用 transactional outbox;需要“消费 + 数据库副作用”一致时,使用业务幂等表、唯一键和可重试事务。把 Kafka 的 exactly-once 直接解释为端到端业务只执行一次,会掩盖外部系统边界。
Consumer Group、Offset 与 Rebalance
消费者循环不是简单的 poll。成员加入 group 后经历发现 coordinator、Join/Sync 或新 consumer protocol 的分配、poll records、处理、提交 offset、heartbeat 与 rebalance。成员增加、退出、session timeout、订阅变化和 partition 增加都可能触发重新分配。
Kafka 4.3 同时支持 classic 与 consumer 两种 group protocol,客户端默认值仍是 classic;只有显式配置 group.protocol=consumer 才会启用新协议。新协议把心跳、会话超时与分配策略的控制更多移到 broker,并以增量分配减少全局停顿,但迁移前仍应检查客户端版本、server-side assignor 和回调行为。官方Consumer Rebalance Protocol给出了在线与停机切换路径。
消费者配置基线:
bootstrap.servers=kafka-a:9092,kafka-b:9092,kafka-c:9092
group.id=orders.projector.v1
client.id=orders-projector-${INSTANCE_ID}
enable.auto.commit=false
auto.offset.reset=earliest
isolation.level=read_committed
max.poll.records=100
max.poll.interval.ms=300000
group.protocol=consumergroup.protocol=consumer 下不要设置客户端 session.timeout.ms、heartbeat.interval.ms 或 partition.assignment.strategy:Kafka 4.3 不支持这些客户端配置,成员会话超时与心跳间隔分别由 broker 的 group.consumer.session.timeout.ms 和 group.consumer.heartbeat.interval.ms 控制,分配器由 broker 的 group.consumer.assignors 与可选客户端 group.remote.assignor 决定。group.consumer.session.timeout.ms 是只读 broker 配置,修改 server 配置后需要重启对应节点;先在每个 broker 的实际加载配置中确认取值,再用隔离 group 复制失联回收过程:
# 终端 A、B 各启动一个成员;consumer-secure.properties 放认证配置,不放 session.timeout.ms
kafka-console-consumer.sh --bootstrap-server kafka-a:9092 \
--consumer.config consumer-secure.properties \
--consumer-property group.protocol=consumer \
--topic orders.events.v1 --group orders.session-probe.v1
# 终端 C 记录初始成员数,然后 kill -9 终端 B 中的消费者进程并重复执行
kafka-consumer-groups.sh --bootstrap-server kafka-a:9092 \
--command-config consumer-secure.properties \
--describe --group orders.session-probe.v1 --members第一次输出应有两个成员;强杀后,成员数应在接近 group.consumer.session.timeout.ms 的时间窗口内降为一个,并由存活成员接管 partition。消费者配置审查也应确认没有把 session.timeout.ms、heartbeat.interval.ms 和客户端 assignor 带入 consumer protocol。若仍使用 group.protocol=classic,会话存活才由客户端 session.timeout.ms 与 heartbeat.interval.ms 配合控制,并受 broker 的 group.min.session.timeout.ms、group.max.session.timeout.ms 范围约束。两种协议的超时所有权不同,切换时必须保留成员退出、进程强杀、rebalance 时长和重复处理证据,不能只看消费者最终重新上线。
enable.auto.commit=false 只是关闭定时自动提交,应用仍必须选择提交时点。批量 poll 后逐条处理却一次提交整个批次,进程在中途失败可能跳过尚未处理的后半批;处理全部成功后提交则会在失败时重放已完成的前半批。解决方案是批次事务、逐 partition 追踪安全 offset、减小 batch 或确保每条业务操作幂等。
max.poll.interval.ms 限制两次 poll 的最大间隔。业务处理、GC 或下游超时超过该值时,成员会失去分配;旧消费者随后提交可能遇到 generation/member 错误。不要只把超时调大,应把长任务移出 poll 线程、限制 max.poll.records、设置下游超时并监控处理 P99。
Rebalance 期间要停止接收被撤销 partition 的新任务,等待在途处理到安全点,提交对应 offset,再释放 partition 级状态。回调阻塞过久会延长停顿;异步线程若没有 partition ownership fencing,旧成员可能在失去分配后继续写业务状态。
重复、顺序、重试与死信要由应用建模
Kafka 的常见交付语义是至少一次:producer 不确定结果会重试,consumer 在业务完成后、offset 提交前崩溃会重读。Event 必须有稳定 id、schema version、业务 key、发生顺序或实体 version;consumer 用唯一约束或 inbox 收敛重复。
顺序只在 partition 内成立。需要同一订单有序时,key 使用订单 id,并避免随意增加 partition 后未评估 key 重映射。重试如果把记录写到另一个 retry topic,跨 topic 后的全局顺序已改变;原 partition 内阻塞重试保序,却可能让一条毒消息阻塞整个分区。
Kafka broker 不会像 RabbitMQ 那样自动对普通 consumer 执行 nack/DLX。应用通常建立:
orders.events.v1
-> 可重试错误 -> orders.events.retry.1m.v1 -> 延迟后回主流
-> 永久错误 -> orders.events.dlq.v1DLQ 记录应包含原 topic、partition、offset、key、event id、schema version、失败分类、重试次数和 payload 引用。敏感 payload 不应直接复制到无保护的错误流。重放工具必须指定 offset 范围、速率、幂等策略和审计批次;“消费 DLQ 后重新 produce”不是无风险操作。
积压用 group lag 表达,而不是 topic 消息总数。单个慢 partition 会决定 group 恢复时间,即使总 lag 不大。扩 consumer 只有在仍有未并行的 partition 时有效;热点 key 造成单 partition 倾斜时,应先修复 key/partition 设计或拆分热点,盲目扩实例不会增加该 partition 的并发。
ISR、选主与故障恢复的判断标准
生产 topic 常用 replication factor 3、min.insync.replicas=2、producer acks=all。当一个 follower 落后并离开 ISR,剩余 leader 与一个 follower 仍可满足最小 ISR;再失去一个同步副本时写入会失败,以避免在副本不足时继续确认。这个失败是数据安全策略生效,不应通过临时改成 acks=1 掩盖。
查看状态:
kafka-topics.sh --bootstrap-server kafka-a:9092 \
--describe --topic orders.events.v1
kafka-metadata-quorum.sh --bootstrap-server kafka-a:9092 \
describe --status诊断关注 leader 是否存在、replicas 是否跨故障域、ISR 是否缩小、under-replicated partitions、offline partitions、replica fetch lag、broker 磁盘和网络。Controller quorum 健康与业务 partition ISR 是两条不同复制链,必须分别观察。
Unclean leader election 允许不在 ISR 的副本成为 leader,可能用可用性换取已确认数据丢失,默认不应作为“快速恢复”按钮。Leader 不可用时先恢复仍持有最新日志的 ISR 副本,确认故障域与磁盘证据,再按业务 RPO 决定是否接受不干净选主。
Broker 下线或磁盘迁移涉及 partition reassignment。再分配会消耗磁盘和网络,应设置 throttle、分批执行并用 kafka-reassign-partitions.sh --verify 证明完成;不能在高峰一次迁移全部 partition。扩 broker 也不会自动让已有 partition 均衡到新节点,必须有明确再分配计划。
Retention、Compaction 与“记录为什么消失”
cleanup.policy=delete 按时间或大小删除旧 log segment。retention.ms 是时间边界,retention.bytes 是每 partition 大小边界;任一条件触发都可能清理。删除是 segment 级异步过程,因此不能把 retention 当精确计时器,也不能保证某条记录保留到某个毫秒。
cleanup.policy=compact 按 key 保留最终状态,旧值可在 compaction 后消失;key 为 null 的记录不能提供有意义的状态合并。写入 key 对应的 tombstone(null value)表示删除,tombstone 又受 delete retention 控制。Compaction 不等于“每个 key 永远只剩一条”,清理过程中仍可能看到多个版本。
可以同时配置 delete,compact,让每个 key 的最新状态也受总保留窗口限制。选择标准是消费模型:事件审计流通常使用 delete 并明确回放窗口;状态 changelog 使用 compact;既要最新状态又不能无限增长时组合两者。不要为了省磁盘把业务审计流改成 compact,否则历史变化会被清除。
开发环境可执行反向实验:把独立实验 topic 的 retention.ms 和 segment.ms 调小,写入记录,等待 segment 滚动和清理线程运行,再比较 earliest 与 latest offset。清理时间并非精确,所以判断证据是 earliest offset 前移和旧记录无法 fetch,而不是只等待固定秒数。
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create --topic te.orders.retention.lab \
--partitions 1 --replication-factor 1
docker exec te-kafka /opt/kafka/bin/kafka-configs.sh \
--bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name te.orders.retention.lab \
--alter \
--add-config retention.ms=60000,segment.ms=10000实验只对专用 topic 执行。共享 topic 改 retention 可能不可逆删除数据,必须先确认 owner、磁盘压力、消费组最长 lag 和回放需求。
容量从吞吐、保留与恢复窗口推导
粗略原始存储需求:
原始数据 = 峰值每秒记录数 × 平均记录字节 × 保留秒数
集群占用 ≈ 原始数据 × 副本因子 ÷ 压缩比 × 索引与安全系数平均值会低估大消息和峰值,应使用 P95/P99 大小、峰值持续时间和压缩后实测。还要预留 segment、index、transaction index、复制赶超、再分配和磁盘故障恢复空间。磁盘到高水位后再扩容,副本迁移本身可能没有足够空间完成。
吞吐容量至少分别测 producer bytes/records、broker request latency、network、disk utilization、partition count、consumer fetch 和端到端延迟。单 partition 吞吐取决于 record 大小、batch、compression、acks、副本、磁盘和网络;不能用脱离硬件与配置的万能数字决定 partition 数。
积压恢复:
净追赶速率 = 消费处理速率 - 新到达速率
恢复时间 = group lag / 净追赶速率若净追赶速率小于等于零,消费者永远追不上。要么提高有效 partition 并行度和消费者处理能力,要么对生产者限流或降级。SLO 应同时约束最大记录年龄和 lag,而不是只看 lag 条数;一万条大事件与一万条小事件的恢复成本不同。
Partition 也是元数据、文件和调度成本。Topic 无 owner、环境前缀失控、每个测试动态建 topic,会逐步推高 controller 内存、broker 文件句柄和巡检负担。团队必须设置创建配额、命名规则、默认 partition 上限和过期回收流程。
Listener、安全与最小权限
开发机的 PLAINTEXT listener 只能绑定 loopback。共享环境至少使用网络隔离,并根据威胁模型选择 SASL_SSL 或 mTLS;TLS 保护传输,SASL 或证书建立 principal,ACL 决定 principal 能操作哪些资源。Kafka 的SASL 指南和ACL 指南应与所用认证机制一起落地。
KRaft 不会因为执行了 kafka-acls.sh 就自动启用授权。所有 broker、controller 以及合并角色节点都要配置把 ACL 存入 KRaft metadata log 的 StandardAuthorizer;修改后滚动重启并检查每个节点的实际加载配置:
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
super.users=User:kafka-broker;User:break-glass-admin
allow.everyone.if.no.acl.found=falsesuper.users 不是应用账号白名单,也不能把示例中的 User:kafka-broker 原样照搬。它必须替换成 inter-broker listener 实际认证得到的 broker principal;User 大小写敏感,多个值用分号分隔,名称必须与认证和 principal mapping 后的日志身份完全一致。KRaft 中,客户端先把 CreateTopics、DeleteTopics 等管理请求发给 broker,broker 再通过 controller listener 转发 Envelope:controller 先鉴权已认证的 broker principal,再鉴权 Envelope 中携带的客户端 principal。因此 broker principal 需要通过受控的 super.users 或等价 ACL 完成控制面通信,客户端仍必须拥有自身操作所需 ACL;自定义 principal.builder.class 还必须实现 KafkaPrincipalSerde,否则 controller 无法还原被转发身份。紧急管理员可以是 super user,普通应用绝不能进入该列表。误设 allow.everyone.if.no.acl.found=true 会让没有任何匹配 ACL 的资源默认放行。
一个 consumer 通常需要目标 topic 的 Read/Describe 和 group 的 Read;producer 需要 topic 的 Write/Describe,幂等生产还涉及 cluster 的 IdempotentWrite,transaction producer 还需要 transactional id 权限。不要给应用 User:* 或 topic * 的 All 权限。
示意 ACL:
kafka-acls.sh --bootstrap-server kafka-a:9092 \
--add --allow-principal User:orders-producer \
--producer --idempotent \
--topic orders.events.v1
kafka-acls.sh --bootstrap-server kafka-a:9092 \
--add --allow-principal User:orders-projector \
--consumer \
--topic orders.events.v1 \
--group orders.projector.v1授权验证必须同时包含允许和拒绝。下面两个 properties 文件分别认证为 orders-projector 和没有任何 ACL 的 acl-denied-probe,其他连接参数保持相同;先向目标 Topic 写入一条测试记录,再执行:
kafka-console-consumer.sh --bootstrap-server kafka-a:9092 \
--consumer.config orders-projector.properties \
--topic orders.events.v1 --group orders.projector.v1 \
--from-beginning --max-messages 1 --timeout-ms 10000
echo "allowed_exit=$?" # 预期 0
kafka-console-consumer.sh --bootstrap-server kafka-a:9092 \
--consumer.config acl-denied-probe.properties \
--topic orders.events.v1 --group orders.projector.denied-probe.v1 \
--from-beginning --max-messages 1 --timeout-ms 10000
echo "denied_exit=$?" # 预期非 0拒绝身份必须收到 TOPIC_AUTHORIZATION_FAILED / TopicAuthorizationException 或 GROUP_AUTHORIZATION_FAILED / GroupAuthorizationException,并以非零状态退出;授权日志也应能关联到该 principal 与资源。若反例仍能读取,先核对客户端最终认证出的 principal,再检查所有节点的 StandardAuthorizer、super.users、匹配该资源的 literal/prefixed/wildcard ACL 和 allow.everyone.if.no.acl.found。可用 kafka-acls.sh --list --topic orders.events.v1 --resource-pattern-type match 查看所有会影响该 Topic 的规则,不能把只列出 exact ACL 当成安全闭环。
JAAS 文件、SCRAM 密码、keystore/truststore 密码、私钥、OAuth client secret 和托管服务 API key 不进入仓库。客户端配置用 secret 引用,日志与截图隐藏 bootstrap 域名、principal、topic、group 和 payload。证书/密码轮换应支持新旧凭证重叠,滚动验证后撤销旧凭证。
Kafka payload 可能长期保留并被多个 group 重放。敏感数据治理要在生产前决定字段最小化、加密、删除权、retention、访问审计和 schema 兼容;不能依靠“消费者已经处理”实现删除。按 key compaction 也不是隐私删除承诺,因为旧 segment、备份和跨集群副本可能仍存在。
项目接入要把生产与消费策略写进配置
应用配置基线:
messaging:
kafka:
bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS}
security-protocol: ${KAFKA_SECURITY_PROTOCOL:SASL_SSL}
topic: orders.events.v1
producer:
client-id: orders-api
acks: all
idempotence: true
delivery-timeout: 120s
compression: zstd
consumer:
group-id: orders.projector.v1
client-id: orders-projector-${INSTANCE_ID}
auto-commit: false
offset-reset: earliest
isolation-level: read_committed
max-poll-records: 100
max-poll-interval: 300s部署 readiness 不应只做 TCP 检查。Producer 服务应能获取目标 topic 元数据并确认 partition leader;consumer 服务应能完成认证、加入 group 并获得 assignment;平台探针则观察 KRaft quorum、offline partition 和 under-replicated partition。不要让每个实例启动时自动创建 topic,否则拼写错误、partition 数和 retention 会绕过治理。
事件契约至少包含稳定 event id、业务 key、event type、schema version、发生主体和追踪标识。Schema 不兼容会让 consumer 无限失败并积压;发布流程应在 CI 做兼容检查,consumer 对未知版本进入隔离流,而不是提交 offset 后静默丢弃。
用证据分类常见故障
| 现象 | 第一证据 | 常见根因 | 首个安全动作 |
|---|---|---|---|
| Bootstrap 成功,produce 超时 | metadata 中 leader 地址、client log | advertised listener 不可达、leader offline | 从客户端网络解析并连接返回地址 |
NotEnoughReplicas | topic describe、ISR、min ISR | follower 落后或 broker 故障 | 恢复 ISR,不降低 acks 掩盖故障 |
| Lag 增长 | group describe、每 partition lag | 消费慢、热点 partition、rebalance | 找最大 lag partition,算净追赶速率 |
| 消费不到历史记录 | committed offset、group id | 复用旧 group;误解 earliest | 用新 group 验证,reset 前先预览 |
| 频繁 rebalance | member 状态、poll 间隔、处理 P99 | poll 阻塞、实例抖动、超时不匹配 | 缩小批次、隔离长任务、检查 readiness |
| 重复处理 | event id、offset、事务/提交日志 | offset 提交前崩溃、producer 重试 | 业务幂等,不盲目跳 offset |
| 记录提前消失 | topic configs、earliest offset | retention/size/segment 清理 | 停止改配置,核对备份与上游重放能力 |
| 新 broker 空闲 | partition placement | 只扩节点,未做 reassignment | 制定带 throttle 的分批再分配 |
看到 UNKNOWN_TOPIC_OR_PARTITION 时,先确认 topic 名、权限和 metadata 传播,不打开自动创建;看到 OFFSET_OUT_OF_RANGE 时,说明 committed offset 已落在当前 log start/end 之外,通常由 retention 或数据迁移引起,应按业务选择 earliest/latest/备份恢复并记录数据缺口。
团队治理把高风险操作从日常账号剥离
每个 topic 应登记 owner、用途、key、schema、partition、replication factor、min ISR、cleanup policy、retention、峰值流量、生产者、consumer group、敏感级别、RPO/RTO、告警和删除条件。每个 group 还要登记 offset 管理者、最大允许 lag、重放路径和幂等能力。
普通应用账号不能创建/删除 topic、修改 retention、执行 partition reassignment、重置其他 group offset 或访问无关 topic。只读控制台也要隐藏 payload;管理 UI 不是授权边界,最终权限以 broker ACL 为准。
升级先验证 broker/client protocol 兼容、feature level、KRaft quorum、事务与 consumer protocol,再按 controller、broker、客户端的受控顺序滚动。数据目录不能用新旧不兼容镜像来回启动。降级能力必须在升级前由官方 upgrade notes 证明,不能假设换回旧镜像就能回滚 metadata。
成本治理同时看 broker 存储、跨可用区副本流量、跨区域复制、托管吞吐单元、长期保留和无效 consumer。降低 retention 能省存储,却可能让故障恢复窗口小于下游 RTO;增加副本提高容错,却增加写带宽和磁盘。架构师要把这些取舍变成预算和 SLO,而不是孤立参数。
清理实验与可恢复回滚
先确认实验 group 没有在运行,并记录 topic describe 与 group offset。删除两个实验 group:
docker exec te-kafka /opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--delete --group te.orders.projector.v1
docker exec te-kafka /opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--delete --group te.orders.audit-replay.v1删除开发 topic:
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--delete --topic te.orders.events.v1
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--delete --topic te.orders.retention.lab确认 topic 不在列表后停止容器:
docker exec te-kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 --list
docker compose down只有个人临时环境、且确认 volume 不含其他 topic、offset 和排障证据时,才删除数据:
docker compose down -v共享集群删除 topic 是不可逆数据变更,必须有 owner 确认、备份/上游重放判断和延迟执行窗口。Offset reset、retention 缩短、unclean leader election 与 reassignment 同样需要变更记录和回滚设计。
何时继续用 Kafka,何时换一种工具
Kafka 适合高吞吐、按 key 分区有序、保留回放、多 consumer group 独立读取和事件流集成。需要单条任务灵活路由、逐消息 ack/nack、低积压和 broker 自动 DLX 时,RabbitMQ 可能更直接;需要轻量 request/reply 时,NATS 的操作成本可能更低;需要复杂多租户与存储计算分离时,可以评估 Pulsar。
继续使用 Kafka 的条件是团队能明确 key 与 partition、接受至少一次并实现幂等、为 retention 支付存储、治理 consumer group、维护 KRaft/ISR 与安全边界。若系统只是为了“异步一下”却没有回放、分区和多订阅者需求,Kafka 的 controller、broker、partition、ACL、容量和升级成本可能高于收益。
上线判断应落在可观察事实上:KRaft quorum 有多数派;关键 partition 的 ISR 满足策略;producer 能区分成功、可重试和不确定结果;consumer 能在 rebalance 后安全恢复;offset reset 与 DLQ 重放受控;listener 在每个网络域可达;retention 大于恢复窗口;容量能承受峰值和副本追赶;凭证、payload 与高风险命令都有责任人。做到这些,Kafka 才是一条可恢复的事件日志,而不只是一个暴露 9092 的容器。
