Upgrade to Pro — share decks privately, control downloads, hide ads and more …

ZITADEL イベントストアの Projection 実行経路と排他制御

ZITADEL イベントストアの Projection 実行経路と排他制御

社内勉強会の LT 資料です。

ZITADEL イベントストアについて書いたブログはこちら:
https://zenn.dev/xtm_blog/articles/a1af216ddbef4f

Avatar for Kenta Okamoto

Kenta Okamoto

August 06, 2026

More Decks by Kenta Okamoto

Other Decks in Programming

Transcript

  1. ブログのおさらい ZITADEL はイベントストアにイベントを記録し、Projection で参照用テーブルに反映している CQRS とは 書き込み (Command) と読み取り (Query)

    の責務を分離する Event Sourcing とは 変化をイベントとして記録する Projection とは イベントを畳み込んで、参照用テーブルへ反映する処理 zenn.dev/xtm_blog/articles/a1af216ddbef4f 03 / 14
  2. 用語 instance ZITADEL のテナント。テナントビルのようなもの ZITADEL Instance auth.example.com Organization: A Organization:

    B Policies Private Labeling / Password / Login Policies Private Labeling / Password / Login Project: X Roles / Applications Granted Project: Z Organization A の Project Grant で貸し出される Users Role Assignments Project: Z Roles / Applications Users Role Assignments インスタンスの中に Organization・Project・User が入る。ZITADEL のプロセスそのものではない 07 / 14
  3. PATH 01 schedule 60 秒ごとに実行される Handler.Start schedule ActiveInstancer triggerInstances go

    h.schedule(ctx) loop [RequeueEvery (60 秒) ごと] ActiveInstances() instance ID の一覧 triggerInstances(instances) instance ごとに Trigger 失敗なら RetryFailedAfter 待って成功するまでリトライ 09 / 14
  4. PATH 02 subscribe インメモリの pub/sub。チャネルが満杯の場合は破棄される 取りこぼす前提の使い方で、最終的に schedule が整合性を担保する command →

    eventstore.Push Eventstore.PushWithClient Subscription.Events チャネル Handler.subscribe Push(cmds...) events2 に INSERT es.notify(mappedEvents) queue chan Event (バッファ 100) 満杯なら default: で破棄 instance 単位に重複排除 重複排除したうえで instance ごとに Trigger を呼ぶ 10 / 14
  5. PATH 03 trigger 読み取り整合が必要な API は同期的に trigger を呼ぶ 例:GET /v2/users/{user_id}

    REST クライアント Queries.GetUserByID Trigger ×2 PostgreSQL GET /v2/users/{user_id} shouldTriggerBulk が true → triggerBatch() par [UserProjection] lockInstance → advisory lock → users14 更新 [LoginNameProjection] lockInstance → advisory lock → login_names3* 更新 wg.Wait() で両方の完了を待つ SELECT (users14 に login_names を JOIN) 200 (JSON) 2 つの Trigger が終わるまでレスポンスは返らない 11 / 14
  6. 全体像 キック経路 schedule subscribe 同期 Trigger Handler.Trigger 不可 skip 取得

    ワーカープール ① lockInstance 再実行 processEvents ① lockInstance プロセス内ロック processEvents ② advisory lock + ③ 行ロック 再実行 additionalIteration 終了 08 / 14
  7. 多重処理を防ぐ排他制御 3 つのロックで排他制御を実現している プロセス内ロック プロセス内でインスタンス単位でロック PostgreSQL の advisory lock pg_try_advisory_xact_lock

    行ロック current_states / projection 名 × インスタンス ID 単位でロック (projection をどこまで処理したかを記録するテーブル)の行の更新をロック 12 / 14
  8. 例: ZITADEL のプロセスが複数動いているときのロックのイメージ ZITADEL A ZITADEL B PostgreSQL ① lockInstance

    取得 ① lockInstance 取得 (別プロセスなので通る) ② pg_try_advisory_xact_lock(users14, inst1) true ② 同じロックを試みる false → ログのみ出して正常終了 (エラーにしない) ③ current_states を FOR NO KEY UPDATE → reduce → UPSERT COMMIT (② advisory lock も自動解放) 13 / 14