十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

智能任务协同Agent:生产级多Agent系统架构与实操

智能任务协同Agent:生产级多Agent系统架构与实操 1. 什么是“智能任务协同Agent”它不是又一个AI玩具“智能任务协同Agent”这六个字最近在技术社区、招聘JD和开源项目列表里高频出现但很多人点开后发现——要么是PPT里的概念图要么是调用几行LangChain代码就宣称“已实现Agent”要么干脆是套壳的ChatUI。我从2022年就开始带团队落地真实业务中的多Agent系统做过电商履约链路的自动异常协同、制造业产线工单的跨系统调度、还有金融风控中三方数据源的动态拉取与交叉验证。这些都不是Demo而是每天处理上万次任务请求、SLA要求99.95%可用性的生产系统。所以今天不讲“Agent是什么”我们直接拆解一个真正能扛住业务压力的智能任务协同Agent到底长什么样、为什么必须这么设计、你在自己搭的时候最容易在哪一步崩盘。核心关键词“智能任务协同”四个字每个词都带着硬约束“智能”不是指它会聊天而是指它能在信息不全、规则模糊、依赖方响应延迟甚至失败的情况下自主判断下一步该问谁、查什么、等多久、降级到哪条备选路径“任务”不是API调用而是有明确输入契约、输出承诺、超时阈值、重试策略、回滚机制的原子工作单元“协同”不是多个Agent排排坐而是存在主控AgentOrchestrator、执行AgentWorker、状态观察AgentWatcher、异常仲裁AgentAdjudicator的职责分层且它们之间通过结构化信道交换带版本号、带签名、带上下文快照的消息最后那个“Agent”在我们团队的定义里必须满足三个铁律可中断、可审计、可回放——任意时刻暂停能导出完整执行栈当前上下文快照任意一次失败能定位到是哪个Agent在哪个上下文片段下触发了哪条决策分支任意一次成功运行能用相同输入相同环境变量100%复现全过程。Python是实现载体不是目的。我们用Python是因为它的生态对快速构建状态机、序列化上下文、集成异步IO、对接各类数据库和消息队列足够成熟而不是因为它“适合写AI”。如果你正打算用这个标题去写简历、启动内部项目、或者接外包需求先问自己三个问题你的任务有没有明确的失败定义协同方是否来自不同系统比如CRM、ERP、短信网关、人工审核后台上下文变化是否会导致前序决策失效例如库存数在你调度途中被抢光如果三个答案都是“是”那接下来的内容就是你省掉三个月踩坑时间的实操地图。2. 整体架构设计为什么必须放弃“单Agent全能幻想”2.1 单Agent模式的致命缺陷从“能跑通”到“不敢上线”的鸿沟很多初学者的第一个Agent项目是用LangChain或LlamaIndex写一个“客服问答Agent”它接收用户问题调用几个工具查知识库、查订单、发邮件然后拼答案返回。这确实能跑通但把它放到真实业务流里立刻暴露三个无法绕过的硬伤第一状态不可见。当用户问“我的订单为什么还没发货”Agent查完订单状态又去调用物流接口再比对仓库系统整个过程像黑箱。运维人员看不到它卡在哪个环节、等哪个依赖、重试了几次。我们曾遇到一个案例Agent因物流API限流连续重试12次每次间隔3秒导致整个订单查询链路阻塞4分钟——而监控系统只看到“API响应慢”根本不知道是Agent在死循环。第二错误不可隔离。单Agent里所有工具调用共享同一上下文和内存空间。一旦某个工具比如调用旧版SOAP接口抛出未捕获异常整个Agent实例崩溃之前已获取的订单信息、用户画像、优惠券状态全部丢失。更糟的是它可能把部分中间结果写入下游系统比如先发了确认短信再因库存不足回滚订单造成数据不一致。第三协同无契约。所谓“协同”在单Agent里只是函数调用顺序。但真实业务中“协同”意味着A系统说“我需要B系统提供X字段格式为JSON Schema v2.1超时3s最多重试2次失败后走人工兜底流程”。单Agent无法表达这种契约更无法在B系统升级Schema后自动检测兼容性。提示别被“AutoGen支持多Agent”这类宣传误导。AutoGen的GroupChat本质仍是单进程内多个角色函数轮转没有独立生命周期、没有跨进程消息总线、没有持久化状态存储——它适合教学演示不适合日均百万级任务调度。2.2 我们采用的四层协同架构每个Agent都是可插拔的“微服务”我们落地的智能任务协同Agent系统采用清晰分层的四角色架构所有Agent实例均独立部署、独立扩缩容、独立监控告警角色职责关键能力典型技术栈Orchestrator主控Agent接收原始任务请求解析意图拆解为子任务DAG分发给Worker聚合结果处理超时/失败/降级动态DAG编排、上下文快照存档、跨Worker事务协调Python Celery Redis Streams SQLAlchemyWorker执行Agent执行具体原子任务如“调用支付网关”、“生成PDF合同”只关心输入输出契约不感知全局流程工具调用沙箱、输入校验、输出标准化、失败原因分类上报Python Pydantic Requests TenacityWatcher观察Agent持续监听任务状态变更、外部事件如库存更新、人工审核完成、超时信号向Orchestrator推送事件事件驱动、低延迟订阅、状态变更diff检测Python Kafka Consumer APSchedulerAdjudicator仲裁Agent当Orchestrator判定任务失败时分析失败根因是网络抖动依赖方宕机还是业务规则变更决定走哪条兜底路径失败模式识别、知识库检索、人工介入触发Python FAISS Llama-3-8B-Instruct本地微调这个架构的核心思想是把“智能”从执行层剥离交给Orchestrator和Adjudicator做决策让Worker专注可靠执行。Worker不需要大模型一个用Pydantic严格校验输入、用Tenacity配置指数退避重试、用OpenTelemetry打点的纯Python函数就够了真正的“智能”体现在Orchestrator如何根据Watcher推送的“库存突降至0”事件动态修改DAG跳过发货子任务直连客服系统生成补偿方案。2.3 上下文管理不是“传参”而是“状态版本控制”很多人以为上下文管理就是把用户问题、历史对话、工具返回结果塞进一个dict传给下一个Agent。这是灾难的开始。我们强制要求所有上下文必须是带版本号、带签名、带TTL的不可变对象。具体实现每次上下文变更比如Worker返回新数据Orchestrator不直接修改原对象而是生成新版本from dataclasses import dataclass from typing import Dict, Any, Optional import hashlib import time dataclass(frozenTrue) class ContextSnapshot: version: str # SHA256(content timestamp parent_version) created_at: float ttl: float # 秒级TTL过期自动归档 content: Dict[str, Any] parent_version: Optional[str] None def __post_init__(self): # 冻结对象禁止运行时修改 object.__setattr__(self, version, self._calc_version()) def _calc_version(self) - str: content_str str(sorted(self.content.items())) raw f{content_str}{self.created_at}{self.parent_version or } return hashlib.sha256(raw.encode()).hexdigest()[:16]为什么必须这样因为当Adjudicator要分析失败原因时它需要精确比对“失败时刻的上下文”和“3分钟前成功时刻的上下文”找出差异字段比如inventory_count从100变成0。如果上下文是可变对象这种diff毫无意义。我们线上系统每秒产生约200个上下文快照全部存入TimescaleDB按task_id version索引支持毫秒级回溯。3. 核心细节解析任务调度、上下文管理、Agent通信的实操陷阱3.1 任务调度别迷信“智能调度”先搞定确定性优先级队列“智能任务调度”常被误解为用大模型预测哪个任务该优先处理。错。在我们系统里95%的调度决策基于确定性规则只有5%留给模型。原因很简单业务方需要可解释、可审计、可压测的调度逻辑。我们采用三级优先级队列P0队列实时强保SLA1s的任务如“支付结果通知”。使用Redis Sorted Setscoreunix_timestamp绝对时间戳Worker用ZPOPMIN争抢超时未ACK则自动释放。P1队列业务优先级按业务维度加权如VIP用户订单权重×3促销活动订单权重×2。权重由Orchestrator在入队时计算并写入消息头Worker消费时按权重比例分配算力。P2队列资源敏感型如“生成月度报表”需占用大量CPU。我们用Kubernetes HPA配合自定义指标如cpu_usage_percent当集群CPU70%时自动将P2任务延迟至夜间执行。关键陷阱不要用时间戳做唯一排序依据。我们曾在线上遇到NTP时钟漂移导致两个Worker拿到相同score的任务引发重复执行。解决方案是score int(time.time()) * 1000000 task_id_hash % 1000000确保即使时间戳相同task_id哈希也能打破平局。3.2 上下文管理三类必存字段与自动清理策略上下文不是越大越好。我们强制规定每个ContextSnapshot必须包含且仅包含三类字段契约字段Contract Fields由Orchestrator在任务创建时注入永不变更。包括task_idUUIDv4、requester_id发起方ID、deadline_unix绝对截止时间、retry_policyJSON Schema定义的重试规则。执行字段Execution FieldsWorker执行过程中写入只增不改。包括worker_id执行者ID、start_time、end_time、output_schema_version输出数据格式版本、tool_call_duration_ms。观测字段Observation FieldsWatcher监听到的外部事件带时间戳。包括inventory_update_at、manual_review_complete_at、payment_gateway_latency_ms。自动清理策略所有快照按created_at分区存入TimescaleDB设置数据保留策略最近7天全量保留支持任意维度查询7-30天聚合为每日统计如“平均上下文大小”、“P0任务占比”30天以上自动归档至冷存储S3 Glacier仅保留task_id和final_status用于审计注意我们禁用任何ORM的懒加载lazy loading特性。ContextSnapshot的content字段在序列化时必须是纯dict禁止嵌套SQLAlchemy对象。否则Worker反序列化时会意外触发数据库查询拖垮整个系统。3.3 Agent间通信为什么坚持用Kafka而非HTTP或gRPC选型对比表基于我们压测10万QPS场景协议端到端延迟P99故障隔离性消息追溯能力运维复杂度适用场景HTTP REST120ms差调用方直连一方宕机全链路雪崩弱需额外埋点低内部调试、低频管理接口gRPC45ms中需客户端负载均衡中需集成OpenTelemetry高需维护proto文件、证书高性能内部服务调用Kafka28ms强生产者/消费者完全解耦强offset可查、消息可重放中需K8s部署Kafka集群生产环境Agent通信主干道我们所有Agent间的业务消息非心跳、非配置必须走Kafka。Topic命名规范{env}.{domain}.{role}_to_{role}如prod.order.orchestrator_to_worker。每条消息强制包含correlation_id: UUIDv4贯穿整个任务生命周期trace_id: OpenTelemetry trace ID用于链路追踪context_version: 对应ContextSnapshot.version确保Worker执行时上下文准确payload: 序列化后的任务指令Pydantic Model实操心得Kafka消费者组Consumer Group的group.id必须按Agent角色环境版本号组合如worker-prod-v2.3.1。这样升级Worker时新版本可以先以新group.id消费验证无误后再将老版本group.id下线实现零停机灰度。4. 实操过程从零搭建可运行的最小协同系统含完整代码4.1 环境准备精简到极致的Python依赖我们拒绝“pip install langchain”式的大包依赖。最小可行系统仅需6个核心包# requirements.txt celery5.3.6 # 任务分发与Worker管理 redis4.6.0 # Broker和Result Backend kafka-python2.0.2 # Kafka Producer/Consumer pydantic2.6.4 # 输入输出校验与序列化 tenacity8.2.3 # 可配置重试策略 sqlalchemy2.0.29 # 上下文快照持久化安装命令避免版本冲突# 创建干净虚拟环境 python -m venv agent_env source agent_env/bin/activate # Linux/Mac # agent_env\Scripts\activate # Windows # 强制指定源国内加速 pip install -i https://pypi.tuna.tsinghua.edu.cn/simple/ -r requirements.txt # 验证关键组件 python -c import celery, kafka, pydantic; print(OK)注意不要用conda安装这些包。Celery与conda的环境隔离机制有兼容性问题曾导致我们线上Worker静默退出。坚持用venvpip。4.2 Orchestrator核心代码DAG拆解与上下文快照Orchestrator是系统大脑其核心逻辑是将高层任务拆解为可执行的DAG节点并为每个节点生成带版本的上下文# orchestrator.py from celery import Celery from pydantic import BaseModel, Field from typing import List, Dict, Any, Optional import json import hashlib import time app Celery(orchestrator) app.conf.broker_url redis://localhost:6379/0 app.conf.result_backend redis://localhost:6379/0 class TaskRequest(BaseModel): task_type: str Field(..., description任务类型如ORDER_SHIPMENT) payload: Dict[str, Any] Field(..., description原始请求数据) class TaskNode(BaseModel): node_id: str Field(..., description节点唯一ID) worker_type: str Field(..., description目标Worker类型) input_context: Dict[str, Any] Field(..., description输入上下文) timeout_sec: int Field(30, description超时时间) class ContextSnapshot(BaseModel): version: str created_at: float ttl: float content: Dict[str, Any] parent_version: Optional[str] None def create_context_snapshot( content: Dict[str, Any], parent_version: Optional[str] None, ttl: float 3600.0 ) - ContextSnapshot: 生成带版本号的上下文快照 created_at time.time() # 计算版本content timestamp parent_version 的SHA256前16位 content_str json.dumps(content, sort_keysTrue) raw f{content_str}{created_at}{parent_version or } version hashlib.sha256(raw.encode()).hexdigest()[:16] return ContextSnapshot( versionversion, created_atcreated_at, ttlttl, contentcontent, parent_versionparent_version ) app.task(bindTrue, max_retries3, default_retry_delay60) def handle_task(self, request_data: dict): 主控任务入口 try: # 1. 解析请求 req TaskRequest(**request_data) # 2. 生成初始上下文快照 initial_context { task_id: self.request.id, task_type: req.task_type, request_payload: req.payload, deadline_unix: time.time() 300, # 默认5分钟截止 retry_policy: {max_attempts: 3, backoff_factor: 2} } snapshot create_context_snapshot(initial_context, ttl3600) # 3. 拆解DAG简化版线性链 dag_nodes [] if req.task_type ORDER_SHIPMENT: # 步骤1检查库存 dag_nodes.append(TaskNode( node_idf{self.request.id}_check_stock, worker_typeinventory_checker, input_context{order_id: req.payload.get(order_id)}, timeout_sec10 )) # 步骤2调用物流 dag_nodes.append(TaskNode( node_idf{self.request.id}_call_logistics, worker_typelogistics_api, input_context{order_id: req.payload.get(order_id)}, timeout_sec15 )) # 4. 分发子任务 results [] for node in dag_nodes: # 为每个节点生成专属上下文快照 node_context { **snapshot.content, current_node: node.node_id, node_input: node.input_context } node_snapshot create_context_snapshot(node_context, snapshot.version) # 发送至Kafka此处简化为Celery任务 result dispatch_to_worker.delay( worker_typenode.worker_type, context_snapshotnode_snapshot.dict(), timeout_secnode.timeout_sec ) results.append(result) # 5. 聚合结果实际中用Chord或Canvas return {status: dispatched, subtasks: [r.id for r in results]} except Exception as exc: # 重试逻辑由Celery框架处理 raise self.retry(excexc) app.task def dispatch_to_worker(worker_type: str, context_snapshot: dict, timeout_sec: int): 分发到具体Worker实际中发Kafka消息 # 这里应发送Kafka消息简化为打印 print(f[Orchestrator] Dispatch to {worker_type}: {context_snapshot[version]}) return {dispatched: True, context_version: context_snapshot[version]}4.3 Worker实现沙箱化执行与结构化失败上报Worker必须极度轻量只做三件事校验输入、执行工具、结构化输出。失败时绝不能抛原始异常必须转换为标准错误码# worker.py from celery import Celery from pydantic import BaseModel, ValidationError from tenacity import retry, stop_after_attempt, wait_exponential import requests import logging app Celery(worker) app.conf.broker_url redis://localhost:6379/0 # 定义库存检查Worker的输入契约 class InventoryCheckInput(BaseModel): order_id: str sku_code: str class InventoryCheckOutput(BaseModel): available_quantity: int reserved_quantity: int status: str # IN_STOCK, LOW_STOCK, OUT_OF_STOCK class WorkerError(BaseModel): error_code: str # 如 INVENTORY_API_TIMEOUT, INVALID_SKU error_message: str retryable: bool # 是否可重试 app.task(bindTrue, max_retries3, default_retry_delay60) def inventory_checker(self, context_snapshot: dict): 库存检查Worker try: # 1. 校验输入上下文必须包含node_input if node_input not in context_snapshot[content]: raise ValueError(Missing node_input in context) # 2. 解析并校验输入数据 input_data InventoryCheckInput(**context_snapshot[content][node_input]) # 3. 执行工具调用带重试 retry( stopstop_after_attempt(3), waitwait_exponential(multiplier1, min2, max10) ) def call_inventory_api(): response requests.get( fhttp://inventory-api/check?sku{input_data.sku_code}order{input_data.order_id}, timeout5 ) response.raise_for_status() return response.json() api_result call_inventory_api() # 4. 校验输出并返回 output InventoryCheckOutput(**api_result) return { success: True, output: output.dict(), context_version: context_snapshot[version] } except ValidationError as e: # 输入校验失败不可重试立即上报 return { success: False, error: WorkerError( error_codeINPUT_VALIDATION_FAILED, error_messagestr(e), retryableFalse ).dict() } except requests.Timeout: # 超时可重试 raise self.retry(excException(Inventory API timeout)) except requests.RequestException as e: # 其他网络错误可重试 raise self.retry(exce) except Exception as e: # 未知错误记录日志不重试 logging.error(fUnexpected error in inventory_checker: {e}) return { success: False, error: WorkerError( error_codeUNKNOWN_ERROR, error_messagestr(e), retryableFalse ).dict() } # 启动Worker命令 # celery -A worker worker --loglevelinfo --concurrency44.4 Watcher实现事件驱动的状态监听Watcher不主动拉取而是被动监听Kafka事件流当检测到关键状态变更时向Orchestrator推送信号# watcher.py from kafka import KafkaConsumer from json import loads import threading import time class TaskWatcher: def __init__(self, bootstrap_servers: str, topic: str): self.consumer KafkaConsumer( topic, bootstrap_serversbootstrap_servers, value_deserializerlambda x: loads(x.decode(utf-8)), auto_offset_resetlatest, enable_auto_commitTrue, group_idwatcher-prod-group ) self.running False def start(self): self.running True t threading.Thread(targetself._poll_loop) t.daemon True t.start() def _poll_loop(self): while self.running: try: for message in self.consumer.poll(timeout_ms1000).values(): for record in message: self._handle_event(record.value) except Exception as e: print(fWatcher poll error: {e}) time.sleep(1) def _handle_event(self, event: dict): 处理事件库存更新、人工审核完成等 if event.get(event_type) INVENTORY_UPDATE: # 推送至Orchestrator简化为Celery任务 from orchestrator import handle_task handle_task.apply_async( kwargs{request_data: {task_type: REACTIVATE_TASK, payload: event}}, countdown0.1 # 延迟100ms避免与原任务竞争 ) # 使用示例 if __name__ __main__: watcher TaskWatcher( bootstrap_serverslocalhost:9092, topicprod.order.inventory_events ) watcher.start() print(Watcher started...) while True: time.sleep(3600) # 保持主线程存活5. 常见问题与排查技巧实录那些文档里不会写的血泪教训5.1 “Agent执行一半卡住日志里什么都没有”——如何定位无声失败这是最令人抓狂的问题。表面看Worker进程在跑但任务既不成功也不失败。我们的排查清单检查Celery Worker的prefetch_count默认值是4意味着Worker会预取4个任务到内存。如果某个任务执行时间极长比如生成PDF卡在字体渲染它会阻塞后续所有任务。解决方案celery -A worker worker --concurrency2 --prefetch-multiplier1验证Redis连接池耗尽高并发下大量Worker同时连接Redis超过max_connections限制。现象是Worker日志出现ConnectionError: Error 113 connecting to localhost:6379. No route to host.。解决方案在Celery配置中显式设置broker_pool_limit 1000检查上下文快照TTL如果Orchestrator生成的快照TTL设为300秒但Worker执行耗时600秒当Worker试图读取快照时数据库已将其标记为过期。解决方案TTL必须大于最大预期执行时间 × 2实操心得我们在每个Worker启动时强制执行一次“健康检查任务”该任务会模拟最耗时的操作如调用一次PDF生成并测量真实耗时动态调整TTL。代码片段# 在worker.py顶部 import time HEALTH_CHECK_DURATION 0 app.task def health_check(): global HEALTH_CHECK_DURATION start time.time() # 执行一次真实工具调用 requests.get(http://pdf-api/generate?templatetest, timeout30) HEALTH_CHECK_DURATION time.time() - start return HEALTH_CHECK_DURATION5.2 “上下文版本对不上Worker报错找不到对应快照”——分布式时钟同步陷阱在K8s集群中不同Node的系统时钟可能偏差达500ms。当Orchestrator在Node A生成快照timestamp1710000000.123Worker在Node B消费时timestamp1710000000.654两者计算的version完全不同。解决方案所有时间戳必须来自统一NTP源并在生成快照时强制使用time.time()而非datetime.now().timestamp()。后者在某些Linux发行版中受系统时区影响。更彻底的方案在Kafka消息头中增加server_timestamp字段由Producer所在服务器写入Consumer以此为准。我们线上系统强制要求所有Kafka Producer配置from kafka import KafkaProducer producer KafkaProducer( bootstrap_servers[localhost:9092], value_serializerlambda v: json.dumps(v).encode(utf-8), # 关键启用时间戳拦截器 interceptor_classes[TimestampInterceptor] ) class TimestampInterceptor: def on_send(self, record): record.headers.append((server_timestamp, str(time.time()).encode())) return record5.3 “任务失败后Adjudicator总是选错兜底路径”——知识库冷启动的实战方法Adjudicator依赖本地微调的Llama模型做失败归因但新业务上线时模型没见过“支付网关返回code503”这种错误。我们的冷启动三步法人工标注前100个失败样本运营同学从ELK日志中筛选典型失败标注根因如“503网关过载走备用通道”、“401token过期刷新重试”用LoRA微调Llama-3-8B仅训练attention层的adapter显存占用从80GB降到12GB2小时即可完成部署影子模式Shadow ModeAdjudicator同时输出“主模型决策”和“规则引擎决策”当两者不一致时将请求路由至人工审核队列并收集反馈。持续3天后模型准确率从62%提升至89%注意绝不允许Adjudicator直接修改生产数据。它的输出只能是{action: SWITCH_TO_BACKUP_CHANNEL, reason: 503 indicates gateway overload}具体执行仍由Orchestrator按预设规则完成。5.4 “Kafka消息积压Consumer Lag飙升”——Agent通信的容量规划公式消息积压不是配置问题而是容量规划失误。我们用这个公式预估所需Consumer数量N_consumers (QPS × avg_processing_time_sec × safety_factor) / 60其中safety_factor取2.0应对流量峰值avg_processing_time_sec取P95值。例如QPS500P95处理时间1.2秒则N_consumers (500 × 1.2 × 2) / 60 20。这意味着至少需要20个Consumer实例。线上监控必须包含kafka_consumer_lag每个Consumer Group的lagkafka_topic_partition_count确保Topic分区数 ≥ Consumer数celery_worker_active_tasksWorker实际并发数当kafka_consumer_lag 10000且持续5分钟自动触发告警并扩容Consumer。6. 我在实际项目中反复验证的三个关键认知第一个认知“智能”的价值不在首次成功率而在失败后的恢复效率。我们统计过一个设计良好的协同Agent系统首次执行成功率约82%但经过Adjudicator介入后的最终成功率可达99.97%。这意味着把精力花在“如何优雅失败”上比花在“如何一次成功”上回报率高得多。具体做法是为每个Worker定义清晰的失败分类网络层、业务层、数据层Adjudicator只针对业务层失败做智能决策网络层失败一律走重试数据层失败一律走人工。第二个认知上下文管理的终极目标不是“记住一切”而是“精准遗忘”。我们线上系统每天产生1200万上下文快照但99.3%在7天内被自动清理。真正需要长期保存的只有task_id、final_status、root_cause_code这三个字段。其他所有中间状态都是为了支撑“此刻”的决策决策完成后即失去价值。强行保留只会拖慢查询、增加存储成本、提高安全风险。第三个认知Agent协同的瓶颈永远不在AI模型而在基础设施的可观测性。当一个任务失败时工程师最需要的不是“模型认为原因是什么”而是“在哪个时间点、哪个Agent、执行了哪行代码、调用了哪个URL、收到了什么响应、上下文版本是多少”。我们投入了40%的开发时间在日志埋点、链路追踪、上下文快照索引上。没有这套可观测性基建再聪明的Agent也只是黑箱。最后分享一个小技巧在Orchestrator的handle_task函数开头强制添加一行日志logging.info(f[TASK_START] {self.request.id} | {req.task_type} | {json.dumps(req.payload)[:100]}...)这行日志能帮你瞬间定位90%的“任务没进来”问题——如果日志里没有这条记录说明请求根本没到达Orchestrator问题出在API网关或负载均衡层。别急着调大模型先看这一行。
返回列表