ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Day3 Agent异步化

Day3 Agent异步化 Day 1 的目标是跑通会议文本 → 提取 → 人审 → 本地任务草稿的最小闭环Day 2 又给 LLM、ASR 和 Tool 加上了统一的 Harness 防护层。但当时图还是同步执行同步节点靠run_sync()去桥接异步 Harness图还能在内部兜底创建全局 HarnessPostgreSQL checkpoint 也还没有证明能跨进程恢复。它们在 Demo 阶段不明显到了 API 服务、Worker 重启和真实外部依赖场景就会变成难定位的问题。Day 3 的目标不是增加功能而是把这条链路改得可替换、可追踪、可恢复。结论本次改造完成了四件事LangGraph、Node 与 Harness 统一进入异步调用链main.py成为唯一生产组合根Graph 与 Node 不再自行寻找或创建 Harness每次外部调用都能用thread_id、meeting_id、request_id、tenant和caller_node关联审计日志不记录会议正文使用真实 PostgreSQL验证进程 A 在人审中断后退出进程 B 用同一个thread_id恢复且不会重复执行 parse 或重复生成本地草稿。改造后的调用链main.py组合根 → 创建 Harness → 创建 AsyncPostgresSaver → build_graph(harness..., checkpointer...) → app.astream(...) → async asr_node / parse_node → await harness.execute_asr / execute_llm → Provider timeout retry breaker audit这里的“异步”不代表五个节点会同时跑。会议输入、ASR、解析、风险审核和任务草稿仍然是顺序流程它解决的是每一层都运行在同一个事件循环中取消、超时、重试和追踪不再穿过同步桥接层。1. 为什么要从同步图改成真正的异步图改造前的调用链是app.stream() → 同步 Node → run_sync() → 新事件循环或辅助线程 → await Harnessrun_sync()本身并不是错误它适合给旧的同步调用方做兼容。但 Graph、Node 和 Harness 都已经要面对超时、取消与恢复时再跨一层事件循环会让问题变得很绕例如超时究竟发生在 Node、桥接线程还是 Provider取消能否传到真正的请求后续接入 FastAPI 时是否又嵌套了一层 event loop。现在入口和节点统一改为 async# main.pyasyncdefmain()-None:harnessget_harness()checkpointerawaitget_checkpointer()appbuild_graph(harnessharness,checkpointercheckpointer)asyncforeventinapp.astream(input_state,config):print(event)# parse_node.py简化asyncdefparse_node(state,config,harness):returnawaitharness.execute_llm(prompt...,caller_nodeparse_node,thread_idconfig[configurable][thread_id],meeting_idstate[meeting_id],request_idstate[request_id],)Harness 仍然是统一入口所以 timeout、retry、circuit breaker、Provider failover 与 audit 的职责没有散落回 Node。变化只是Node 不再把协程“转回同步”而是直接await。2. 为什么 Node 不应该调用全局get_harness()全局单例的问题不是“它一定不能用”而是依赖被隐藏了。只看parse_node(state)时看不出它会调用模型、读取哪份配置也很难在测试里替换成一个确定的假模型。此前 Graph 的行为类似active_harnessharnessorget_harness()调用方忘记注入 Harness 时Graph 会悄悄创建生产 Harness。单测可能因此碰到真实配置生产环境也会把装配错误延后到运行中才暴露。现在 Graph 的函数签名明确要求两个依赖defbuild_graph(*,harness:AgentHarness,checkpointer:Any):...main.py是唯一的生产组合根它创建 Harness 和 Checkpointer再交给 GraphGraph 通过闭包把 Harness 传给真正需要外部调用的asr_node、parse_node。Node 不导入 Provider bootstrap不调用get_harness()也不直接调用 SDK。测试时则注入最小 FakeHarnessharnessRecordingFakeHarness()appbuild_graph(harnessharness,checkpointerMemorySaver())asyncfor_inapp.astream(audio_input,config):passassert[call[operation]forcallinharness.calls][asr,llm]这个测试没有联网也不会调用真实 ASR 或 LLM它只检查图是否使用了传入的 Harness以及节点是否传了正确的关联字段。依赖从“运行时猜出来”变成“构图时一眼可见”。3. PostgreSQL checkpoint 为什么必须做跨进程恢复测试MemorySaver很适合普通单测但它只存在于当前 Python 进程内。同一进程里调用get_state()只能证明内存里的对象还在不能证明 Worker 重启、进程崩溃或部署滚动更新后可以恢复。Day 3 的集成测试使用AsyncPostgresSaver和AsyncConnectionPool。同步版PostgresSaver的异步 checkpoint 方法不能满足astream()与aget_state()的调用方式因此不能拿来冒充异步恢复验收。测试将自己作为两个独立 Python 子进程启动进程 A thread_idrecovery-*** → input → parse(1) → review → interrupt → checkpoint 写入 PostgreSQL → 退出 进程 B全新 Graph / Harness / Checkpointer → 使用同一个 thread_id 调用 aget_state() → 读取 A 的 decision、todo 与 interrupt → Command(resumeapprove) → create → END验收不只看“B 能继续”还断言A 中parse_calls 1B 恢复时parse_calls 0不会重复调用 LLM 解析B 最终只生成 1 条本地任务草稿A、B 使用的是不同 Python 进程、不同 Graph/Harness/Checkpointer 对象但使用同一个 PostgreSQL 和thread_id。在 Windows 上psycopg的异步连接不支持默认 Proactor event loop因此入口使用 Selector event loop 策略。这是运行环境兼容处理不是业务流程的一部分。集成测试默认跳过避免没有数据库的日常回归失败验收时必须显式开启$env:RUN_POSTGRES_INTEGRATION 1.\venv\Scripts\python.exe-m pytest-q-m integration-s这次真实执行结果为1 passed。这比保存一张单次终端截图更有意义同一命令可以在未来重复运行持续验证“恢复不会重跑已完成节点”。4. 为什么审计字段不能只留一个request_id一次外部调用至少有四个不同维度字段它回答的问题thread_id这是哪条可恢复的工作流执行meeting_id这是哪场业务会议request_id这是哪一次端到端请求tenant属于哪个租户或组织caller_node是哪个 Node 发起的调用thread_id不能替代meeting_id。同一场会议可能有重试、重新处理或多个恢复线程反过来同一工作流线程也不应该被误当成业务会议的永久 ID。这些字段由 Node 传给 HarnessHarness 将它们从通用kwargs中取出写入ExecutionRequestExecutor 在选择 Provider 后将它们传给 PipelinePipeline 创建只包含安全关联字段的HookContext。此外LLM 的完整messages其中包含 system prompt 和会议正文被单独保存在请求对象中只给 Provider 调用模型使用。它不进入审计上下文。Console Audit 只输出白名单字段{caller_node:parse_node,event:after,latency_ms:12.8,meeting_id:meeting-test,provider:fake:fake-model,request_id:request-test,status:success,target:llm:fake:fake-model,tenant:tenant-test,thread_id:thread-test}其中target是受保护资源标签例如llm:openai-us:gpt-4o用于区分 timeout、retry 和 breaker 作用的对象provider是本次实际尝试的 Provider例如openai-us:gpt-4ostatus表示成功、失败或熔断latency_ms是本次执行耗时。它们都是结构化字段不是 API Key、Prompt、会议正文、音频 URL 或数据库连接串。Provider ID 与 API Key 在配置中是两个不同字段审计日志只记录前者。对应测试会构造“绝密会议正文”作为输入再断言HookContext包含关联字段、资源、Provider、状态和耗时同时序列化后的审计记录不含这段正文。验证结果与当前边界验收项真实结果验证方式异步图与 Node通过app.astream() FakeHarness 测试Harness 显式注入通过Graph 缺少依赖时构图失败节点边界守卫测试通过审计关联字段通过test_audit_correlation.py断言字段与正文隔离PostgreSQL 跨进程恢复通过RUN_POSTGRES_INTEGRATION1的 A/B 子进程测试1 passed默认单元回归通过11 passed, 1 skipped开启集成测试的完整回归通过12 passed当前还有两个明确边界create_node生成的是本地演示任务草稿不代表 Jira、飞书或其他外部系统已经创建成功。未来真实写操作必须走await harness.execute_tool(..., idempotency_key...)。本文没有把单次人工运行日志当成证据文件可重复运行的 A/B 集成测试才是恢复能力的主要验证手段。这次改造带来的好处调用链更简单Graph、Node、Harness 都在同一个 async 事件循环中不再依赖同步桥接。依赖更透明生产 Harness 只在组合根创建测试可以稳定注入 FakeHarness。恢复更可信不是“同一进程内还能读状态”而是新进程可以从 PostgreSQL 恢复并避免重复 parse。排障更可控可以按thread_id、request_id、Provider、状态和耗时定位一次调用而无需记录会议正文。为真实 Tool 铺路先解决恢复与关联问题再接入有幂等要求的外部写操作避免网络超时造成重复创建。参考资料LangGraph InterruptsLangGraph PersistencePsycopg async connections
返回列表