Outboxパターンでイベント発行を保証する
目標
二重書き込みの不整合を自分で作ってみた後、アウトボックスパターンで消失をなくし、残る重複をコンシューマーの冪等性で吸収するまでの全区間を、手を動かして完成させます。
なぜ重要なのか
DBへの保存とイベントの発行を並べて呼び出すコードは、どこにでもあり、普段はうまく動作します。問題は、ブローカーが3秒間不安定になるその瞬間にだけ現れ、そのとき生じた不整合は静かに残ります。注文はあるのに在庫が減っていない状態を、数日後の精算で発見する、という具合です。アウトボックスは、この問題を、「ブローカーをトランザクションに参加させる」ではなく、「発行の意図をDBに一緒に書き込む」と、ひっくり返して解きます。その代わり、新しい性質が生まれます。消失はなくなりますが、重複は生じます。リレーが発行の直後、状態の更新の前に落ちると、再起動後に同じイベントをもう一度送ります。そのため、このラボでは、重複をバグではなく設計の前提として扱い、コンシューマー側で吸収するところまでを、1つのラボの中で終えます。
ステップ
/root/outbox/init.pyで/root/outbox/app.dbを作成してください。orders(id, sku, qty)とoutbox(event_id, aggregate_id, seq, event_type, payload, status)の2つのテーブルがあり、statusのデフォルト値はPENDINGでなければなりません。/root/outbox/dualwrite.pyは、注文を1件保存した後、発行に失敗します。実行後、/root/outbox/dualwrite.outの1行目にINCONSISTENT orders=<n> published=<m>を書いてください。nとmが異なっていなければなりません。/root/outbox/place_order.pyは、ordersとoutboxに同じトランザクションで書き込みます。3回実行すると、ordersが3行、outboxが3行になります。/root/outbox/relay.pyは、status='PENDING'の行をRedisのリストoutbox.eventsにRPUSHして、その行をPUBLISHEDに変更します。2回実行しても、リストの長さが増えてはいけません。- 同じ
aggregate_idのイベントが、seqの昇順でキューに入らなければなりません。/root/outbox/order_check.outにORDER OKを残してください。 /root/outbox/relay_crash.pyは、発行だけを行って、状態を更新しません。実行後に正常なリレーをもう一度動かすと、同じevent_idがキューに2回入ります。/root/outbox/atleastonce.outにDUPLICATE event_id=<id> count=2を書いてください。/root/outbox/consumer.pyは、キューを空にしながら、event_idを基準に重複をふるい落として処理します。処理結果をRedisのハッシュprocessedに残してください。重複があっても、processedのサイズは、ユニークなイベントの数と同じでなければなりません。/root/outbox/report.txtに、orders=<n>、outbox=<n>、enqueued=<n>、processed_unique=<n>の4行を書いてください。enqueuedは、processed_unique以上でなければなりません。
参考
PENDINGの問い合わせ用の部分インデックス:CREATE INDEX ... ON outbox(status) WHERE status='PENDING'- Redisの確認:
redis-cli LLEN outbox.events、redis-cli HLEN processed - よくあるミス1: リレーが、発行と状態の更新を1つの原子的な単位にしようと頑張ることです。不可能です。重複を受け入れて、コンシューマーでふるい落とすのが正解です。
- よくあるミス2:
seqなしでタイムスタンプでソートすることです。同じミリ秒に2つのイベントが生じると、順序が逆転します。
注文とアウトボックスのテーブルを作る
/root/outbox/init.pyで/root/outbox/app.dbを作成してください。orders(id, sku, qty)とoutbox(event_id, aggregate_id, seq, event_type, payload, status)の2つのテーブルがあり、statusのデフォルト値はPENDINGでなければなりません。
python3のsqlite3モジュールで十分です。アウトボックスの行には、イベントの識別子、集約の識別子、タイプ、本文、状態が必要です。
二重書き込みの不整合を再現する
/root/outbox/dualwrite.pyは、注文を1件保存した後、発行に失敗します。実行後、/root/outbox/dualwrite.outの1行目にINCONSISTENT orders=<n> published=<m>を書いてください。nとmが異なっていなければなりません。
DBへの保存は成功させて、ブローカーへの発行だけを失敗させれば済みます。2つのストアの件数を数えて、互いに異なることをファイルに残してください。
1つのトランザクションに折り込む
/root/outbox/place_order.pyは、ordersとoutboxに同じトランザクションで書き込みます。3回実行すると、ordersが3行、outboxが3行になります。
2つのINSERTを、同じコネクション、同じコミットの中に入れます。コミットの前に例外が起きたら、両方ともないのが正常です。
リレーでブローカーに移す
/root/outbox/relay.pyは、status='PENDING'の行をRedisのリストoutbox.eventsにRPUSHして、その行をPUBLISHEDに変更します。2回実行しても、リストの長さが増えてはいけません。
PENDINGの行を読んでRedisのリストに押し込み、状態を変更します。何回実行しても、すでに移したものをもう一度移さないようにしなければなりません。
集約ごとの順序を守る
同じaggregate_idのイベントが、seqの昇順でキューに入らなければなりません。/root/outbox/order_check.outにORDER OKを残してください。
同じ注文のイベントは、発生した順序どおりに出ていかなければなりません。ソートの基準を何にするかを考えてみてください。時刻は同じになることがあります。
リレーのクラッシュで重複を作る
/root/outbox/relay_crash.pyは、発行だけを行って、状態を更新しません。実行後に正常なリレーをもう一度動かすと、同じevent_idがキューに2回入ります。/root/outbox/atleastonce.outにDUPLICATE event_id=<id> count=2を書いてください。
発行はしたのに、状態の更新ができずに落ちる状況を再現します。もう一度動かすと、同じイベントが2回キューに入ります。
コンシューマーに重複排除を付ける
/root/outbox/consumer.pyは、キューを空にしながら、event_idを基準に重複をふるい落として処理します。処理結果をRedisのハッシュprocessedに残してください。重複があっても、processedのサイズは、ユニークなイベントの数と同じでなければなりません。
すでに処理したイベントの識別子を覚えておけば済みます。Redisの集合のデータ構造や、SET NXが向いています。
全区間を数値で報告する
/root/outbox/report.txtに、orders=<n>、outbox=<n>、enqueued=<n>、processed_unique=<n>の4行を書いてください。enqueuedは、processed_unique以上でなければなりません。
注文の数、アウトボックスの行数、キューに入った数、実際に処理されたユニークな数を、1つのファイルにまとめます。最初の3つの値と、最後の値の関係が、このパターンのすべてです。