拆掉Agent的while循环:Durable Execution架构实践

发布时间:2026/10/4 6:57:22
拆掉Agent的while循环:Durable Execution架构实践 1. 项目概述当 Agent 不再“死磕” while 循环我把 Agent 的 while 循环拆掉了——这句话刚在内部技术分享会上说出来底下就有同事下意识摸了摸键盘仿佛在确认自己的 IDE 还没崩。不是夸张而是真实反应。过去三年我带团队落地了 7 个生产级 AI Agent 系统从金融风控辅助决策到工业设备故障预判所有早期版本的主控逻辑里都嵌着一个看似朴素、实则危险的while True:块。它像老式电风扇的调速旋钮转得越久电机越烫一旦卡住整台机器就停摆。我们曾为一个日均调用量 230 万次的客服意图路由 Agent在凌晨三点紧急回滚只因那个 while 循环在某次模型响应延迟突增时把线程池耗尽连带压垮了下游的 Redis 缓存集群。这不是理论风险是血淋淋的 SRE 报告里的第 17 条事故根因。这个标题说的“拆掉”不是简单删掉几行代码而是对 Agent 架构范式的重新锚定把“持续运行”的责任从单个进程的循环体内移交到系统层的持久化执行引擎上。核心关键词——Agent、while 循环、架构、Durable Execution、微服务——在这里不是并列关系而是一条因果链因为要构建真正可用的 Agent不是 Demo就必须解决 while 循环带来的状态脆弱性而解决它必然导向 Durable Execution 这一设计原则最终这又天然倒逼整个系统向微服务化演进。它不依赖 MATLAB OOP 那种面向对象的算法封装思想也不需要 Spring Cloud 那套复杂的分布式事务协调它更接近 Linux 内核处理中断的方式——把“等待”这件事交给调度器而不是让进程自己傻等。适合谁看如果你正被这些问题困扰Agent 在长流程中偶发丢失上下文、重试后状态错乱想加个“暂停-恢复”功能却无从下手并发量一上来Agent 实例就内存暴涨甚至 OOM或者你只是好奇为什么开源社区里最火的 Agent 框架如 LangChain 的早期版本、AutoGen 的默认模式默认都带着 while 循环而生产环境的大厂却几乎没人这么用——那这篇就是为你写的。它不讲抽象理论只讲我在铜陵学院李光耀老师那套“基于 MATLAB OOP 架构的多算法融合数字图像处理系统”里学到的教训再精妙的算法如果跑在不牢靠的底盘上结果就是一张模糊的图。2. 架构设计与思路拆解为什么 while 循环是“伪实时”的陷阱2.1 传统 while 循环 Agent 的三重幻觉我们先直面那个被无数教程和 Demo 捧上神坛的while True:结构def run_agent(): state initialize_state() while True: user_input get_user_input() # 阻塞等待 state update_state(state, user_input) action decide_action(state) result execute_action(action) state update_state(state, result) send_response(result)初看很美逻辑清晰、控制流明确、像教科书一样“实时”。但这是三层精心包装的幻觉。第一层幻觉它很“实时”。真相是get_user_input()这一行代码在 Web 场景下通常是 HTTP 请求的阻塞读取在 CLI 场景下是 stdin 的阻塞读取。它根本不是“实时”而是“同步等待”。一旦上游用户、API 网关网络抖动这个循环就卡在 I/O 上整个 Agent 实例等于“假死”。而真正的实时系统比如电信信令网或高频交易系统其“实时”指的是确定性的响应时间上限如 10ms而非“永不阻塞”。while True提供的只是“永不退出”离“实时”差了十万八千里。第二层幻觉它很“可控”。开发者以为break就能优雅退出continue就能跳过异常。但现实是当execute_action()调用一个外部 API而该 API 因为熔断策略返回 429Too Many Requests时state已经被update_state修改了一半。此时continue只是让循环继续带着一个半残缺的状态进入下一轮结果就是“用户问‘查订单’Agent 却开始执行‘取消订单’”。我在一个物流跟踪 Agent 里亲眼见过这种 bug一次网络超时导致状态机从WAITING_FOR_TRACKING_NUM错误地跳到了PROCESSING_SHIPMENT后续所有操作都建立在错误前提上。第三层幻觉它很“轻量”。while True看似没有额外依赖启动快。但它的轻量是虚假的。每个循环迭代都在内存里维护着完整的state对象、action上下文、result缓存。当一个 Agent 需要处理一个平均耗时 8 分钟的复杂诊断流程比如分析 500 行日志 调用 3 个模型 生成 PDF 报告这个state对象会不断膨胀。我们做过压力测试一个基于while True的诊断 Agent在并发 200 路时JVM 堆内存峰值稳定在 4.2GBGC 频率高到影响响应。而它的核心业务逻辑其实只需要不到 200MB。多出来的 4GB全是为“维持循环”这个动作支付的昂贵租金。2.2 Durable Execution把“状态”从内存里解放出来拆掉 while 循环不是要 Agent “停下来”而是让它“可暂停、可恢复、可追踪、可审计”。这就是Durable Execution持久化执行的核心思想将 Agent 的每一次关键决策点Decision Point都视为一个独立的、有明确输入输出、可持久化存储的“执行单元”Execution Unit。这借鉴了工作流引擎如 Temporal、Cadence和函数计算如 AWS Lambda的设计哲学。想象一下快递分拣中心包裹即用户请求进来不是由一个工人拿着它从头跑到尾而是被自动送到一个个工位即 Execution Unit。在“扫码”工位它被扫描、记录 ID、更新状态为“已接收”在“分拣”工位它被识别目的地、贴上标签、状态更新为“分拣中”在“装车”工位它被放入对应车厢状态变为“已发出”。每个工位只做一件事做完就交出包裹并把当前状态ID、时间戳、操作人、结果写入中央数据库。如果某个工位的工人临时请假包裹就在那里等着不会消失也不会乱跑。Agent 的 Durable Execution 就是这个逻辑用户发起一个“帮我分析这份财报”的请求系统立刻生成一个唯一的execution_id: exec-7a3f9b21。第一个 Execution Unit 是LOAD_DOCUMENT它从对象存储下载 PDFOCR 识别文字将原始文本存入数据库状态更新为LOADED。下一个 Unit 是EXTRACT_KEY_METRICS它从数据库读取exec-7a3f9b21的文本调用财务模型提取营收、利润等指标将结构化 JSON 存入数据库状态更新为METRICS_EXTRACTED。如此往复直到GENERATE_REPORT完成状态变为COMPLETED。整个过程没有一个while True在后台狂转。每个 Unit 都是一个短生命周期的、无状态的函数。它们的“状态”全部外置到数据库如 PostgreSQL 的 JSONB 字段或专用的工作流状态库。Agent 的“大脑”不再是那个永不停歇的循环而是中央的Orchestrator编排器—— 一个轻量级服务它只做两件事监听数据库里状态变更的事件如LOADED - METRICS_EXTRACTED然后根据预定义的规则触发下一个 Unit 的执行。2.3 微服务化不是为了时髦而是为了生存有人会问把一个 Agent 拆成十几个微服务是不是过度设计我的回答是在生产环境里这不是设计选择而是生存必需。while True架构本质上是单体Monolith的。所有逻辑——输入解析、状态管理、模型调用、结果渲染——都挤在一个进程里。这带来三个致命问题故障域过大一个execute_action()里的内存泄漏会让整个 Agent 实例崩溃所有正在处理的请求全部失败。伸缩性僵硬如果LOAD_DOCUMENT单元是 I/O 密集型大量文件读写而GENERATE_REPORT是 CPU 密集型大模型推理你无法单独给前者加机器、给后者加 GPU。你只能给整个单体加资源造成巨大浪费。演进成本高想把EXTRACT_KEY_METRICS单元从 Python 换成 Rust为了性能你得重构整个单体还得确保所有state格式兼容。而在微服务架构下你只需部署一个新的metrics-extractor-rs服务修改 Orchestrator 的路由规则旧服务可以并行运行灰度切换。这正是为什么“微服务”会成为热搜词。它不是银弹但它把一个庞大、脆弱的系统分解成一组小的、独立的、可替换的乐高积木。每个积木即每个 Execution Unit 的实现服务可以用最适合的语言编写Python 处理胶水逻辑Rust 做高性能计算Go 写高并发 API用最适合的资源运行CPU 实例跑模型GPU 实例跑推理内存实例跑缓存有自己独立的监控、告警、发布流水线。我们在一个为制造业客户做的预测性维护 Agent 中就采用了这种模式。vibration-analyzer服务用 C 编写直接调用 FFT 库和failure-predictor服务用 PyTorch 训练的 LSTM 模型完全解耦。当客户要求将预测模型升级为 Transformer 架构时我们只替换了failure-predictor的 Docker 镜像整个系统零停机。如果它们还锁在同一个while True循环里那次升级会是一场为期三天的灾难性发布。3. 核心细节解析与实操要点从概念到落地的七道坎3.1 执行单元Execution Unit的边界划分什么该拆什么不该拆这是落地的第一道坎也是最容易犯错的地方。拆得太细服务间通信开销RPC、序列化会吃掉大部分性能拆得太粗又失去了微服务的弹性优势。我的经验是遵循“单一职责 状态变更”双准则。单一职责每个 Unit 必须只做一件明确的事且这件事有清晰的输入和输出。例如✅ 好的 Unitparse_email_body输入原始邮件 HTML输出纯文本内容、call_llm_for_summary输入文本输出摘要 JSON。❌ 坏的 Unitprocess_user_request输入用户消息输出最终响应—— 这还是一个单体只是名字换了。状态变更只有当一个 Unit 的执行会导致 Agent 整体状态发生可观测、可持久化的改变时它才值得成为一个独立的 Unit。状态变更必须满足原子性要么全部成功状态更新数据落库要么全部失败状态回滚。幂等性同一个execution_id下重复调用该 Unit结果必须一致比如send_notification发送邮件第二次调用应检测到已发送直接返回成功。我们曾在一个法律咨询 Agent 中把check_conflict_of_interest检查利益冲突和generate_legal_opinion生成法律意见强行合并为一个 Unit。结果发现check_conflict_of_interest需要调用外部律师数据库慢可能超时而generate_legal_opinion是本地模型快。合并后一次数据库超时就会拖垮整个快速生成流程。后来拆开check_conflict_of_interest设置了 5 秒超时和重试机制generate_legal_opinion则永远在 200ms 内完成系统稳定性提升了 40%。提示画一张“状态迁移图”是避免边界错误的最有效方法。横轴是execution_id的生命周期INITIALIZED → LOADING → PARSING → ... → COMPLETED/FAILED纵轴是所有可能的 Unit 名称。每条箭头代表一个 Unit 的触发条件。如果发现某个 Unit 的箭头指向了多个不同状态或者多个 Unit 共享同一个状态那边界就划错了。3.2 状态State的建模与存储JSONB 还是专用 Schema状态是 Durable Execution 的心脏。它必须能承载任意复杂的数据结构又要能被高效查询比如“查出所有卡在PENDING_PAYMENT状态的订单”。我们对比了三种主流方案方案优点缺点我们的选型纯 JSONB (PostgreSQL)开发极快Schema Free支持 GIN 索引加速 JSON 查询复杂嵌套查询性能下降明显无法强制数据类型约束历史版本追溯困难✅初期首选。用jsonb_path_exists()和jsonb_path_query()做基础过滤足够支撑 90% 场景。混合 Schema (表JSONB)关键字段如status,created_at,user_id走强类型列保证查询性能和约束非结构化数据如model_output,debug_trace放 JSONB设计稍复杂需要定义哪些是“关键字段”✅中期主力。我们将execution_id,status,next_unit,retry_count,last_updated作为表字段其余全放payloadJSONB。专用工作流状态库 (如 Temporal)开箱即用的重试、超时、信号、查询 API企业级可靠性学习成本高运维复杂对简单场景是杀鸡用牛刀⚠️大型项目后期。当我们 Agent 日均执行量突破 500 万次且需要精确的“定时唤醒”如“3 天后自动催缴”时才引入。一个关键细节永远不要在 State 里存“函数引用”或“闭包”。我见过最惨的案例是有人把lambda x: x.upper()这样的匿名函数序列化成字符串存进数据库。当服务升级 Python 版本后反序列化直接报错。State 必须是纯数据POJO所有逻辑都放在 Unit 的代码里。3.3 编排器Orchestrator的实现事件驱动还是轮询Orchestrator 是整个系统的“交通指挥中心”。它决定“下一步该做什么”。实现方式有两种主流路径事件驱动推荐每个 Unit 执行完毕向消息队列如 Kafka、RabbitMQ发布一个UnitCompleted事件包含execution_id和新status。Orchestrator 订阅此 Topic收到事件后查询数据库获取当前完整状态匹配预定义的“状态机规则”然后调用下一个 Unit 的 API。✅ 优势实时性高毫秒级解耦彻底天然支持水平扩展起 10 个 Orchestrator 实例消费同一个 Topic。❌ 劣势引入消息中间件增加了运维复杂度需要处理消息重复、丢失等分布式难题。轮询简易版Orchestrator 定期如每 500ms查询数据库找出所有status IN (LOADED, PARSED)的记录然后批量触发对应的 Unit。✅ 优势零外部依赖50 行代码就能跑起来适合 PoC 或小流量场景。❌ 劣势有延迟最长 500ms轮询本身是无意义的数据库压力高并发下容易成为瓶颈。我们的实践是起步用轮询上线前切事件驱动。轮询让我们在 2 小时内就跑通了第一个端到端流程验证了核心逻辑而上线前我们用 Kafka 替换轮询将平均端到端延迟从 1.2 秒降到了 320 毫秒并轻松支撑了 10 倍的并发增长。注意无论哪种方式Orchestrator 本身必须是无状态的。它的所有决策依据都来自数据库或消息队列而不是自己的内存变量。这样才能保证任意一个实例宕机其他实例能无缝接管。3.4 错误处理与重试不是“try-catch”那么简单while True里错误处理往往就是一层try-except捕获异常后continue。Durable Execution 下错误是系统的一等公民必须被分类、记录、并触发明确的恢复策略。我们定义了三级错误Transient Error瞬时错误网络超时、数据库连接池满、外部 API 返回 429。策略指数退避重试第一次 100ms 后重试第二次 200ms第三次 400ms... 最多 5 次。Business Error业务错误用户输入格式错误、权限不足、数据不存在。策略标记为FAILED并写入error_message字段通知用户。这类错误重试毫无意义。Fatal Error致命错误代码逻辑 Bug、不可恢复的序列化失败。策略标记为CRASHED触发告警人工介入。系统绝不自动重试防止雪崩。关键技巧重试必须带上retry_count。每次 Unit 执行前Orchestrator 会检查数据库里该execution_id的retry_count。如果已达上限如 5 次就不再触发直接标记为FAILED。这个计数器必须是数据库原子操作UPDATE ... SET retry_count retry_count 1 WHERE execution_id ? AND retry_count 5否则在并发下会失效。我们在一个电商比价 Agent 中曾因忘记加retry_count限制导致一个商品链接失效的错误引发了 1200 次重试瞬间打垮了比价服务的限流阀值。那次事故后“重试必带计数器”成了我们 Code Review 的第一条铁律。4. 实操过程与核心环节实现手把手搭建你的第一个 Durable Agent4.1 环境准备与工具选型务实主义者的清单别被“微服务”吓到。一个最小可行的 Durable Agent只需要 3 个组件全部用开源、成熟、易上手的工具数据库PostgreSQL 14。理由JSONB 支持完美事务 ACID 保障强LISTEN/NOTIFY机制可模拟轻量级事件总线省去 Kafka。安装命令Ubuntusudo apt update sudo apt install -y postgresql postgresql-contrib sudo -u postgres psql -c CREATE DATABASE durable_agent;消息/事件总线可选但强烈推荐RabbitMQ。比 Kafka 更轻量管理界面友好对新手极其友好。Docker 一键启动docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 -e RABBITMQ_DEFAULT_USERadmin -e RABBITMQ_DEFAULT_PASSsecret rabbitmq:3-management访问http://localhost:15672账号 admin/secret创建一个名为agent_events的 Exchange。编程语言与框架Python 3.10 FastAPI SQLAlchemy Celery。FastAPI 做 API 网关SQLAlchemy 操作 PostgreSQLCelery 作为分布式任务队列它内置了 RabbitMQ 支持且自带重试、定时、监控等企业级特性远超手写轮询。提示不要纠结于“是否必须用 Kubernetes”。一个docker-compose.yml文件就能跑起整个开发环境。K8s 是为千级 Pod 准备的不是为你的第一个 Agent 准备的。4.2 数据库 Schema 与初始化让状态“活”起来创建executions表这是整个系统的基石-- 创建 executions 表 CREATE TABLE executions ( id SERIAL PRIMARY KEY, execution_id VARCHAR(64) UNIQUE NOT NULL, -- exec-7a3f9b21 status VARCHAR(32) NOT NULL DEFAULT INITIALIZED, next_unit VARCHAR(128), -- 下一个要执行的 Unit 名如 load_document payload JSONB NOT NULL DEFAULT {}, -- 所有状态数据 retry_count INTEGER NOT NULL DEFAULT 0, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), error_message TEXT ); -- 为常用查询字段创建索引 CREATE INDEX idx_exec_status ON executions(status); CREATE INDEX idx_exec_next_unit ON executions(next_unit); CREATE INDEX idx_exec_updated_at ON executions(updated_at);初始化一个测试执行记录INSERT INTO executions (execution_id, status, next_unit, payload) VALUES (exec-test-001, INITIALIZED, receive_input, {user_id: u123, request: Analyze this PDF});4.3 编写第一个 Execution Unitreceive_input这是一个最简单的 Unit负责接收用户输入并更新状态。我们用 FastAPI 写一个 HTTP 接口# unit_receive_input.py from fastapi import FastAPI, HTTPException from sqlalchemy import create_engine, text import json app FastAPI() # 数据库连接生产环境请用连接池 engine create_engine(postgresql://admin:secretlocalhost:5432/durable_agent) app.post(/units/receive_input) async def receive_input(execution_id: str, user_input: str): try: with engine.connect() as conn: # 1. 检查 execution_id 是否存在且状态为 INITIALIZED result conn.execute(text( SELECT id, status FROM executions WHERE execution_id :eid ), {eid: execution_id}) row result.fetchone() if not row or row[1] ! INITIALIZED: raise HTTPException(400, fInvalid execution_id {execution_id} or status not INITIALIZED) # 2. 更新状态为 RECEIVED并存入用户输入 conn.execute(text( UPDATE executions SET status RECEIVED, payload payload || :new_payload::jsonb, next_unit parse_input, updated_at NOW() WHERE execution_id :eid ), { eid: execution_id, new_payload: json.dumps({user_input: user_input, received_at: 2023-10-27T10:00:00Z}) }) conn.commit() return {status: success, execution_id: execution_id, next_unit: parse_input} except Exception as e: raise HTTPException(500, fFailed to receive input: {str(e)})启动它uvicorn unit_receive_input:app --host 0.0.0.0 --port 80014.4 编写 Orchestrator用 Celery 触发下一个 UnitOrchestrator 的核心是监听状态变更并触发下一个 Unit。我们用 Celery 的task来实现 Unit 的调用并用 PostgreSQL 的LISTEN/NOTIFY机制来监听状态变更免去 Kafka简化 PoC# orchestrator.py from celery import Celery import psycopg2 import json import threading # Celery 配置使用 RabbitMQ app Celery(orchestrator) app.config_from_object(celeryconfig) # celeryconfig.py 包含 broker_url 等 # 定义一个 Celery Task用于调用任意 Unit app.task(bindTrue, autoretry_for(Exception,), retry_kwargs{max_retries: 3, countdown: 60}) def call_unit(self, execution_id: str, unit_name: str, **kwargs): 通用 Unit 调用函数 # 这里是伪代码实际应调用对应 Unit 的 HTTP API 或 RPC print(fCalling unit {unit_name} for {execution_id}) # 例如requests.post(fhttp://unit-{unit_name}:8000/units/{unit_name}, json{...}) return {result: done} # 监听 PostgreSQL 的状态变更 def listen_to_db(): conn psycopg2.connect(dbnamedurable_agent useradmin passwordsecret) conn.set_isolation_level(psycopg2.extensions.ISOLATION_LEVEL_AUTOCOMMIT) cursor conn.cursor() cursor.execute(LISTEN execution_status_change;) def handle_notify(): conn.poll() while conn.notifies: notify conn.notifies.pop(0) print(Got NOTIFY:, notify.payload) payload json.loads(notify.payload) # 解析 payload触发对应 task if payload.get(status) RECEIVED: call_unit.delay(payload[execution_id], parse_input) # 启动监听线程 t threading.Thread(targethandle_notify, daemonTrue) t.start() if __name__ __main__: listen_to_db() # 启动 Celery worker: celery -A orchestrator worker --loglevelinfo现在当你调用curl -X POST http://localhost:8001/units/receive_input -d {execution_id:exec-test-001, user_input:Hello World}Orchestrator 就会收到通知并自动触发parse_inputUnit。while True的循环已经彻底消失了。5. 常见问题与排查技巧实录那些文档里不会写的坑5.1 问题状态“幽灵更新”——数据库里看到状态变了但 Orchestrator 没反应现象手动执行UPDATE executions SET statusRECEIVED WHERE execution_idexec-test-001;但在 Orchestrator 的日志里没有任何call_unit被触发。根因与排查检查LISTEN是否生效在psql里执行\set QUIET on然后LISTEN execution_status_change;。再手动NOTIFY execution_status_change, {execution_id:exec-test-001,status:RECEIVED};。如果psql窗口没打印Asynchronous notification execution_status_change received说明LISTEN没成功。检查连接隔离级别LISTEN/NOTIFY要求连接的isolation_level必须是READ COMMITTED默认且不能在事务块里BEGIN; ...; END;。我们的listen_to_db()函数里设置了conn.set_isolation_level(...)但如果在其他地方用了事务就会失效。检查通知 PayloadNOTIFY的 payload 必须是合法 JSON 字符串。NOTIFY execution_status_change, invalid json;不会报错但psql也收不到。解决方案在UPDATE语句后强制发送NOTIFY而不是依赖 ORM 自动触发。在unit_receive_input.py的UPDATE后加上conn.execute(text(NOTIFY execution_status_change, :payload), { payload: json.dumps({execution_id: execution_id, status: RECEIVED}) })5.2 问题Unit 执行“假成功”——HTTP 返回 200但数据库状态没更新现象调用call_unit.delay(...)后Celery 日志显示Task succeeded但executions表里status还是RECEIVEDnext_unit也没变。根因与排查检查 Unit 的事务提交这是最高频的错误。在unit_receive_input.py的UPDATE后必须调用conn.commit()。如果忘了事务在函数结束时自动回滚所有更改丢失。print(Before commit)和print(After commit)是最朴实的调试手段。检查 SQL 语法payload payload || :new_payload::jsonb这种拼接如果:new_payload是空字符串或None会导致整个UPDATE失败。务必在json.dumps()前做空值检查。检查数据库连接Celery Worker 和 Unit 服务是否连接的是同一个数据库实例开发时容易配错localhost指向本机和host.docker.internal指向宿主机。解决方案在所有 Unit 的入口处添加一个“健康检查”app.post(/units/receive_input) async def receive_input(...): try: # ... 数据库操作 ... conn.commit() # 强制查询刚更新的记录验证 check conn.execute(text(SELECT status FROM executions WHERE execution_id :eid), {eid: execution_id}).fetchone() if not check or check[0] ! RECEIVED: raise Exception(fDB update failed. Expected RECEIVED, got {check[0] if check else None}) return {...} except Exception as e: # 记录详细错误包括 conn.url logger.error(fDB op failed: {e}, conn: {conn.engine.url}) raise5.3 问题重试风暴——一个失败触发了上千次重试现象某个 Unit 因代码 Bug 一直抛异常Celery 日志里刷屏式出现Retry in 60s... Retry in 120s...几分钟内生成了 2000 个重试任务压垮了 RabbitMQ。根因与排查检查autoretry_for的异常范围autoretry_for(Exception,)是最危险的配置它会重试所有异常包括ValueError、TypeError这些绝对不该重试的编程错误。应该只重试requests.exceptions.Timeout、psycopg2.OperationalError等明确的瞬时错误。检查retry_kwargs的max_retriesmax_retries3是安全的但max_retriesNone或max_retries100就是灾难。必须显式设置一个合理的上限我们团队规定是 5。检查重试的countdowncountdown60是固定 60 秒但countdown60 * (2 ** self.request.retries)才是指数退避。固定值在高并发下会造成“重试洪峰”。解决方案采用“防御性重试”模式app.task(bindTrue, autoretry_for(requests.exceptions.Timeout, requests.exceptions.ConnectionError), retry_kwargs{max_retries: 5}, reject_on_worker_lostTrue) def call_unit(self, execution_id: str, unit_name: str): try: # ... 执行逻辑 ... except (requests.exceptions.Timeout, requests.exceptions.ConnectionError) as exc: # 显式触发重试带指数退避 countdown 60 * (2 ** self.request.retries) raise self.retry(excexc, countdowncountdown) except Exception as exc: # 其他所有异常直接失败不重试 logger.error(fFatal error in {unit_name}: {exc}) raise5.4 问题状态机“死锁”——两个 Unit 相互等待现象Unit A的逻辑是“调用Unit B然后根据B的结果决定下一步”而Unit B的逻辑是“调用Unit A的某个子功能”。系统启动后execution_id的状态卡在RUNNING_A永远不前进。根因与排查 这是典型的分布式系统“循环依赖”问题。while True里循环体内的函数调用是同步的栈帧清晰死锁很容易被发现。但在 Durable Execution 下Unit A和Unit B是两个独立进程它们的“调用”是异步的 HTTP 请求或消息投递。A发出请求后就返回状态变成WAITING_FOR_BB收到请求后又发一个请求给A的子服务自己也变成WAITING_FOR_A_SUB。两个服务都在等对方的响应形成分布式死锁。解决方案严格禁止 Unit 之间的直接、同步调用。所有 Unit 必须是“无依赖”的。如果A需要B的结果正确的流程是A执行完毕状态更新为A_COMPLETED并把需要B处理的数据写入payload。Orchestrator 监听到A_COMPLETED触发B。B从数据库读取execution_id的完整payload拿到A的输出进行处理。实操心得在设计阶段就用白板画出所有 Unit 的“输入来源”和“输出去向”。如果箭头形成了闭环立刻重构。我们团队有个硬性规定任何 Unit 的代码里不允许出现httpx.post(http://unit-.*)这样的硬编码 URL。所有跨 Unit 通信

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询