TT Lab
はじめる
学ぶ 学習パス コース

注文が二度届き、一度は消えた

同じキーは同じパーティションへ — ログを実際に操作する

TT Labで続きを見る

このラボは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つが、後のモジュールの重複・消失の話の土台です。

ステップ

  1. ordersトピックを、パーティション3つで作成してください。
  2. /root/kafka/orders.txtファイルの9行(키:값の形式で、プレースホルダーはキーと値です)を、kafka-console-producer.shでキーをパースして入れてください。複数のパーティションに分かれて入らなければなりません。
  3. パーティション・オフセット・キーが見えるように、最初からすべて読み、/root/kafka/keyed.txtファイルに保存してください。同じキーは、1つのパーティションにだけなければなりません。
  4. コンシューマーグループ(order-svc)で、最初から9件を読み、そのグループをdescribeした結果を、/root/kafka/group.txtファイルに保存してください。
  5. /root/kafka/more.txtファイルの3行をさらに入れて(読みはせずに)、order-svcをもう一度describeし、/root/kafka/lag.txtファイルに保存してください。LAGが見えなければなりません。
  6. 2つ目のグループ(analytics)で、最初から12件を読んでください。order-svcのオフセットは、そのままでなければなりません。
  7. 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が入ったパーティションです)。

参考

パーティション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)をパーティション数で割った余りでパーティションを決めるので、パーティション数が変わると同じキーが別のパーティションに行くことがあり、既存のデータは移動されません。