如何配置Kafka broker实现消息仅保留1秒的快速过期清理
Kafka 消息1秒自动清理配置方案
原配置不生效的核心原因
log.cleanup.policy配置被注释,即使默认值为delete,也存在集群自定义过默认策略的风险- 日志清理线程默认运行间隔为5000ms,远大于设置的1秒保留时间,导致过期消息无法及时被清理
- Kafka仅对已关闭的非活跃日志段执行清理,未配置日志段滚动规则时,活跃日志段迟迟不关闭,内部的过期消息无法被删除
- 已创建的
test-topic会继承创建时的broker配置,后续修改的全局broker配置不会自动同步到已有topic
第一步:修正broker全局配置
调整server.properties的相关配置,补全缺失项:
auto.create.topics.enable=true default.replication.factor=3 min.insync.replicas=2 num.io.threads=8 num.network.threads=5 num.partitions=1 num.replica.fetchers=2 replica.lag.time.max.ms=10000 socket.receive.buffer.bytes=102400 socket.request.max.bytes=104857600 socket.send.buffer.bytes=102400 unclean.leader.election.enable=true zookeeper.session.timeout.ms=18000 # 启用删除策略,取消原注释 log.cleanup.policy=delete # 消息保留时间1秒,原有配置保留 log.retention.ms=1000 # 清理线程运行间隔调整为500ms,加快清理频率 log.cleaner.backoff.ms=500 # 日志段滚动时间设为1秒,让活跃段快速关闭可被清理 log.roll.ms=1000 # 限制单日志段大小,避免数据量小时无法触发滚动 log.segment.bytes=10485760 log.flush.interval.ms=1000
修改后重启所有broker节点生效。
第二步:更新已有topic的配置
执行命令手动修改已创建的test-topic配置,覆盖创建时的旧参数:
# 对接ZooKeeper的版本使用该命令 ./bin/kafka-configs.sh --zookeeper <你的zk地址> --alter --topic test-topic --add-config retention.ms=1000,cleanup.policy=delete,segment.roll.ms=1000 # 新版无需对接ZK的Kafka使用该命令 ./bin/kafka-configs.sh --bootstrap-server broker1:9092,broker2:9092,broker3:9092 --alter --topic test-topic --add-config retention.ms=1000,cleanup.policy=delete,segment.roll.ms=1000
第三步:验证生效
执行以下命令查看topic配置是否更新成功:
./bin/kafka-configs.sh --bootstrap-server broker1:9092,broker2:9092,broker3:9092 --describe --topic test-topic
注意:Kafka的消息清理是异步操作,存在最多数百毫秒的误差,无法保证完全精确的1秒到期即删,属于正常现象
内容的提问来源于stack exchange,提问作者Jong - Sung Kim
相关产品推荐
相关产品推荐

