Skip to content

DATABRICKS

Read original ↗

Databricks — Build durable agents with Temporal and Lakebase

Summary

A Databricks reference implementation shows how to build a durable long-running AI agent — a personal-loan underwriting agent that gathers evidence, applies governed policy, produces a recommendation, and then may wait days for a human underwriter — by pairing Temporal for durable execution with Lakebase Postgres for queryable operational state. The central design idea is a two-store split with a projection contract: Temporal's Event History is the replay-driving control-flow record, and Lakebase holds the application-facing view (run status, transcript, evidence, review state, metrics). The two systems do not share a transaction — Lakebase writes run as Temporal at-least-once Activities, and idempotency is enforced through deterministic identifiers, Postgres primary keys/unique constraints, and guarded state-transition writes (compare-and-set). Governed underwriting policy flows in from Unity Catalog via a continuous synced table, and Lakebase Change Data Feed provides the return path back to Unity Catalog Delta history tables for audit.

Key takeaways

  1. A cloud agent outlives the process that started it, so progress must survive independently of the executing process. Recovery needs both the results of completed operations and the control-flow state (which operations were scheduled, which results recorded, what the agent is waiting for). A conversation transcript alone is insufficient — the Event History carries the replay-critical control flow. (Source: this article.)
  2. Temporal vocabulary maps cleanly onto agents. A Workflow is the durable control flow for one agent run; an Activity is a call to a model, tool, or database whose result is recorded in Event History and can be retried; a Signal is an async command sent to a running Workflow (e.g. the underwriter's decision). Replay returns recorded Activity results instead of re-running them, so a completed credit check stays completed and a recorded model response stays fixed. (Source: this article — durable execution.)
  3. Retry Policies are set per operation. Model-calling Activities: up to 4 attempts within a 3-minute schedule-to-close timeout. Tool-calling Activities: up to 3 attempts, 60-second start-to-close timeout. Lakebase Activities: up to 5 attempts, 15-second start-to-close timeout.
  4. The Temporal↔Lakebase boundary is at-least-once, not transactional, so every external effect must be safe to repeat. A Lakebase write can commit before the Worker reports Activity completion; if the connection drops in that gap, Temporal has no recorded result and schedules another attempt — both attempts are the same logical write. The fix: stable identities (run_id, message_id, tool_call_id, event_id, review_id, decision_id) enforced by Postgres PKs/unique constraints, plus guarded upserts whose predicate only allows a nonterminal row to transition (a retry against an already-terminal row affects zero rows and raises no error). (concepts/idempotent-operations, patterns/conditional-write.)
  5. Guarded zero-row results must be classified, not swallowed. The sample's LakebaseWriteResult returns affected row count but the Activity wrapper does not turn zero into a failure; the article explicitly flags that production code should classify zero as an expected no-op only after confirming the stored terminal state, otherwise raise or record a conflict. (Caveat from the article.)
  6. Lakebase is the async-projected read model. Event History gives execution/debug semantics; the application needs indexed relational queries — list cases by user+status, load one transcript with evidence, find reviews waiting on a person, aggregate metrics. Lakebase stores this in a normalized Postgres schema (agent_runs, agent_messages, agent_tool_calls, agent_review_decisions, plus Events and metrics). The API can briefly show older state while a write retries, then the accepted row becomes queryable. (patterns/async-projected-read-model.)
  7. Human review is durable and stale-command-resistant. On a recommendation, the Workflow derives review_id from run_id+turn, writes a pending review, sets projection to AWAITING_REVIEW, and calls workflow.wait_condition — Temporal keeps the open Workflow without holding a Worker. The API precheck validates Lakebase shows the run awaiting review and the submitted review_id matches the current round; the Workflow independently validates the command and ignores stale or duplicate decisions even when the Lakebase projection lags. A 202 from the API confirms only that Temporal received the Signal — business acceptance is async. (concepts/human-in-the-loop.)
  8. Governed policy is served without redeploying Workers. Underwriting thresholds live in Unity Catalog; a continuous Lakebase synced table (agent_policy.underwriting_policy_limits) makes them a read-only Postgres copy that policy_lookup queries by normalized loan purpose. Policy owners edit the Unity Catalog source, the sync pipeline propagates, and a later run reads the new value with no Worker or API deploy. (concepts/policy-as-data.) The demo falls back to fixture policy (recorded as fixture_fallback) when Lakebase is disabled; a regulated Workflow might instead fail closed — the application must make that choice explicitly.
  9. Change Data Feed is the audit return path. Each agent_ops table is prepared with REPLICA IDENTITY FULL; once an admin enables the feature, Lakebase captures inserts/updates/deletes from the Postgres WAL and batches them (~every 15 s) into Unity Catalog-managed Delta history tables named lb_<table>_history, reconstructing each run's policy source, evidence, attempts, review wait, recommendation, and human decision. (concepts/change-data-capture, concepts/audit-trail.)
  10. Operational model separates the tiers. React/FastAPI scales with request load; Temporal Workers scale with Workflow/Activity Task backlog; Lakebase autoscaling adjusts DB compute. The Lakebase client uses OAuth M2M auth and refreshes its SQLAlchemy connection pool before the one-hour DB credential expires — without rotation a long-running Worker hits DB failures on a predictable schedule. (concepts/oauth-token-lifecycle.)

Systems / concepts / patterns extracted

Operational numbers

  • Model Activities: 4 attempts / 3-min schedule-to-close.
  • Tool Activities: 3 attempts / 60-s start-to-close.
  • Lakebase Activities: 5 attempts / 15-s start-to-close.
  • Change Data Feed flush cadence: ~15 s (Public Preview at time of writing).
  • DB credential lifetime: 1 hour (pool refreshed before expiry).
  • Sample borderline applicant: 665 credit score, $76,000 verified annual income, $2,400 monthly debt, 1 non-material delinquency flag.
  • Test suite: 21 passing tests (Workflow sequencing, review behavior, OAuth connection construction, idempotent persistence, metrics contracts, API Workflow start, Worker settings); plus a crash-recovery script.

Caveats

  • Applicant/provider data are fixtures; the repo does not validate lending models, regulatory compliance, production security controls, regional availability, or performance at scale.
  • The local crash exercise ran with Lakebase disabled, so it isolates Temporal recovery only.
  • Change Data Feed still requires manual enablement + verification in the target environment; the repo configures the source schema but does not include an observed end-to-end CDF run.
  • The guarded-write wrapper does not classify zero-row results as failures — production code must add that (see takeaway 5).

Source

Last updated · 766 distilled / 2,225 read