相同键进入相同分区——动手操作日志
本实验在 VM 中运行
Ubuntu VM 上以 KRaft 单节点方式运行着 Apache Kafka 4.3.1(systemd
kafka 服务,localhost:9092)。/opt/kafka/bin 已在 PATH 中,所以可以直接使用 kafka-topics.sh
之类的工具。首次启动大约需要 4 分钟。
目标
创建主题并放入带键的事件,观察相同的键进入同一个分区、分区内部 顺序得到保持。读取消费者组的偏移量和延迟(lag),确认两个 组彼此独立地读取,然后增加分区数,亲自观察键的分布 发生变化。
为什么重要
要读懂“订单到了两次,又有一次消失了”这样的事故,首先必须知道 Kafka 承诺了什么、没有承诺什么。Kafka 只承诺分区内部的 顺序。相同的键会进入同一个分区,所以单个订单的事件 顺序会得到保持,但不同订单之间则不会。而且被消费的消息不会被删除 ——每个消费者组都各自有一个表示“读到哪里了”的整数(偏移量), 所以一个组漏掉的内容,另一个组可以从头重新读取。 这两点是后面各模块中关于重复与丢失的讨论的基础。
步骤
- 创建主题
orders,设置 3 个分区。 - 用
kafka-console-producer.sh把/root/kafka/orders.txt的九行(키:값格式,占位符依次为键、值)解析键之后放入。它们必须分散到多个分区中。 - 从头读取全部内容,使分区、偏移量、键都能看到,并保存到
/root/kafka/keyed.txt。相同的键必须只出现在一个分区里。 - 用消费者组
order-svc从头读取九条,并把对该组执行 describe 的结果保存到/root/kafka/group.txt。 - 再放入
/root/kafka/more.txt中的三行(不要读取),再次对order-svc执行 describe,并保存到/root/kafka/lag.txt。必须能看到 LAG。 - 用第二个组
analytics从头读取十二条。order-svc的偏移量必须保持不变。 - 把
orders增加到 6 个分区,放入order-1:refunded后从头读取,并保存到/root/kafka/repartition.txt。然后在/root/kafka/repartition-report.txt中写四行:partitions_before=3、partitions_after=6、order1_before=<3단계에서 order-1 이 있던 파티션>、order1_after=<refunded 가 들어간 파티션>(占位符依次为第 3 步中 order-1 所在的分区、refunded 所在的分区)。
参考
- 解析键:
--reader-property parse.key=true --reader-property key.separator=:(在 4.3 中--property已被弃用)。 - 读取时显示元数据:
--formatter-property print.key=true --formatter-property print.partition=true --formatter-property print.offset=true。输出形如Partition:0\tOffset:0\t키\t값(占位符依次为键与值)。 - 消费者用
--from-beginning --max-messages N --timeout-ms 8000读取 N 条后就结束。如果没有它,就得按 Ctrl-C。 - 组的状态:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group <이름>(占位符为名称)。消费者结束之后会显示“has no active members”,但偏移量表仍然照常显示。 - 增加分区:
kafka-topics.sh ... --alter --topic orders --partitions 6。不支持缩减(运维文档)。 - 常见错误 1:只给
--timeout-ms而不给--max-messages,会在最后一条消息之后再多等那么长时间。两个都要给。 - 常见错误 2:以为修改分区数之后旧数据会被迁移。Kafka 不会重新分配已有的数据。
有 3 个分区的主题
创建主题 orders,设置 3 个分区。
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic ... --partitions 3。创建之后,用 --describe 确认 PartitionCount。
放入带键的事件
用 kafka-console-producer.sh 把 /root/kafka/orders.txt 的九行(키:값 格式,占位符依次为键、值)解析键之后放入。它们必须分散到多个分区中。
加上 --reader-property parse.key=true --reader-property key.separator=:,并把文件作为标准输入传入。如果不解析键,整行都会成为值,键是 null,分区就与键无关了。
亲眼看看分区和偏移量
从头读取全部内容,使分区、偏移量、键都能看到,并保存到 /root/kafka/keyed.txt。相同的键必须只出现在一个分区里。
在 --from-beginning --max-messages 9 --timeout-ms 8000 的基础上加上 --formatter-property print.key=true --formatter-property print.partition=true --formatter-property print.offset=true,并把输出重定向到文件。不同分区之间的顺序会混在一起,但同一个分区内部,偏移量是递增的。
消费者组的偏移量
用消费者组 order-svc 从头读取九条,并把对该组执行 describe 的结果保存到 /root/kafka/group.txt。
给了 --group order-svc,读取的位置就会提交给 broker。结束之后,kafka-consumer-groups.sh --describe --group order-svc 的 CURRENT-OFFSET 在每个分区上都应该与 LOG-END-OFFSET 相等,LAG 应该是 0。
不读取,延迟就会堆积
再放入 /root/kafka/more.txt 中的三行(不要读取),再次对 order-svc 执行 describe,并保存到 /root/kafka/lag.txt。必须能看到 LAG。
用与第 2 步相同的方法放入,但不要运行消费者。LOG-END-OFFSET 会上升,CURRENT-OFFSET 保持不变,这个差就是 LAG。
各个组互相独立
用第二个组 analytics 从头读取十二条。order-svc 的偏移量必须保持不变。
--group analytics --from-beginning --max-messages 12。被消费的消息不会被删除,所以新的组可以从头读取全部内容,对其他组的偏移量没有任何影响。
增加分区,键的分布就会改变
把 orders 增加到 6 个分区,放入 order-1:refunded 后从头读取,并保存到 /root/kafka/repartition.txt。然后在 /root/kafka/repartition-report.txt 中写四行:partitions_before=3、partitions_after=6、order1_before=<3단계에서 order-1 이 있던 파티션>、order1_after=<refunded 가 들어간 파티션>(占位符依次为第 3 步中 order-1 所在的分区、refunded 所在的分区)。
执行 --alter --topic orders --partitions 6 之后,用解析键的方式放入一行,并像第 3 步那样用 --max-messages 13 读取。正如运维文档所说,默认分区器是 hash(key) % 分区数,所以分区数一变,相同的键就可能进入另一个分区,而已有的数据不会被迁移。