何を残し、何を捨てるか — retention・compact・min.insync.replicas
一言でいうと
トピックの設定は、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
で大きなレコードが拒否されるのを確認した後、実行中のトピックの保持を変更します。