TT Lab
开始
学习 学习路径 课程

订单重复到达,又消失了一次

留下什么、丢弃什么——retention、compact、min.insync.replicas

在 TT Lab 中继续学习

一句话总结

主题配置回答三个问题:要保留多久(retention.ms、 retention.bytes),是否只为每个键保留最后一个值(cleanup.policy=compact), 以及要把一次写入算作成功,需要多少个副本确认(min.insync.replicas 与 acks=all)。删除和压缩都是以段(segment)为单位进行的,不会动活动段 ——“明明删了却还在”,大部分都是这个原因。

为什么需要它

Kafka 消费之后不会删除,所以总得删掉点什么。主题配置文档 中的 cleanup.policy 设有两种策略。delete(默认)会丢弃达到保留时间或大小 上限的旧段,compact 则开启为每个键保留最新值的日志 压缩。如果把两者一起写(delete,compact),旧段按保留规则 丢弃,剩下的段再做压缩。空列表表示无限期保留。

为什么需要压缩,设计文档 的“Log Compaction”一节用一个例子作了说明。用户 123 的邮箱改了三次时, 按时间保留会把较早的变更整个丢掉,这样即使从头读取,也无法 还原当前状态。压缩则一定会为每个键保留最后一次更新,让日志成为所有键的最终 值的快照——订阅 DB 变更、事件溯源、状态日志都建立在这之上。

工作原理

按时间保留是以段为单位的。 retention.ms 的默认值是 604800000(7 天), -1 表示无限制。文档把它称为“关于消费者必须多快读取的 SLA”。retention.bytes 是每个分区的大小上限,默认是 -1(没有)。 但是被删除的单位不是记录,而是段文件。只有达到 segment.ms(默认 7 天)或 segment.bytes(默认 1GiB)而使段发生滚动,旧段才会成为 待删除的候选,而活动段永远不会被删除。而且 broker 只按 broker 配置中的 log.retention.check.interval.ms(默认 300000,即 5 分钟)的周期检查。所以 即使给了 retention.ms=1분(韩文,意为“1 分钟”),如果段没有滚动,或者检查周期还没到, 数据就仍然留着。实验 VM 把检查周期缩短到了 10 秒。

events (retention.ms=20000, segment.ms=10000)
  세그먼트 0: e1 e2 e3        ← 12초 뒤 e4 가 들어오며 굴러간다 (닫힘)
  세그먼트 4: e4              ← 활성. 절대 지워지지 않는다
  20초 + 검사 주기 뒤 → 세그먼트 0 삭제 → 가장 앞 오프셋이 3 이 된다

压缩会为每个键保留最后一个值。 设计文档的保证——压缩不会改变顺序, 只会去掉记录,而偏移量绝不会改变(缺失的偏移量按与其后一个 偏移量相同的位置处理)。从头读取的消费者会按写入顺序 看到所有键的最终状态。键的值为 null 的记录是删除标记(tombstone),它会使该键的 旧记录被删除,自己也会在 delete.retention.ms(默认 1 天)之后消失。 什么时候压缩,由 min.cleanable.dirty.ratio(默认 0.5——日志有一半是重复的才会 启动)、min.compaction.lag.ms、max.compaction.lag.ms 决定,这里活动 段同样不在对象之内。实验把 dirty ratio 降到 0.01,使其很快就被压缩。

过大的记录。 max.message.bytes(默认 1048588)是 broker 接受的记录 批次的最大大小。超过的批次会以 RecordTooLargeException 被 拒绝,而日志里什么都不会留下——这已经通过实测确认过了。

min.insync.replicas 与 acks=all。 min.insync.replicas 条目——生产者 以 acks=all 发送时,写入要想成功必须确认的最小 ISR 数量(包括领导者)。 如果 ISR 少于这个数,生产者会收到 NotEnoughReplicas 或 NotEnoughReplicasAfterAppend 异常。文档举出的典型配置是副本 因子 3、min.insync.replicas=2、acks=all——必须有过半数确认写入才 算成功,而且无论 acks 如何,数据都会被复制到整个 ISR,在满足这个条件之前, 消费者是看不到的。默认值是 1,所以没有任何保证。本课程的 VM 只有 一个 broker,无法重现这种交互(实测中,在单节点上也 没有出现错误),只能通过测验来确认。unclean.leader.election.enable(默认 false)决定是否把 ISR 之外的副本作为最后手段选为领导者,开启的话可能会 丢失数据。

运行中修改。 运维文档 中的“Modifying topics”——分别用 kafka-configs.sh --entity-type topics --entity-name X --alter --add-config k=v 로 설정을 더하고 --delete-config k 来添加和删除配置。不需要 重新创建主题。

在现场相遇的样子

“把保留期缩短到 1 天了,磁盘却没有变小”这类反馈,看一下段的大小就能解开。 对于 segment.bytes 是 1GiB、每天积累 200MB 的分区,要过五天 段才会滚动,在那之前什么都不会被删除。答案是同时缩短 segment.ms。

“明明是压缩主题,却看到了旧值”,要确认三件事:它是否在活动段中 (不会被压缩)、是否超过了 dirty ratio(默认 0.5),以及消费者从头 读取时是否还没有到达 head(压缩发生在 tail)。这三种 都是正常行为。

下一项实验要做什么

在保留 20 秒的主题上让段滚动,看到最前面的偏移量上升;在压缩 主题上看到 alice 缩减为只剩最后一个值;通过 max.message.bytes 确认大记录会被拒绝,然后修改一个运行中主题的保留时间。