跳到主要内容

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关键配置说明
事件流与审计流deleteretention.ms、retention.bytes保留期应覆盖最大消费延迟和回放窗口
状态快照与配置compactmin.compaction.lag.msConsumer 必须正确处理 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 或已验证的补偿机制。