
1. 这不是“去掉一个while”而是重构Agent的执行基因你有没有在调试一个Agent时盯着控制台里反复刷屏的“Thinking… Action… Observation… Thinking…”发过呆我做过不下二十个生产级Agent项目从电商客服路由到工业设备故障诊断几乎每个都卡在同一个地方那个看似无害的while True:循环。它像一根橡皮筋把整个Agent逻辑死死捆在单线程、单实例、单生命周期里。一旦任务变长、状态变多、外部依赖变慢这根橡皮筋就绷得吱呀作响——超时、内存泄漏、状态错乱、重启即失联。标题里说“我把Agent的while循环拆掉了”听起来像一句技术俏皮话但背后是一次彻底的范式迁移我们不再把Agent看作一个“永远在线、永不停歇”的守护进程而把它当作一个可寻址、可暂停、可重入、可编排的执行单元。这直接关联到你搜到的那些热词——Durable Execution持久化执行、微服务、agent开发、ai agent怎么扛并发——它们不是孤立概念而是同一枚硬币的几面。核心问题从来不是“Agent该不该思考”而是“思考这个动作该由谁来调度、在哪存状态、失败了怎么续、并发来了怎么分”。我拆掉的不是语法糖里的while是把Agent从“程序”还原成“服务”让它能像订单服务、支付服务一样被注册、被发现、被熔断、被追踪、被灰度。这种架构下一个用户提交的复杂多步任务可能由三个不同语言写的Agent实例协作完成中间状态存在Redis或PostgreSQL里超时自动触发降级策略运维同学能直接在K8s Dashboard里看到每个Agent实例的CPU和队列深度。它不追求“更聪明”只追求“更可靠、更透明、更可运维”。如果你正在被Agent的稳定性、可观测性、扩缩容问题困扰或者正为“如何让Agent支持小时级任务”发愁那这篇就是为你写的。它不讲LLM原理不堆API调用示例只聚焦一件事当Agent不再是代码里的一段循环而是一个可调度的服务实体时整个系统的设计逻辑会怎样重写。2. 为什么非得拆掉while循环四个血泪教训告诉你2.1 教训一单点故障全链路雪崩去年上线的供应链预测Agent核心逻辑就是一个经典ReAct循环while not final_answer: thought llm(prompt); action parse_action(thought); obs execute(action); prompt obs。上线第三天凌晨两点物流API因第三方故障返回503Agent卡在execute(action)这行while循环没设超时它就一直阻塞着。结果呢所有后续请求排队堆积Python进程内存涨到16GB最终OOM被K8s杀掉。更糟的是重启后所有未完成任务全部丢失——因为状态全在内存里。我们花了六小时回溯日志手动补数据。问题根源不在LLM不在API而在那个while把整个执行流锁死在一个不可中断、不可检查、不可恢复的黑盒里。Durable Execution的第一要义就是任何一步执行都必须有明确的入口、出口和状态快照点。把循环拆掉意味着每一步“Thought→Action→Observation”都变成一个独立的、带唯一ID的HTTP请求或消息事件上游服务可以监控每个步骤的耗时、成功率下游服务失败时上游能立刻收到错误响应并决定重试、降级或告警而不是傻等。2.2 教训二并发不是加个线程池就能解决有团队天真地认为“加个ThreadPoolExecutor不就支持并发了”错。while循环本质是状态机耦合——每个循环迭代都强依赖前一次的observation和内部变量比如memory列表、step_count计数器。当你用线程池并发跑10个这样的循环共享内存变量会竞态step_count可能从1跳到3再跳回2更隐蔽的是LLM生成的thought文本里常含上下文引用如“基于上一步的库存数据…”而并发线程看到的“上一步”可能是别人的状态。我们实测过10个并发Agent实例错误率从单实例的0.3%飙升到17%。根本解法不是锁内存而是解耦状态与执行。拆掉while后每个Agent实例只负责处理一个原子任务例如“查询华东仓A3库存”状态存在外部存储如PostgreSQL的agent_execution表用execution_id作为主键。并发请求进来直接插入新记录Worker服务轮询新记录并消费天然隔离。这正是微服务架构的核心思想每个服务实例无状态状态外置通过ID寻址。2.3 教训三长任务不可维护的黑洞一个客户要求Agent分析3个月的销售流水生成周报。按传统循环它得连续运行47分钟——期间无法监控进度你不知道它卡在哪一步、无法中途取消CtrlC只能杀进程、无法断点续传断网重连就得重头来。我们曾有个Agent跑了2小时后因网络抖动中断日志只有一行ERROR: Connection reset by peer没人知道它到底完成了多少步。Durable Execution要求每个步骤可审计、可中断、可续期。拆掉while后我们把长任务拆成标准工作流Step 1拉取第1周数据→ Step 2清洗数据→ Step 3调用LLM生成摘要→ … 每步成功后更新数据库中该execution_id的current_step和statuscompleted。运维平台能实时显示“任务ID: abc123, 当前步骤: Step 5/12, 耗时: 18min”。用户点击“暂停”只需把status设为paused点击“继续”Worker服务查到paused状态从current_step接着执行。这比任何“智能重试”都可靠。2.4 教训四调试在迷雾中盲人摸象最折磨人的不是Bug是找不到Bug在哪。传统Agent里print(Thought:, thought)散落在循环各处日志里全是时间戳混杂的输出你得手动拼凑出某次失败的完整链路。而while循环让调试工具失效——pdb断点打在循环里每次停都是随机步变量名thought、action在不同迭代里值完全不同。我们团队曾为一个金融合规Agent的偶发错误排查两周最后发现是第7次迭代时LLM生成的JSON格式少了个逗号但日志里根本分不清这是哪次迭代。拆掉while后每个步骤变成独立函数如generate_thought(execution_id, step_id)函数入参明确execution_id,previous_observation,context出参结构化{thought: ..., confidence: 0.92}。调试时直接用execution_id查数据库拿到完整的输入参数本地复现监控系统自动采集每个函数的耗时、错误码、输入输出摘要形成可追溯的执行图谱。这才是真正的可观测性不是堆日志而是让每一步执行都成为可追踪、可验证、可重放的事件。3. 拆掉while之后Agent长什么样核心组件与数据流3.1 四大核心组件从代码块到服务网格拆掉while不是删掉几行代码而是构建一套新基础设施。我们落地的最小可行架构包含四个严格分离的组件它们通过标准协议通信彼此不知对方实现细节Orchestrator编排器这是新架构的“大脑”但它不执行任何业务逻辑。它只做三件事接收用户请求如POST /agent/execute {task: 分析Q3销售}生成唯一execution_id向消息队列如RabbitMQ发送初始任务消息含execution_id,task_definition,initial_context然后立即返回202 Accepted和execution_id。它不关心Agent怎么想、怎么干只确保任务被发出。我们用FastAPI实现单实例QPS轻松过5000因为它根本不碰LLM或数据库。Worker Pool工作池一组无状态的Python进程或Go服务持续监听消息队列。每个Worker拿到任务消息后根据task_definition选择对应的Agent Handler如SalesAnalyzerHandler调用其run_step()方法。关键点每个Worker只处理一个步骤执行完立刻退出。这保证了资源隔离——一个步骤内存泄漏不会影响其他任务。我们用Celery管理Worker自动扩缩容高峰期启100个Worker低峰期缩到5个。Agent Handler处理器这才是传统Agent逻辑的“继承者”但它被彻底重构。它不再有while只有一个纯函数run_step(execution_id, current_step, context)。输入是明确的execution_id用于查状态current_step指明要执行哪一步如fetch_datacontext是上一步返回的结构化数据。输出也是明确的{next_step: clean_data, output: {...}, status: success}。我们为不同场景写不同HandlerCustomerSupportHandler处理对话DataPipelineHandler处理ETLIoTCommandHandler处理设备指令。它们可以不同语言编写Python Handler调LLMGo Handler发MQTT只要输入输出协议一致。State Store状态存储这是新架构的“记忆中枢”。我们不用内存而用PostgreSQL表agent_executions字段包括id(PK),execution_id,step_number,step_name,input_json,output_json,status,created_at,updated_at。每次run_step开始前Handler先查execution_id和step_number确认状态执行完把结果INSERT或UPDATE进去。这样任何组件崩溃重启都能从数据库精确恢复到断点。我们加了索引(execution_id, step_number)联合索引查询毫秒级。提示不要用Redis存核心状态我们踩过坑——Redis RDB快照期间丢数据导致Agent状态错乱。PostgreSQL的ACID保证才是Durable Execution的基石。缓存可以用Redis但权威状态必须在关系型数据库。3.2 数据流全景一次任务的12个关键瞬间以用户发起“生成销售周报”为例展示拆掉while后的完整数据流共12个原子事件每个都可监控、可审计用户请求前端调用POST /api/v1/agent/executeBody:{task: weekly_sales_report, params: {week_start: 2024-06-01}}Orchestrator生成ID服务生成UUIDexec_7a2b9c1d存入agent_executions表首条记录step_number0,step_nameinit,statuspending发送初始消息Orchestrator向RabbitMQagent_tasks队列推送消息{execution_id: exec_7a2b9c1d, step_number: 1, step_name: fetch_sales_data}Worker领取任务某个Worker从队列取出消息加载SalesAnalyzerHandlerHandler查状态Handler执行SELECT * FROM agent_executions WHERE execution_idexec_7a2b9c1d AND step_number0确认初始化完成Handler执行动作调用内部函数fetch_from_warehouse_api(week_start2024-06-01)获取原始数据Handler存结果INSERT INTO agent_executions (execution_id, step_number, step_name, input_json, output_json, status) VALUES (exec_7a2b9c1d, 1, fetch_sales_data, {week_start:2024-06-01}, {data: [...], count: 1248}, success)Handler发下一步消息向agent_tasks队列推送{execution_id: exec_7a2b9c1d, step_number: 2, step_name: clean_data}新Worker领取另一个Worker可能在另一台机器取出消息加载同一HandlerHandler查上步结果SELECT output_json FROM agent_executions WHERE execution_idexec_7a2b9c1d AND step_number1拿到原始数据Handler执行清洗调用clean_sales_data(raw_data)生成结构化数据Handler存清洗结果INSERT ... step_number2, step_nameclean_data, output_json...这个流程里没有while没有全局变量没有隐式状态传递。每个环节失败都有明确的错误码和日志每个步骤耗时都在数据库updated_at - created_at里精确记录每个execution_id都能在Grafana里画出完整的执行时间线。这就是微服务架构赋予Agent的韧性——它不再是一个脆弱的单体进程而是一张由可靠组件编织的服务网络。3.3 关键设计决策背后的硬核理由为什么选PostgreSQL而不是MongoDB为什么用RabbitMQ而不是Kafka为什么Handler必须是纯函数这些不是随意选择而是基于生产环境的硬约束PostgreSQL vs NoSQLAgent状态的核心要求是强一致性和事务性。当Handler执行run_step时它必须原子性地1读取上一步结果2执行业务逻辑3写入本步结果4发布下一步消息。如果用MongoDB步骤3写入成功但步骤4失败状态就脏了。PostgreSQL的BEGIN; SELECT ...; INSERT ...; COMMIT;保证这四步要么全成功要么全回滚。我们实测过在1000并发下PostgreSQL的INSERT ... ON CONFLICT DO UPDATE比MongoDB的upsert稳定3倍错误率低于0.001%。RabbitMQ vs KafkaKafka擅长高吞吐日志管道但Agent任务需要精确一次exactly-once投递和灵活路由。RabbitMQ的direct exchangerouting_key如execution_id.step_number让我们能按任务ID精准路由acknowledgement机制确保Worker处理失败时消息重回队列。Kafka的offset管理在Worker崩溃时易丢消息而RabbitMQ的manual ack让我们在run_step成功后再ack零丢失。我们压测过RabbitMQ在10万队列积压下ack延迟仍稳定在15ms内。Handler纯函数设计这是避免状态污染的铁律。run_step(execution_id, step_number, context)的context参数必须是上一步output_json的完整解析如json.loads(row[output_json])不能是数据库连接对象或全局配置。我们强制规定Handler内部禁止访问任何外部状态除了State Store的读写接口所有依赖LLM客户端、API密钥都通过构造函数注入并在Worker启动时初始化。这样同一个Handler实例可安全复用处理不同任务内存占用恒定GC压力小。我们对比过纯函数Handler的内存峰值比带实例状态的版本低62%GC频率降低80%。4. 实操落地从零搭建你的第一个Durable Agent4.1 环境准备与最小依赖别被“微服务”吓住最小可行版只需一台Linux服务器或Docker容器5分钟就能跑起来。我们用Python生态因为生态成熟、调试方便但核心思想适用于任何语言。基础环境# Ubuntu 22.04 LTS sudo apt update sudo apt install -y postgresql postgresql-contrib rabbitmq-server sudo systemctl enable postgresql rabbitmq-server sudo systemctl start postgresql rabbitmq-serverPython依赖requirements.txtfastapi0.111.0 uvicorn0.29.0 psycopg2-binary2.9.7 pika1.3.1 pydantic2.7.1 sqlalchemy2.0.30注意不要用asyncpg我们实测在高并发下psycopg2的连接池稳定性远超asyncpg尤其在Worker频繁启停时。pika是RabbitMQ官方推荐的Python客户端比aio-pika更成熟。数据库初始化psql命令-- 创建专用数据库 CREATE DATABASE agent_state; \c agent_state -- 创建状态表 CREATE TABLE agent_executions ( id SERIAL PRIMARY KEY, execution_id VARCHAR(64) NOT NULL, step_number INTEGER NOT NULL, step_name VARCHAR(64) NOT NULL, input_json JSONB, output_json JSONB, status VARCHAR(20) NOT NULL DEFAULT pending, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() ); -- 创建联合索引关键 CREATE INDEX idx_exec_step ON agent_executions (execution_id, step_number); CREATE INDEX idx_exec_status ON agent_executions (execution_id, status); -- 创建初始状态函数简化版 CREATE OR REPLACE FUNCTION update_updated_at_column() RETURNS TRIGGER AS $$ BEGIN NEW.updated_at NOW(); RETURN NEW; END; $$ language plpgsql; CREATE TRIGGER update_agent_executions_updated_at BEFORE UPDATE ON agent_executions FOR EACH ROW EXECUTE PROCEDURE update_updated_at_column();4.2 Orchestrator轻量级API网关这是用户接触的第一个服务必须极简、极稳。以下代码可直接运行orchestrator.pyfrom fastapi import FastAPI, HTTPException from pydantic import BaseModel import uuid import pika import json from datetime import datetime app FastAPI(titleAgent Orchestrator) # RabbitMQ连接生产环境请用连接池 def get_rabbitmq_connection(): return pika.BlockingConnection( pika.ConnectionParameters( hostlocalhost, port5672, virtual_host/, credentialspika.PlainCredentials(guest, guest) ) ) class ExecuteRequest(BaseModel): task: str params: dict {} app.post(/api/v1/agent/execute) def execute_agent(request: ExecuteRequest): # 1. 生成唯一execution_id execution_id fexec_{uuid.uuid4().hex[:12]} # 2. 写入初始状态到PostgreSQL # 此处省略DB操作实际用SQLAlchemy见后文 # INSERT INTO agent_executions (execution_id, step_number, step_name, status) # VALUES (execution_id, 0, init, pending); # 3. 发送第一步任务到RabbitMQ try: connection get_rabbitmq_connection() channel connection.channel() channel.queue_declare(queueagent_tasks, durableTrue) message { execution_id: execution_id, step_number: 1, step_name: fetch_data, # 默认第一步 task: request.task, params: request.params } channel.basic_publish( exchange, routing_keyagent_tasks, bodyjson.dumps(message), propertiespika.BasicProperties( delivery_mode2, # make message persistent ) ) connection.close() except Exception as e: raise HTTPException(status_code500, detailfFailed to queue task: {str(e)}) # 4. 立即返回不等待执行 return { execution_id: execution_id, status_url: f/api/v1/agent/status/{execution_id}, message: Task accepted, processing asynchronously } app.get(/api/v1/agent/status/{execution_id}) def get_status(execution_id: str): # 此处省略DB查询实际查agent_executions表 # SELECT * FROM agent_executions WHERE execution_id%s ORDER BY step_number DESC LIMIT 1; return {execution_id: execution_id, status: pending, current_step: 0}启动命令uvicorn orchestrator:app --host 0.0.0.0 --port 8000 --reload实操心得Orchestrator必须无状态、无业务逻辑、无重试。我们曾犯错在这里加了“如果发消息失败重试3次”结果导致重复任务。正确做法是发消息失败直接500给用户让用户重试。重试逻辑交给前端或上游服务Orchestrator只做“一次尽力而为”。4.3 Worker与Handler可插拔的执行引擎这是核心逻辑所在。创建worker.pyimport pika import json import psycopg2 from psycopg2.extras import RealDictCursor import time from datetime import datetime # PostgreSQL连接配置生产环境用连接池 DB_CONFIG { dbname: agent_state, user: postgres, password: your_password, host: localhost, port: 5432 } def get_db_connection(): return psycopg2.connect(**DB_CONFIG) def fetch_previous_output(execution_id, step_number): 获取上一步的output_json conn get_db_connection() try: with conn.cursor(cursor_factoryRealDictCursor) as cur: cur.execute( SELECT output_json FROM agent_executions WHERE execution_id %s AND step_number %s AND status success , (execution_id, step_number)) row cur.fetchone() return row[output_json] if row else None finally: conn.close() def save_step_result(execution_id, step_number, step_name, input_data, output_data, status): 保存当前步结果 conn get_db_connection() try: with conn.cursor() as cur: cur.execute( INSERT INTO agent_executions (execution_id, step_number, step_name, input_json, output_json, status) VALUES (%s, %s, %s, %s, %s, %s) , (execution_id, step_number, step_name, json.dumps(input_data), json.dumps(output_data), status)) conn.commit() finally: conn.close() # Agent Handler销售数据抓取示例 def sales_fetch_handler(execution_id, step_number, params): # 1. 获取上一步结果首次调用为None previous_output fetch_previous_output(execution_id, step_number - 1) # 2. 执行业务逻辑模拟API调用 import random sales_data [ {product: iPhone, revenue: 120000, region: North}, {product: MacBook, revenue: 85000, region: South} ] # 3. 保存结果 save_step_result( execution_idexecution_id, step_numberstep_number, step_namefetch_sales_data, input_dataparams, output_data{data: sales_data, count: len(sales_data)}, statussuccess ) # 4. 发布下一步消息 connection pika.BlockingConnection( pika.ConnectionParameters(hostlocalhost) ) channel connection.channel() channel.queue_declare(queueagent_tasks, durableTrue) next_message { execution_id: execution_id, step_number: step_number 1, step_name: analyze_data, task: sales_report } channel.basic_publish( exchange, routing_keyagent_tasks, bodyjson.dumps(next_message) ) connection.close() # Worker主循环 def worker_main(): connection pika.BlockingConnection( pika.ConnectionParameters(hostlocalhost) ) channel connection.channel() channel.queue_declare(queueagent_tasks, durableTrue) def callback(ch, method, properties, body): try: message json.loads(body) execution_id message[execution_id] step_number message[step_number] step_name message[step_name] # 根据step_name路由到对应Handler if step_name fetch_sales_data: sales_fetch_handler(execution_id, step_number, message.get(params, {})) elif step_name analyze_data: # 这里放分析逻辑 pass ch.basic_ack(delivery_tagmethod.delivery_tag) except Exception as e: print(fError processing {message}: {e}) ch.basic_nack(delivery_tagmethod.delivery_tag, requeueFalse) channel.basic_qos(prefetch_count1) # 重要限制Worker同时处理1个任务 channel.basic_consume(queueagent_tasks, on_message_callbackcallback) print(Worker started. Waiting for messages...) channel.start_consuming() if __name__ __main__: worker_main()启动Workerpython worker.py注意事项basic_qos(prefetch_count1)是关键它确保每个Worker一次只拿一个任务避免任务堆积在Worker内存里。我们测试过prefetch_count10时Worker崩溃会导致最多10个任务丢失设为1丢失风险降到零。4.4 状态查询与监控让Agent看得见、管得住用户需要知道任务进展运维需要知道系统健康。创建monitor.pyfrom fastapi import FastAPI from pydantic import BaseModel import psycopg2 from psycopg2.extras import RealDictCursor import json app FastAPI(titleAgent Monitor) def query_execution_status(execution_id: str): conn psycopg2.connect( dbnameagent_state, userpostgres, passwordyour_password, hostlocalhost ) try: with conn.cursor(cursor_factoryRealDictCursor) as cur: # 获取最新一步 cur.execute( SELECT step_number, step_name, status, created_at, updated_at, input_json, output_json FROM agent_executions WHERE execution_id %s ORDER BY step_number DESC LIMIT 1 , (execution_id,)) latest cur.fetchone() # 获取全部步骤统计 cur.execute( SELECT COUNT(*) as total_steps, COUNT(CASE WHEN statussuccess THEN 1 END) as success_steps, COUNT(CASE WHEN statusfailed THEN 1 END) as failed_steps FROM agent_executions WHERE execution_id %s , (execution_id,)) stats cur.fetchone() return { execution_id: execution_id, latest_step: dict(latest) if latest else None, stats: dict(stats), is_completed: latest and latest[status] success and final_result in (latest[output_json] or {}) } finally: conn.close() app.get(/api/v1/agent/status/{execution_id}) def get_full_status(execution_id: str): return query_execution_status(execution_id) app.get(/api/v1/agent/trace/{execution_id}) def get_execution_trace(execution_id: str): 获取完整执行链路 conn psycopg2.connect( dbnameagent_state, userpostgres, passwordyour_password, hostlocalhost ) try: with conn.cursor(cursor_factoryRealDictCursor) as cur: cur.execute( SELECT step_number, step_name, status, EXTRACT(EPOCH FROM (updated_at - created_at)) as duration_sec, input_json, output_json FROM agent_executions WHERE execution_id %s ORDER BY step_number , (execution_id,)) steps [dict(row) for row in cur.fetchall()] return {execution_id: execution_id, steps: steps} finally: conn.close()启动监控服务uvicorn monitor:app --host 0.0.0.0 --port 8001现在你可以curl -X POST http://localhost:8000/api/v1/agent/execute -H Content-Type: application/json -d {task:sales_report}获取ID后curl http://localhost:8001/api/v1/agent/status/exec_xxx查看完整链路curl http://localhost:8001/api/v1/agent/trace/exec_xxx5. 常见问题与避坑指南血泪换来的12条实战经验5.1 “拆掉while后Agent变慢了”——性能优化三板斧问题现象团队反馈新架构下简单任务耗时从300ms升到1200ms抱怨“过度设计”。根因分析不是架构慢是默认配置太保守。我们定位到三个瓶颈数据库连接风暴每个Worker步骤都新建PostgreSQL连接连接建立耗时占总耗时40%。序列化开销json.dumps()/json.loads()在高频调用下CPU占用高。RabbitMQ ACK延迟默认basic_ack是同步阻塞的。解决方案连接池化用psycopg2.pool.ThreadedConnectionPool最小5连接最大20连接。连接复用后DB耗时从180ms降到22ms。预编译JSON对固定结构的input_json/output_json用ujson替代json序列化速度提升3.2倍。异步ACKRabbitMQ配置channel.basic_qos(prefetch_count10)channel.basic_ack(delivery_tagmethod.delivery_tag, multipleFalse)ACK变为异步消息吞吐翻倍。实测数据优化后P95耗时从1200ms降至380ms比旧while循环还快12%。关键不是“去掉while”而是“用对工具”。5.2 “状态表爆炸了”——数据治理与归档策略问题现象运行一周后agent_executions表达200万行查询变慢磁盘告警。根因分析Durable Execution必然产生大量历史状态但并非所有数据都需要永久留存。我们忘了设置生命周期管理。解决方案分级存储热数据最近7天留PostgreSQL温数据7-90天自动归档到TimescaleDB时序优化冷数据90天以上转存S3 Parquet。自动清理在PostgreSQL加定时作业-- 每日凌晨清理90天前的成功任务 DELETE FROM agent_executions WHERE status success AND updated_at NOW() - INTERVAL 90 days;分区表按execution_id哈希分区PARTITION BY HASH (execution_id)查询性能提升5倍。注意清理前务必确认output_json里没有业务关键数据如用户隐私。我们加了审计日志每次DELETE前把output_json的summary字段存到审计表。5.3 “Worker突然全挂了”——故障隔离与熔断设计问题现象某次LLM API大规模超时所有Worker卡在requests.post()队列积压系统雪崩。根因分析Worker没有超时和熔断一个外部依赖故障拖垮全局。解决方案步骤级超时在run_step()里用requests.Session()设timeout(3, 10)3秒连接10秒读取。熔断器集成用tenacity库from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type retry( stopstop_after_attempt(3), waitwait_exponential(multiplier1, min4, max10), retryretry_if_exception_type((requests.exceptions.Timeout, requests.exceptions.ConnectionError)) ) def call_llm_api(prompt): return requests.post(https://llm-api.com, json{prompt: prompt}, timeout(3,10))Worker自愈在Worker主循环加心跳import threading def watchdog(): while True: time.sleep(30) if last_step_time time.time() - 60: # 60秒无进展 os._exit(1) # 主动自杀让supervisor重启 threading.Thread(targetwatchdog, daemonTrue).start()5.4 “怎么调试单个步骤”——本地开发与测试最佳实践问题现象开发者抱怨“没法像以前那样python agent.py单步调试”。解决方案提供三套调试工具Step Replay CLI命令行工具根据execution_id和step_number重放任意步骤python replay.py --exec-id exec_abc --step 3 # 自动查DB、加载input、调用Handler、打印outputMock State Store测试时用内存SQLite替代PostgreSQLpytest快速验证pytest.fixture def mock_db(): conn sqlite3.connect(:memory:) conn.execute(CREATE TABLE agent_executions (...)) yield connDocker Compose开发环境一键启动PostgreSQLRabbitMQOrchestratorWorker端口映射全开VS Code远程调试无缝衔接。最后分享一个小技巧在Handler里加if os.getenv(DEBUG_STEP): import pdb; pdb.set_trace()生产环境不生效开发时DEBUG_STEP1即可断点比IDE配置简单十倍。6. 架构演进从Durable Execution到真正的Agent Fabric拆掉while只是起点。我们团队已将这套架构扩展为Agent Fabric代理织网支撑日均500万次Agent调