Databricks、Temporal と Lakebase を活用した堅牢なエージェント構築を公開
本文の状態
日本語全文を表示中
詳細モードで約23分の本文を読めます。
同じ出来事の情報源
この情報源を基点に整理
Databricks AI Engineering
Databricks AI Engineering は、長期実行型クラウドエージェントの信頼性を高めるため、Temporal と Lakebase を組み合わせた参照実装を発表し、失敗耐性と監査機能を強化する。
AI深層分析を開く2026年9月9日 02:49
AI深層分析
キーポイント
長期実行エージェントの課題と要件
クラウドエージェントはプロセスやコンテナの寿命を超えて実行されるため、回復力、再試行、長時間待機、運用可視性、ガバナンス、監査という 6 つの要件を満たす必要がある。
Temporal と Lakebase の連携アーキテクチャ
Temporal が永続的な実行を担い、Lakebase Postgres が照会可能な運用状態を管理する構成により、作業の継続性とデータの整合性を確保する。
Databricks エコシステムとの統合機能
Unity Catalog の下書きポリシーを Lakebase に同期し、変更データフィードを通じて Delta 履歴テーブルに自動反映することで、既存の Databricks 管理環境とシームレスに連携する。
Temporal と Lakebase の役割分担
Temporal はワークフローのイベント履歴を通じて耐障害性を担保し、Lakebase はアプリケーション向けのステータスや証拠などの状態を管理する。両システムはトランザクションを共有せず、それぞれ異なる消費者のために異なる状態を保持する。
エージェント実行の信頼性確保
Temporal の Activity として Lakebase を呼び出す際、一意な識別子や制約、ガード付き更新を用いて、再試行時でも同じ論理レコードを正確にターゲットする。これにより、ワーカーの入れ替えや長時間待機後もセッションが生存し続ける。
重要な引用
A cloud agent may outlive the request, worker, container, or deployment that started it.
Recovery requires both the results of completed operations and the control-flow state needed to determine what happens next.
Temporal simplifies the management of distributed systems.
Temporal plus Lakebase is most useful when an agentic session must survive worker replacement, accept input after long waits, expose relational state to an application, and apply governed data while it remains open.
編集コメントを表示
編集コメント
エージェント技術が実用段階に進む中で、プロセス障害や長時間待機時の状態管理は最も重要な課題の一つである。今回の参照実装は、Databricks ユーザーにとって即座に適用可能な堅牢な解決策を提供する。
Source Article
元記事を日本語で読む
本文に関係しない購読案内、埋め込み通知、サイト内プロモーションは除いています。
個人ローン審査エージェントは証拠を収集し、ポリシーを適用しますが、レビュー担当者の確認に数日かかることもあります。その間、作業者が再起動したり、ツール呼び出しが失敗したりする可能性があります。そのため、アプリケーションは完了した作業を保存し、実行を再開でき、レビュー担当者が証拠にアクセスできる状態を維持する必要があります。
この参考実装では、永続的な実行に Temporal を、照会可能な運用ステータス管理に Lakebase Postgres を使用しています。同期されたテーブルにより、Unity Catalog 上の審査ポリシーが Lakebase で利用可能になります。Temporal の Activity は証拠、判断結果、メトリクスを Lakebase に書き込みます。有効化すれば、Lakebase の Change Data Feed がこれらの変更を Unity Catalog で管理する Delta ヒストリテーブルに公開できます。この組み合わせは、Databricks 上でエージェントの入力と下流の分析を既に管理している場合に特に有用です。
長期実行型クラウドエージェントの課題
クラウドエージェントは、それを開始したリクエストや作業者、コンテナ、デプロイメントよりも長く存続することがあります。ユーザーがセッションを開始して翌日に戻り、別の作業者で継続することも可能です。デプロイメントの更新やプロセスの障害は日常茶飯事なので、エージェントの進捗は、それを実行しているプロセスとは独立して維持されなければなりません。回復には、完了した操作の結果と、次に何を行うかを決定するために必要な制御フローの状態の両方が必要です。
この審査エージェントにおいては、以下の 6 つの要件が生まれます:
「回復」では、代替ワーカーが最後に完了したステップから再開する必要があります。
「再試行」では、ツール呼び出しやデータベース操作を繰り返し実行しても、副作用が重複しないように耐性を持たせる必要があります。
「長時間の待機」では、エージェントは人間や外部システムの応答を待つ間、ワーカーを開放したままにしておく必要があります。
「運用上の可視化」では、アプリケーションとオペレーターは、現在のステータス、証拠、再試行の状態、および障害の詳細を確認できる必要があります。
「ランタイムガバナンス」では、ポリシーの更新はコードデプロイなしで即時反映可能であり、アプリケーションがいつオープンケースに新しいポリシーを適用するかを定義する必要があります。
「監査」では、システムは各実行に関連する証拠、ポリシー、推奨事項、および人間の判断を保持する必要があります。
会話のトランスクリプトはこの状態の一部しかカバーしていません。回復には、制御フロー履歴も必要です。つまり、どの操作がスケジュールされたか、どの結果が記録されたか、エージェントが何待ちしているか、そしてどのコマンドを受け入れたかを把握する必要があります。
Temporal は分散システムの管理を簡素化します。Temporal を用いて構築する際、Workflow は 1 つのエージェント実行における永続的な制御フローとなります。
Activity は、モデルやツール、データベースへの呼び出しであり、その結果は Workflow のイベント履歴に記録されます。また、Activity は再試行が可能です。
Signal は、実行中の Workflow に対して送られる非同期コマンドです。例えば、審査担当者の判断などがこれに該当します。
Temporal Lakebase AgentWorkflow の参考実装 は、個人向けローンの審査を行うエージェントの動作可能なサンプルです。このエージェントは複数のツールを呼び出し、規制されたポリシーを読み込み、推奨事項を生成した上で、審査担当者の判断を待ちます。
Lakebase Postgres も開発者のこうした課題を解決しますが、Temporal と Lakebase は異なる消費者のために異なる状態を保持します。Temporal のイベント履歴は再生(replay)の駆動源となり、Lakebase はアプリケーションが直面するビュー——現在のランステータス、メッセージ、証拠、レビュー状態、メトリクスなど——を保存します。Unity Catalog はポリシーのソースとして機能し、同期されたテーブルによって Postgres 内でポリシー照会が可能になります。また、Change Data Feed が運用履歴の戻り経路を提供します。これらのシステムはトランザクションを共有しません。Lakebase の書き込みは、少なくとも一度実行(at-least-once execution)の下で Temporal Activity として実行されます。決定論的な識別子、制約、ガード付き更新、そして Postgres の upsert により、Activity の再試行が同じ論理レコードを確実に対象とします。
このアーキテクチャは、2 つの管理システムとその間の投影契約を追加するものです。これらによってエージェントの回復力とスケーラビリティが向上し、運用オーバーヘッドは低く抑えられます。Temporal と Lakebase の組み合わせが最も効果を発揮するのは、アジェンシーセッションがワーカーの入れ替えに耐えなければならず、長時間の待機後に入力を受け付けられ、関係型状態をアプリケーションに公開でき、かつオープンな状態で統制されたデータ適用が可能である場合です。
与信審査ユースケース
私はローン与信審査を選んだ理由として、同じ実行プロセスで証拠収集、ポリシー適用、推奨結果の生成、そして人の確認待ちという一連のステップを踏む必要があるからです。これらのステップ間のどこかでワーカーが失敗する可能性があります。アプリケーションデプロイなしにポリシーが変更されることもあり、ワークフローが閉じる前には UI で最新の証拠を確認できる必要があります。
簡略化のため、実際の信用情報機関や所得証明機関の代わりにモックデータを用いた申請者データを想定しています。各リクエストにはユーザー ID、申請者 ID、金額、目的、モデル選択、およびターン制限が含まれます。FastAPI が run_id を割り当てて LoanUnderwritingWorkflow を開始し、その ID は API 側、Temporal の実行、そして Lakebase の行全体で共通して使用されます。
最初の処理において、credit_check はスコア、取引枠、延滞状況、および現在の債務を返します。income_verification は収入と雇用証明を返し、debt_to_income_calc は債務対所得比率を計算します。
policy_lookup は、ローン目的に応じたポリシーを読み込み、証拠が承認・紹介・不可の各閾値に合致しているかを評価します。
サンプルの境界線上にある申請者は、信用スコアが 665 で、検証済みの年収は 76,000 ドル、月々の債務返済額は 2,400 ドル、かつ非重大な延滞フラグが一つあります。ポリシーの結果には、すべてのルール、閾値、実際の数値、合格・不合格の判定、情報源、推奨事項、そしてその根拠が記録されます。
このモデルは推奨を行うことはできますが、最終決定を下すことはできません。承認、却下、または追加情報の要求は審査担当者が行います。追加情報の要求は、別のユーザーメッセージとして、そしてエージェントの次のターンとして処理されます。
本ケースでは、ツールの呼び出し完了後にワーカーがクラッシュする事象、コミットされた Lakebase への書き込みだがアクティビティの完了情報が失われる事象、数日間レビューされ続ける未完了の案件、 stale なブラウザによる判断、そして実行中のポリシー変更といった課題を扱います。
図 1:参照実装における実行、運用状態、ガバナンスのパス。
アンダーライターエージェントを実装するには、React と FastAPI が HTTP および UI の処理を担当します。具体的には、実行の開始、証拠のレンダリング、ケースの一覧表示、レビュー判断の提出などを行います。Temporal Cloud はイベント履歴を保存し、タスクを配信します。ワーカーはワークフローコードを再生してモデル、ツール、Lakebase アクティビティを実行します。ネットワークやデータベースの I/O は、決定論的なワークフローコードの外側で処理されます。
FastAPI がワークフローを開始すると実行が開始されます。ワーカーはアクティビティをスケジューリングし、Temporal がその結果を記録します。その後、エージェントは最終的に AWAITING_REVIEW ステータスに到達します。審査担当者の回答はシグナルを通じて返され、承認または却下で実行が終了します。追加情報の要求があれば、エージェントのループは再開されます。
Lakebase には2つの運用スキーマが存在します。agent_ops には、FastAPI が SQL で照会できる実行ステータス、メッセージ、ツール呼び出し、レビュー記録、イベント、メトリクスが含まれています。一方、agent_policy には policy_lookup が使用する読み取り専用で同期されたポリシーが格納されています。各 Activity は、Workflow と同じ決定論的な識別子をキーとしてレコードを書き込むため、再試行後にプロジェクションが追いついても、Lakebase が Temporal のリプレイメカニズムの一部になることはありません。
Unity Catalog は、与信限度額の基準となるソースです。同期されたテーブルを継続的に更新することで、実行中のエージェントがこれらの基準に即時アクセスできます。適用された限度額、収集した証拠、そしてその後の人間による判断は agent_ops に記録されます。Change Data Feed を使用すれば、これらの変更点を Unity Catalog が管理する履歴テーブルへ公開し、監査や分析に活用することが可能です。
ワーカー障害後に完了済みの作業を復元する
Temporal は、別のワーカー上でワークフローの状態を再構築するために必要な順序付けされた イベント履歴 を保持しています。この履歴には、アクティビティのスケジューリングと結果、タイマー、そして シグナル**が含まれます。リプレイでは、これらの記録されたイベントに対してワークフローコードを実行し、現在のターン数、承認されたレビュー判断、トークン使用量、収集した証拠といった変数を再構築します。
リプレイ中は、記録済みのアクティビティ結果が返されるため、再度アクティビティを実行する必要はありません。完了済みの与信チェックは引き続き完了状態として扱われ、記録済みのモデル応答もその実行における応答として保持されます。ワーカー障害時にアクティビティが進行中であり、Temporal がその完了を記録していなかった場合、Temporal は再試行をスケジュールします。エージェントにとっては、イベント履歴に既に記録されたモデル応答が維持されることになります。完了の記録が残っていないモデル呼び出しは、プロバイダー側で処理が完了していたとしても、再度実行されることがあります。
リトライポリシーは個々の操作単位で設定され、コード内で再利用可能です。例では、モデル呼び出しを行うアクティビティは 3 分間のスケジュール開始から完了までのタイムアウト期間内に最大 4 回の試行を許可します。ツール呼び出しのアクティビティは最大 3 回まで、開始から完了までのタイムアウトが 60 秒です。Lakebase のアクティビティは最大 5 回の試行が可能で、タイムアウトは 15 秒に設定されています。
外部効果の反復実行を安全にする
注意すべきリスクの一つとして、Lakebase ツールの結果書き込みがワーカーからアクティビティ完了の報告を受ける前にコミットされてしまうケースがあります。この間に接続が切断されると、Temporal は記録された結果を持たないため、別の試行をスケジュールします。これら 2 つの試行は、論理的には同じ書き込み操作を表しています。
Lakebase の各レコードには、安定した識別子が付与されています。run_id が運用スキーマの基準となり、message_id はメッセージを、tool_call_id はツール呼び出しをそれぞれ特定します。
event_id はマイルストーンを、review_id はレビューラウンドを表します。
そして、decision_id はレビューコマンドです。PostgreSQL の主キーと一意制約がこれらの識別子を保証します。
ツール開始の書き込みでは、安定したアイデンティティとターミナル状態のガードが両方とも示されます。
リトライは同じ tool_call_id を対象とします。最終的な述語は、既存の非終端行のみを「開始中」の状態として書き戻すことを許可します。もしその行がすでに成功または失敗している場合、PostgreSQL は影響を受ける行数をゼロにします。ただし、エラーは発生しません。
呼び出し元は、0 行の結果を必ず確認する必要があります。LakebaseWriteResult は影響を受けた行数を返しますが、現在の Activity ラッパーでは 0 を失敗として扱っていません。本番環境のコードでは、保存された終端状態を確認した上でなければ、0 は期待される何もしない操作(no-op)とみなすべきです。それ以外の場合は、競合を検知して例外を発生させるか記録する必要があります。このルールは、ガード付きの実行やレビュー遷移にも同様に適用されます。
同様のアップサート処理は、メッセージ、ツール結果、イベントに対しても行われます。決定論的な ID により、リトライが同じ論理的行に収束します。一方、各ガード付きの書き込みでは、許可される状態遷移が定義されています。API では、書き込みのリトライ中に一時的に古い状態が表示されることがあります。Activity が成功した後には、承認された行をクエリで取得できるようになります。
すべての副作用を持つツールには、同等の契約が必要です。支払い API は整合性キーを受け付け、メールサービスは呼び出し元が提供するメッセージ ID を受け付け、データベースは一意制約を持ちます。外部システムに重複排除機能がない場合、Activity 側で独自の記録を作成するか、照合プロセスを設ける必要があります。リトライのタイミングは Temporal が決定し、外部システムがそのリトライをどのように処理するかは Activity が決定します。
Expose Current State and Operational Metrics
イベント履歴(Event History)は、実行のセマンティクスとデバッグの詳細情報を提供します。アプリケーションでは、現在のランタイムに対してインデックス付きのリレーショナルクエリが必要です。具体的には、ユーザーとステータスでケースをリストアップしたり、証拠を含む 1 つのトランスクリプトを読み込んだり、特定の人物が待機中のレビューを検索したり、実行全体にわたって測定値を集計したりすることが求められます。
Lakebase は、アプリケーションのビューを正規化された Postgres スキーマとして保存します。agent_runs テーブルには、現在のステータス、ワークフロー ID、リクエスト内容、トークン使用量の合計、タイムスタンプ、推奨事項のメタデータが格納されます。また、agent_messages テーブルでは会話履歴(トランスクリプト)が保持されます。
agent_tool_calls には、引数、ステータス、構造化された結果、エラー、およびタイミングが記録されます。また agent_review_decisions は、推奨事項を安定したレビュー ID、レビュアーの指示、根拠、そして決定時刻と結びつけます。
スキーマは、ワークフロー、ターン、アクティビティ試行の各レベルで名前付きイベントとメトリクスも記録します。FastAPI はこれらのテーブルをバックエンドとして持つ run-detail、workflow-metrics、retry-metrics エンドポイントを公開しています。UI では、証拠収集中の実行、レビュー待ちの実行、失敗したツールの再試行中の実行などを並列で表示できます。オペレーターは SQL を使って同じ行を検索することも可能です。
ワークフローが完了する前に、証拠を確認できます。policy_lookup が終了すると、その構造化された結果はツール呼び出しとともに保存されます。実行が AWAITING_REVIEW ステータスに達すると、審査担当者は信用スコア、DTI(債務対所得比率)、閾値、ルール判定結果、推奨の根拠、および推奨を導き出したポリシーソースを確認できます。
人間のレビューを耐久性高く保ち、陳腐なコマンドを拒否する
モデルが推奨事項を返した際、ワークフローは run_id と現在のターンから review_id を導出します。
レビュー待ちのステータスを Lakebase に書き込み、agent.review_pending イベントを記録し、プロジェクションを AWAITING_REVIEW に設定した上で、workflow.wait_condition を呼び出します。これにより Temporal はワークフローを開いたまま維持しますが、ワーカープロセスは占有されません。
API はアンダーライターのアクションをシグナルとして送信します。送信前に、API は Lakebase 上で実行がレビュー待ち状態にあることと、提出された review_id が現在のラウンドと一致していることを確認します。いずれかのチェックに失敗した場合、API はコンフリクト(競合)エラーを返します。また、ワークフローは独自の状態に基づいてコマンドを独立して検証し、古くなったまたは重複した決定は無視するため、Lakebase の投影が遅れていても実行の保護が維持されます。
シグナルを受け取った後、ワークフローは冪等性を持つ Lakebase アクティビティを通じて決定を永続化します。承認または却下で実行は完了し、追加情報の要求により投影状態は RUNNING に戻され、レビュアーの根拠がユーザーメッセージとして追加されて次のターンが始まります。ターンが変更されたため、次回の推奨事項には新しい review_id が付与されます。
API から 202 レスポンスが返されたことは、Temporal がシグナルを受信したことを確認するものです。ビジネス上の承認処理はワークフロー内で非同期に実行されるため、コマンドが API の事前チェックを通過しても、レビュー状態が変更されていれば無視されることがあります。クライアントは Lakebase による投影を更新し、その結果として生じた状態を確認します。
ワーカーの再デプロイなしで統制されたポリシーを提供する
アンダーライティングの閾値は、ワーカーコードとは独立して変更されます。Unity Catalog 内のソーステーブルには、最小信用スコアや自動承認 DTI(債務対所得比率)、ハードディクラインの閾値、ポリシー名など、特定の目的に特化した値が格納されています。
セットアップスクリプトは、agent_policy.underwriting_policy_limits.policy_lookup という名前の一貫した Lakebase 同期テーブルを作成します。このテーブルは、正規化されたローン目的に基づいて読み取り専用 Postgres コピーを照会します。ポリシー所有者は Unity Catalog のソースを更新し、同期パイプラインがその変更を伝播させます。その後、ワーカーや API をデプロイすることなく、次の実行でその値を読み取ることができます。 (原文の技術表記: agent_policy.underwriting_policy_limits. policy_lookup)
ポリシー結果には、適用されたしきい値、各ルールの実際の値、合格/不合格の結果、およびソース情報が含まれます。Lakebase が無効化されている場合や行が利用できない場合は、デモは固定のポリシー(fixture policy)にフォールバックし、その経路を fixture_fallback として記録します。規制対象のワークフローでは、代わりにフェイルクローズ(安全側に失敗する挙動)を行う可能性があります。このフォールバック判断は、アプリケーション側で明示的に行う必要があります。
Unity Catalog へ運用変更を反映する
このリポジトリでは、Lakebase の Change Data Feed に備えて各 agent_ops テーブルに対して REPLICA IDENTITY FULL を設定します。ただし、この機能を有効にするには管理者がスキーマレベルで手動で許可する必要があります。
その後、Lakebase は PostgreSQL の WAL(Write-Ahead Log)から挿入・更新・削除のイベントをキャプチャし、lb__history という名前の Unity Catalog 管理下の Delta ヒストリテーブルへバッチ処理で書き込みます。
Change Data Feed は現在パブリックプレビュー中で、約 15 秒ごとにデータ変更をフラッシュします。この間隔は監査や分析に適していますが、現在の運用状態については UI が Lakebase に直接クエリを実行しています。
履歴テーブルを使えば、実行ごとのポリシーソース、ツールの証拠、Activity の試行、レビュー待機、推奨事項、そして人間の判断を復元できます。このパスのソーススキーマはリポジトリ設定で定義されていますが、実際に観測されたエンドツーエンドの Change Data Feed 実行データは含まれていません。フィードの有効化と宛先テーブルの確認は、デプロイ後に実施するステップです。
システム運用
本番環境では React/FastAPI と Temporal Worker を分離して展開します。API レプリカはリクエスト負荷に応じてスケールし、Worker は Workflow や Activity Task のバックログと設定された並行度に基づいてスケールします。Lakebase の自動スケーリング では、プロジェクトの範囲内でデータベース計算リソースを調整できます。
Temporal Cloud の料金体系 は、アクション数とアクティブなイベント履歴ストレージ、保持中のイベント履歴ストレージに基づいて算出されます。そのため、再試行頻度や長期にわたる履歴の保有状況もコストに影響します。チーム側では、Kubernetes レプリカ、Task Queues、コネクションプールの制限、データベースのバウンドを workload に合わせて設定する必要があります。あるいは、最新のオープンソースリリースを使用して、自前で Temporal Service を構築することも可能です。
Lakebase クライアントは、OAuth マシン間認証を使用します。Databricks の OAuth トークンと生成されたデータベース認証情報は有効期限を持つため、クライアントは 1 時間後に期限切れとなるデータベース認証情報の前に SQLAlchemy の接続プールを自動的に更新します。すべての接続には TLS が使用されています。認証情報のローテーションを行わない場合、長時間稼働するワーカーは予測可能なスケジュールでデータベースエラーに直面することになります。
オペレーターは、ワークフローとアクティビティの履歴を確認するために Temporal を使い、アプリケーションの状態やメトリクスをクエリするには Lakebase を利用し、プロセスやデプロイメントの健全性をチェックするには Kubernetes を用います。これにより、意図的なレビュー待機状態と、アクティビティのリトライ、データベースアクセスの失敗、あるいはツールの失敗を明確に区別することが可能になります。
検証結果と制限事項
テストスイートには、ワークフローシーケンス、レビュー動作、OAuth 接続の構築、冪等性のある永続化、メトリクス契約、API ワークフローの開始、ワーカー設定に関する 21 の合格テストが含まれています。クラッシュ回復スクリプトでは、決定論的プロバイダーを用いたプロセス障害のシナリオも追加されています。
出願者とプロバイダーのデータは固定値(フィクスチャ)です。本リポジトリは、貸付モデルの検証、規制遵守状況、本番環境におけるセキュリティ制御、地域ごとの可用性、あるいは大規模運用時のパフォーマンスについては保証しません。ローカルでのクラッシュ実験では Lakebase を無効化して実行したため、これは Temporal の回復機能のみを隔離して検証するものです。Change Data Feed については、対象の Databricks 環境で有効化し、確認を行う必要があります。
次のステップ
もっと詳しく知りたい場合は、デモを実行して耐久性のある実行の仕組みを直接確認してみましょう。リファレンス実装 を実行し、ワーカーが動作中に停止した状態からエージェントがどのように回復するかを観察してください。また、Lakebase に接続することで、各意思決定の背後にある証拠やポリシー、人間によるレビューの状態を探索できます。
関連記事
今日のまとめ
AIデイリーブリーフで今日の重要ニュースをまとめ読み