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

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

何を残し、何を捨てるか — retention・compact・min.insync.replicas

TT Labで続きを見る

一言でいうと

トピックの設定は、3つの質問に答えます。どれだけ長く残すのか(retention.ms、 retention.bytes)、キーごとに最後の値だけを残すのか(cleanup.policy=compact)、 書き込みを成功と見なすには、レプリカが何個確認する必要があるのか(min.insync.replicasと acks=all)。削除と圧縮はセグメント単位で起き、アクティブなセグメントには手を触れません。 「削除したのに残っている」の大半は、これです。

なぜ必要なのか

Kafkaは、消費しても削除しないので、何かは削除しなければなりません。トピック設定のドキュメント のcleanup.policyには、2つのポリシーがあります。delete(デフォルト)は、保持時間やサイズの 上限に達した古いセグメントを捨て、compactは、キーごとに最新の値を残すログの 圧縮を有効にします。両方を併記すると(delete,compact)、古いセグメントは保持のルールで 捨て、残ったセグメントは圧縮します。空のリストは、無期限の保持です。

圧縮が必要な理由は、設計ドキュメントの 「Log Compaction」の節が例で説明しています。ユーザー123のメールアドレスが3回変わったとき、 時間による保持は、古い変更をまるごと捨ててしまい、最初から読んでも現在の状態を復元 できなくなります。圧縮は、キーごとに最後の更新を必ず残して、ログがすべてのキーの最終的な 値のスナップショットになるようにします。DBの変更の購読、イベントソーシング、状態のジャーナリングが、この上に成り立ちます。

どう動くのか

時間による保持は、セグメント単位です: retention.msのデフォルト値は604800000(7日)で、 -1なら無制限です。ドキュメントは、これを「コンシューマーがどれだけ速く読まなければならないかに関する SLA」と呼んでいます。retention.bytesは、パーティションあたりのサイズの上限で、デフォルトは-1(なし)です。 ところが、削除される単位はレコードではなく、セグメントファイルです。segment.ms(デフォルト 7日)やsegment.bytes(デフォルト1GiB)に達して、セグメントがロールされて初めて、古いセグメントが 削除の候補になり、アクティブなセグメントは決して削除されません。そして、ブローカーは ブローカー設定の log.retention.check.interval.ms(デフォルト300000、5分)ごとにしか検査しません。そのため、 retention.ms=1분(末尾の韓国語の単位は「分」を意味します)を指定しても、セグメントがロールされていなかったり、検査の周期が来ていなかったりすれば、 残っています。ラボのVMは、検査の周期を10秒に縮めてあります。

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

このコードブロックの韓国語コメントは、12秒後にe4が入るとセグメント0が閉じてロールし、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)は、ブローカーが受け付けるレコード バッチの最大サイズです。超えるバッチは、プロデューサーに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は ブローカーが1つなので、この相互作用を再現できず(実測でも、単一ノードでは エラーは出ませんでした)、クイズでだけ確認します。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で取り除きます(2つのコマンドの間の韓国語は、「で設定を追加し、」という意味です)。トピックを 作り直す必要はありません。

現場での姿

「保持を1日に減らしたのに、ディスクが減らない」という報告は、セグメントのサイズを見れば解決します。 segment.bytesが1GiBで、1日に200MBずつ溜まるパーティションは、5日が過ぎて初めて、 セグメントがロールされ、それまでは何も削除されません。segment.msを 併せて縮めるのが答えです。

「圧縮トピックなのに、古い値が見える」は、3つを確認します。アクティブなセグメントにあるか (圧縮されません)、dirty ratioを超えたか(デフォルト0.5)、そして、コンシューマーが最初から 読んでいて、まだheadに到達していないか(圧縮はtailで行われます)。3つとも 正常な動作です。

次のラボですること

20秒保持のトピックで、セグメントをロールさせて、一番前のオフセットが上がるのを見て、圧縮 トピックでaliceが最後の値1つに減るのを見て、max.message.bytes で大きなレコードが拒否されるのを確認した後、実行中のトピックの保持を変更します。