同じキーは同じパーティションへ — ログを実際に操作する
このラボはVM上で動きます
UbuntuのVMに、Apache Kafka 4.3.1がKRaftの単一ノードとして起動しています(systemdの
kafkaサービス、localhost:9092)。/opt/kafka/binがPATHに入っているので、kafka-topics.sh
のようなツールをすぐに使えます。最初の起動には4分ほどかかります。
目標
トピックを作ってキー付きのイベントを入れ、同じキーが同じパーティションに行き、パーティションの 中で順序が守られるのを確認します。コンシューマーグループのオフセットとラグ(lag)を読み、2つの グループが互いに独立して読むことを確認した後、パーティション数を増やすとキーの配置が 変わるのを自分で観察します。
なぜ重要なのか
「注文が2回届いて、1回は消えた」という事故を読み解くには、まずKafkaが 何を約束し、何を約束しないのかを知っておく必要があります。Kafkaはパーティションの中の 順序だけを約束します。同じキーが同じパーティションに行くので、注文1つのイベントの 順序は守られますが、別の注文どうしではそうではありません。そして、消費されたメッセージは削除 されません。コンシューマーグループごとに「どこまで読んだか」という整数1つ(オフセット)を別々に 持つので、あるグループが見逃したものを、別のグループは最初から読み直せます。 この2つが、後のモジュールの重複・消失の話の土台です。
ステップ
ordersトピックを、パーティション3つで作成してください。/root/kafka/orders.txtファイルの9行(키:값の形式で、プレースホルダーはキーと値です)を、kafka-console-producer.shでキーをパースして入れてください。複数のパーティションに分かれて入らなければなりません。- パーティション・オフセット・キーが見えるように、最初からすべて読み、
/root/kafka/keyed.txtファイルに保存してください。同じキーは、1つのパーティションにだけなければなりません。 - コンシューマーグループ(
order-svc)で、最初から9件を読み、そのグループをdescribeした結果を、/root/kafka/group.txtファイルに保存してください。 /root/kafka/more.txtファイルの3行をさらに入れて(読みはせずに)、order-svcをもう一度describeし、/root/kafka/lag.txtファイルに保存してください。LAGが見えなければなりません。- 2つ目のグループ(
analytics)で、最初から12件を読んでください。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 가 들어간 파티션>の4行を書いてください(プレースホルダーは、ステップ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값の形です(出力の後半の2つの列は、キーと値です)。 - コンシューマーは、
--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:
--max-messagesなしで--timeout-msだけを指定すると、最後のメッセージの後にその時間だけさらに待ちます。両方を指定してください。 - よくあるミス2: パーティション数を変えた後、古いデータが移動されると期待することです。Kafkaは既存のデータを再配置しません。
パーティション3つのトピック
ordersトピックを、パーティション3つで作成してください。
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic ... --partitions 3。作成した後、--describeでPartitionCountを確認してください。
キー付きのイベントを入れる
/root/kafka/orders.txtファイルの9行(키:값の形式で、プレースホルダーはキーと値です)を、kafka-console-producer.shでキーをパースして入れてください。複数のパーティションに分かれて入らなければなりません。
--reader-property parse.key=true --reader-property key.separator=:を指定して、ファイルを標準入力として渡してください。キーをパースしないと、1行全体が値になり、キーがnullなので、パーティションがキーと無関係に決まります。
パーティションとオフセットを目で見る
パーティション・オフセット・キーが見えるように、最初からすべて読み、/root/kafka/keyed.txtファイルに保存してください。同じキーは、1つのパーティションにだけなければなりません。
--from-beginning --max-messages 9 --timeout-ms 8000に、--formatter-property print.key=true --formatter-property print.partition=true --formatter-property print.offset=trueを加えて、出力をファイルに送ってください。パーティション間では順序が混ざって出力されますが、1つのパーティションの中ではオフセットが上がります。
コンシューマーグループのオフセット
コンシューマーグループ(order-svc)で、最初から9件を読み、そのグループをdescribeした結果を、/root/kafka/group.txtファイルに保存してください。
--group order-svcを指定すると、読んだ位置がブローカーにコミットされます。終了した後、kafka-consumer-groups.sh --describe --group order-svcのCURRENT-OFFSETが、パーティションごとにLOG-END-OFFSETと同じで、LAGが0でなければなりません。
読まなければラグが溜まる
/root/kafka/more.txtファイルの3行をさらに入れて(読みはせずに)、order-svcをもう一度describeし、/root/kafka/lag.txtファイルに保存してください。LAGが見えなければなりません。
ステップ2と同じ方法で入れますが、コンシューマーは動かさないでください。LOG-END-OFFSETは上がり、CURRENT-OFFSETはそのままなので、その差がLAGです。
グループは互いに独立している
2つ目のグループ(analytics)で、最初から12件を読んでください。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 가 들어간 파티션>の4行を書いてください(プレースホルダーは、ステップ3でorder-1があったパーティションと、refundedが入ったパーティションです)。
--alter --topic orders --partitions 6の後に、キーのパースで1行を入れて、ステップ3と同じように--max-messages 13で読んでください。運用ドキュメントが述べるように、デフォルトのパーティショナーは、hash(key)をパーティション数で割った余りでパーティションを決めるので、パーティション数が変わると同じキーが別のパーティションに行くことがあり、既存のデータは移動されません。