返回
RSS Databricks Blog AI 逐段翻译 发布 2026-09-08 23:46 收录于 09-09

用Temporal和Lakebase构建持久化数据Agent的参考实现

DataHot 速览

Databricks 展示了如何用 Temporal 与 Lakebase 构建持久化数据 Agent,并以个人贷款承销为例。该 Agent 可跨工作进程和容器存活,处理重试、长时间等待和故障恢复。Temporal 负责持久化控制流状态,Lakebase Postgres 存储证据与决策,并通过 Change Data Feed 与 Unity Catalog 同步。整个方案满足了恢复、重试、长等待、运维可见、运行时治理和审计六项要求。

为什么值得关注:对构建可恢复、可治理的生产级 Data Agent 给出了具体架构与实现思路,值得做 Agent 工程化的团队参考。

本文目录 12 节
  1. 长时间运行的云代理面临的挑战
  2. 核保用例
  3. 架构
  4. 在工作节点故障后恢复已完成的工作
  5. 让外部效应安全地重复
  6. 暴露当前状态和操作指标
  7. 保持人工审核的持续性并拒绝过时命令
  8. 无需重新部署Worker即可提供受治理的政策
  9. 将操作更改返回到Unity Catalog
  10. 操作该系统
  11. 证据和限制
  12. 后续步骤

译文

AI 逐段翻译

个人贷款核保代理收集证据、应用政策,并可能等待审核员数天。在此期间,工作进程可能重启,工具调用可能失败。应用程序必须保留已完成的工作、恢复执行,并确保证据可供审核员使用。

参考实现使用Temporal进行持久执行,使用Lakebase Postgres存储可查询的操作状态。来自Unity Catalog的核保政策通过同步表在Lakebase中可用。Temporal活动将证据、决策和指标写入Lakebase;启用后,Lakebase变更数据流可将这些更改发布到Unity Catalog管理的Delta历史表。当Databricks已管理代理的输入和下游分析时,此组合尤为有用。

长时间运行的云代理面临的挑战

云代理可能比启动它的请求、工作进程、容器或部署存活更久。用户可能开始一个会话,第二天返回,并在另一台工作进程上继续。部署和进程失败是常规情况,因此代理的进展必须独立于执行它的进程而持久。恢复需要已完成操作的结果以及确定下一步的控制流状态。

对于此核保代理,这产生了六个要求:

  1. 恢复: 替代工作进程必须从最后完成的步骤恢复。
  2. 重试: 工具调用和数据库操作必须容忍重复执行而不会重复产生副作用。
  3. 长时间等待: 代理必须等待人员或外部系统而无需保持工作进程开放。
  4. 操作可见性: 应用程序和操作员需要当前状态、证据、重试状态和失败详情。
  5. 运行时治理: 政策更新必须无需代码部署即可可用,且应用程序必须定义开放案例何时采纳它们。
  6. 审计: 系统必须保留每次运行相关的证据、政策、建议和人工决策。

对话记录仅覆盖此状态的一部分。恢复还需要控制流历史:已调度哪些操作,记录了哪些结果,代理在等待什么,以及已接受哪些命令。

Temporal简化了分布式系统的管理。使用Temporal构建时,工作流是单个代理运行的持久控制流。活动是对模型、工具或数据库的调用,其结果记录在工作流的事件历史中;活动可以重试。信号是发送到正在运行的工作流的异步命令,例如核保员的决定。Temporal Lakebase AgentWorkflow参考实现是一个可运行的个人贷款核保代理。它调用多个工具,读取治理政策,产生建议,并等待核保员。

Lakebase Postgres也帮助开发者管理这些问题,但Temporal和Lakebase为不同的消费者存储不同的状态。Temporal的事件历史驱动重放。Lakebase存储面向应用程序的视图:当前运行状态、消息、证据、审核状态和指标。Unity Catalog仍然是政策源;同步表使该政策在Postgres中可查询,而变更数据流为操作历史提供返回路径。系统不共享事务。Lakebase写入作为Temporal活动以至少一次执行方式运行。确定性标识符、约束、条件更新和Postgres上插确保重复的活动尝试针对相同的逻辑记录。

此架构添加了两个受管系统及其之间的投影契约。它们共同提高了代理的弹性和可扩展性,同时保持较低的操作开销。当代理会话必须在工作进程替换后存活、在长时间等待后接受输入、向应用程序暴露关系状态并在开放时应用治理数据时,Temporal加Lakebase最为有用。

核保用例

我选择贷款核保是因为同一次运行必须收集证据、应用政策、产生建议并等待人。任何步骤之间工作进程都可能失败。政策可以在不部署应用程序的情况下更改,而UI在工作流关闭前需要当前证据。

模拟申请人代替真实的信用机构和收入提供商,工具序列是确定性的以简化。每个请求包含用户ID、申请人ID、金额、目的、模型选择和回合限制。FastAPI分配run_id,启动LoanUnderwritingWorkflow,并在API、Temporal执行和Lakebase行中使用相同的ID。

在第一轮中,credit_check返回分数、信用额度、违约情况和当前债务。income_verification返回收入和就业证据。debt_to_income_calc计算债务收入比。policy_lookup加载贷款目的的政策,并评估证据是否符合批准、转介和硬拒绝阈值。

样本中的临界申请人的信用评分为665,验证年收入为76,000美元,月债务为2,400美元,并有一个非实质性违约标记。政策结果记录每个规则、阈值、实际值、通过/失败结果、来源、建议和理由。模型可以建议但不能决定。核保员批准、拒绝或请求更多信息。请求更多信息成为另一个用户消息和另一个代理回合。该案例测试了已完成工具调用后的工作进程崩溃、活动完成丢失但Lakebase写入已提交、审核保持开放数天、过期浏览器决策以及执行过程中的政策变更。

架构

image1.jpg

为实现承保代理,React和FastAPI负责HTTP和UI工作:启动运行、渲染证据、列出案例以及提交审核决策。Temporal Cloud存储事件历史并调度任务。工作节点重放工作流代码并执行模型、工具和Lakebase活动;网络和数据库I/O保持在确定性工作流代码之外。

当FastAPI启动一个工作流时,运行开始。工作节点调度活动,Temporal记录其结果,代理最终达到AWAITING_REVIEW。承保人的响应通过信号返回。批准或拒绝关闭运行;请求更多信息则恢复代理循环。

Lakebase持有两个操作模式。agent_ops包含运行状态、消息、工具调用、审核记录、事件和指标,FastAPI可以用SQL查询这些内容。agent_policy包含被policy_lookup使用的只读同步策略。每个活动写入由工作流使用的相同确定性标识符键控的记录,因此投影可以在重试后赶上,而不会使Lakebase成为Temporal重放机制的一部分。

Unity Catalog是承保阈值的来源。一个持续同步的表使它们可用于运行中的代理。应用的阈值、证据和随后的人工决策被写入agent_ops。变更数据馈送可以将这些更改发布到由Unity Catalog管理的历史表中,用于审计和分析。

在工作节点故障后恢复已完成的工作

Temporal保留有序的事件历史,这是在另一个工作节点上重建工作流状态所必需的。该历史包括活动调度和结果、定时器以及 信号。重放根据这些记录的事件运行工作流代码,并重建变量,如当前轮次、接受的审核决策、令牌使用量和收集的证据。

重放期间返回记录的活动结果,而不是再次运行该活动。已完成的信用检查保持完成,已记录的模型响应保持为该执行的结果。如果某个活动在工作节点故障时正在运行,并且Temporal从未记录其完成,则Temporal可以调度另一次尝试。对于代理,这保留了事件历史中已记录的模型响应。如果模型调用的完成未被记录,即使提供商已处理完,该调用仍可能再次运行。

重试策略在单个操作的粒度上分配,并且可以在代码中重用。在示例中,模型调用活动允许在三分内的调度到关闭超时时间内最多尝试四次。工具调用活动允许最多三次尝试,并有60秒的开始到关闭超时。Lakebase活动允许最多五次尝试,并有15秒的开始到关闭超时。

让外部效应安全地重复

一个风险是,Lakebase工具结果写入可以在工作节点报告活动完成之前提交。如果连接在该间隙中断,Temporal没有记录结果,因此调度另一次尝试。两次尝试代表相同的逻辑写入。

每个Lakebase记录都有稳定的标识。run_id锚定操作模式。message_id标识消息,tool_call_id标识工具调用,event_id标识里程碑,review_id标识审核轮次,decision_id标识审核人命令。Postgres主键和唯一约束强制实施这些标识。

工具启动写入显示了稳定的标识和终端状态保护:

重试针对相同的tool_call_id。最终谓词只允许现有的非终端行被写回已启动状态。如果该行已经成功或失败,PostgreSQL影响零行。它不会引发错误。

调用方必须检查零行结果。LakebaseWriteResult返回受影响的行数,但当前的Active包装器不会将零转换为失败。生产代码应在确认存储的终端状态后,将零分类为预期的无操作;否则,应引发或记录冲突。相同的规则适用于受保护的运行和审核转换。

类似的更新插入涵盖消息、工具结果和事件。确定性ID使重试收敛到相同的逻辑行,而每个受保护的写入定义了哪些状态转换是合法的。API可以在写入重试时短暂显示较旧的状态。在活动成功后,接受的行可以查询。

每个有副作用的工具都需要等效的契约。支付API可能接受幂等键,电子邮件服务可能接受调用方提供的消息ID,数据库可能接受唯一约束。如果外部系统不提供去重机制,活动需要自己的记录或对账流程。Temporal决定何时重试。活动决定外部系统如何处理该重试。

暴露当前状态和操作指标

事件历史提供执行语义和调试细节。应用程序需要对当前运行的索引关系查询:按用户和状态列出案例,加载带有证据的一个聊天记录,查找等待人的审核,并跨执行聚合测量。

Lakebase在规范化的Postgres模式中存储该应用程序视图。agent_runs持有当前状态、工作流ID、请求、令牌总数、时间戳和推荐元数据。agent_messages保留聊天记录。agent_tool_calls记录参数、状态、结构化结果、错误和计时。agent_review_decisions将推荐连接到稳定的审核ID、审核人命令、理由和决策时间。

该模式还在工作流、轮次和活动尝试级别记录命名的事件和指标。FastAPI暴露run-detailworkflow-metricsretry-metrics端点,由这些表支持。UI可以显示一个运行收集证据,另一个等待审核,第三个重试失败的工具。操作员可以用SQL查询相同的行。

在工作流完成之前,证据是可用的。之后policy_lookup完成时,其结构化结果会与工具调用一起存储。当运行达到AWAITING_REVIEW状态时,承保人可以看到信用评分、DTI、阈值、规则结果、理由以及生成建议的政策来源。

保持人工审核的持续性并拒绝过时命令

当模型返回建议时,工作流从run_id和当前轮次派生review_id。它将待审核写入Lakebase,记录一个agent.review_pending事件,将投影设置为AWAITING_REVIEW,并调用workflow.wait_condition。Temporal保留打开的工作流,而不会占用Worker进程。

API将承保人的操作作为Signal发送。在发送之前,API检查Lakebase显示运行处于等待审核状态,并且提交的review_id与当前轮次匹配。如果任一检查失败,API返回冲突。工作流根据自身状态独立验证命令,并忽略过时或重复的决定,即使Lakebase投影滞后也能保护执行。

接受Signal后,工作流通过幂等的Lakebase Activity持久化决定。批准或拒绝完成运行。请求更多信息会将投影更改回RUNNING,将审核人的理由作为用户消息附加,并开始下一轮。由于轮次已更改,下一个建议将获得新的review_id

API的202响应确认Temporal已接收Signal。业务接受在工作流中异步发生,因此命令可能通过API预检查,但如果审核状态已更改,仍可能被忽略。客户端刷新Lakebase投影以观察最终状态。

无需重新部署Worker即可提供受治理的政策

承保阈值独立于Worker代码发生变化。Unity Catalog中的源表包含特定于用途的值,例如最低信用评分、自动批准DTI、硬拒绝阈值和政策名称。

设置脚本创建一个名为agent_policy.underwriting_policy_limits. policy_lookup的连续Lakebase同步表,该表按标准化贷款用途查询此只读Postgres副本。政策所有者更新Unity Catalog源;同步管道传播更改,后续运行在无需Worker或API部署的情况下读取它。

政策结果包含应用的阈值、每条规则的实际值和通过/失败结果以及来源。当Lakebase被禁用或行不可用时,演示可以回退到夹具政策,并将该路径记录为fixture_fallback。受监管的工作流可能会改为失败关闭。应用程序必须明确做出该回退决定。

将操作更改返回到Unity Catalog

该存储库通过设置agent_ops表为Lakebase变更数据馈送做好准备REPLICA IDENTITY FULL。管理员仍需为架构启用该功能。然后Lakebase从Postgres预写日志捕获插入、更新和删除,并将它们批量写入Unity Catalog托管的Delta历史表,这些表以lb_<table>_history模式命名。

变更数据馈送当前处于公开预览阶段,大约每15秒刷新更改。该间隔适用于审计和分析,而UI直接查询Lakebase以获取当前操作状态。

历史表可以重建运行的政策来源、工具证据、Activity尝试、审核等待、建议和人工决定。存储库为源架构配置了此路径,但未包含观察到的端到端变更数据馈送运行。启用馈送并验证目标表仍然是部署步骤。

操作该系统

该部署将React/FastAPI与Temporal Worker分开。API副本随请求负载扩展;Worker随工作流和Activity任务积压以及配置的并发性扩展。Lakebase自动扩展在项目边界内调整数据库计算。

Temporal Cloud定价基于Actions加上活动和保留的事件历史存储,因此重试频率和长时间打开的历史也会影响成本。团队仍然需要为其工作负载设置Kubernetes副本、任务队列、连接池限制和数据库边界。或者,您可以使用最新的开源版本搭建自己的开源Temporal Service。

Lakebase客户端使用OAuth机器对机器认证。Databricks OAuth令牌和生成的数据库凭据会过期,因此客户端会在数据库凭据过期前一小时刷新其SQLAlchemy连接池。连接使用TLS。如果没有轮换,长时间运行的Worker将按可预测的时间表遇到数据库故障。

操作员使用Temporal检查工作流和Activity历史,使用Lakebase查询应用程序状态和指标,并使用Kubernetes检查进程和部署健康状况。操作员可以区分故意的审核等待与Activity重试、数据库访问失败或工具失败。

证据和限制

测试套件包含21个通过测试,涵盖工作流顺序、审核行为、OAuth连接构建、幂等持久性、指标契约、API工作流启动和Worker设置。崩溃恢复脚本添加了一个使用确定性提供程序的进程故障演练。

申请人和提供者数据是夹具。该存储库不验证贷款模型、法规遵从性、生产安全控制、区域可用性或规模性能。本地崩溃演练在Lakebase被禁用的情况下运行,因此它隔离了Temporal恢复。变更数据馈送仍需在目标Databricks环境中启用和验证。

后续步骤

想了解更多?尝试自己运行演示,看看持久执行的实际效果。运行参考实现,在运行中停止Worker,观察代理恢复。连接Lakebase以探索每个决定背后的证据、政策和人工审核状态。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存