Topic 规划与生产者可靠性
Topic 分区数决定消费者并行度上限,副本数决定可用性和存储成本,保留策略决定可回放窗口。分区创建后只能增加不能减少;增加分区会改变按 Key 的路由,因此必须评估业务顺序和下游幂等能力。
1. 创建与检查 Topic
# 显式定义分区、副本、最小 ISR 与保留期,避免依赖不明确的 Broker 默认值。
/opt/kafka/bin/kafka-topics.sh --create --if-not-exists \
--bootstrap-server kafka-1.example.internal:9092 \
--topic orders.created.v1 --partitions 12 --replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000 \
--config cleanup.policy=delete
# 检查每个分区的 Leader、Replica 与 ISR;ISR 应包含所有健康副本。
/opt/kafka/bin/kafka-topics.sh --describe \
--bootstrap-server kafka-1.example.internal:9092 --topic orders.created.v1
# 查看 Topic 的最终配置,确认是否存在动态覆盖值。
/opt/kafka/bin/kafka-configs.sh --describe \
--bootstrap-server kafka-1.example.internal:9092 \
--entity-type topics --entity-name orders.created.v1
| 业务类型 | cleanup.policy | 关键配置 | 说明 |
|---|---|---|---|
| 事件流与审计流 | delete | retention.ms、retention.bytes | 保留期应覆盖最大消费延迟和回放窗口 |
| 状态快照与配置 | compact | min.compaction.lag.ms | Consumer 必须正确处理 Tombstone |
| 状态变更加短期回放 | compact,delete | 压缩与时间保留 | 上线前验证删除与压缩的业务语义 |
2. 生产者可靠性
关键事件使用 acks=all、幂等与重试。幂等只能避免单 Producer 会话内的重试重复,跨服务重试和下游写入仍必须采用业务唯一键保障幂等。
# producer.properties
# 等待最小 ISR 全部确认;Topic 的 min.insync.replicas=2 才能形成可靠约束。
acks=all
# 开启幂等,避免可恢复重试导致同一分区内重复写入。
enable.idempotence=true
# 允许可恢复失败重试;应用仍必须处理最终失败和超时。
retries=2147483647
# 同一业务实体使用稳定 Key,确保其消息落入同一分区并保持局部顺序。
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer
# 用固定 Key 写入一条烟雾测试消息,验证 Key 分区和连通性。
printf 'order-1001|created\n' | /opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server kafka-1.example.internal:9092 \
--topic orders.created.v1 \
--property parse.key=true --property key.separator='|'
# 从头读取一条消息并打印 Key;仅用于测试,不应用于生产消费组。
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka-1.example.internal:9092 \
--topic orders.created.v1 --from-beginning --max-messages 1 \
--property print.key=true --property key.separator='|'
3. 分区扩容
# 分区数只能增加;执行前评估 Key 路由变化、消费者并发与下游处理能力。
/opt/kafka/bin/kafka-topics.sh --alter \
--bootstrap-server kafka-1.example.internal:9092 \
--topic orders.created.v1 --partitions 18
# 扩容后检查副本分布,避免新分区持续集中在同一 Broker。
/opt/kafka/bin/kafka-topics.sh --describe \
--bootstrap-server kafka-1.example.internal:9092 --topic orders.created.v1
过多小分区会放大文件句柄、Controller 元数据、Leader 选举与重平衡成本。分区数应依据每分区吞吐、消费者并行度、Broker 磁盘和恢复时长压测确定。
4. 顺序、消息大小与事务边界
Kafka 只保证单分区顺序。同一订单、账户或设备应使用稳定 Key;不同 Key 跨分区没有全局顺序。超大消息会放大页缓存、复制、网络和恢复成本,通常应将大对象放入对象存储,只在 Kafka 传递引用和校验信息。
# 事务生产者只适用于需要原子写多个分区且下游支持 read_committed 的场景。
# transactional.id 必须稳定且按实例唯一,防止存活实例互相围栏。
transactional.id=orders-writer-01
enable.idempotence=true
acks=all
Kafka 事务不能自动解决数据库和 Kafka 的双写一致性。涉及数据库状态和事件发布时,应使用 Outbox、CDC 或已验证的补偿机制。