共计 2120 个字符,预计需要花费 6 分钟才能阅读完成。
在电商订单处理场景中,一个订单可能依次经过风控检测→库存锁定→支付处理→物流分配等多个 Agent 协作;在智能客服系统中,用户问题需经过意图识别→知识检索→话术生成等 Agent 流水线处理。这类业务天然具备流程化特征,但传统硬编码调用方式会导致代码臃肿、扩展困难——这正是链式调用的用武之地。
同步阻塞 vs 异步消息:决策树帮你选型
当业务满足以下条件时建议采用同步调用:
- 调用链深度≤3 层
- 单步骤耗时 <100ms
- 强依赖执行顺序
- 需要实时返回结果
反之则应考虑异步消息队列方案:
- 使用 RabbitMQ 等消息队列解耦 Agent
- 每个 Agent 消费上游消息并生产下游事件
- 通过 correlation_id 维护调用链上下文
(注:此处应为流程图,实际使用需替换真实 URL)
责任链模式的三层实现
基础结构
class Agent(ABC):
@abstractmethod
def set_next(self, agent: 'Agent') -> 'Agent':
pass
@abstractmethod
def handle(self, context: dict) -> Optional[dict]:
pass
class AbstractAgent(Agent):
_next_agent: Agent = None
def set_next(self, agent: Agent) -> Agent:
self._next_agent = agent
return agent
def handle(self, context: dict) -> Optional[dict]:
if self._next_agent:
return self._next_agent.handle(context)
return None
业务 Agent 实现
type PaymentAgent struct {next Agent}
func (p *PaymentAgent) SetNext(a Agent) {p.next = a}
func (p *PaymentAgent) Handle(ctx Context) (bool, error) {if err := validatePayment(ctx); err != nil {return false, err}
if p.next != nil {return p.next.Handle(ctx)
}
return true, nil
}
Redis 上下文共享方案
采用 hash 结构存储全链路上下文:
- 生成全局 chain_id 作为 Redis key
- 每个 Agent 读写独立 field
- 设置 TTL 防止内存泄漏
# 写操作示例
def save_context(chain_id: str, agent_name: str, data: dict):
r = redis.Redis()
r.hset(f"chain:{chain_id}", agent_name, json.dumps(data))
r.expire(f"chain:{chain_id}", 3600)
熔断机制双语言实现
Python 版基于 circuitbreaker:
from circuitbreaker import circuit
@circuit(failure_threshold=3, recovery_timeout=60)
def risk_check(ctx: dict):
# 风控逻辑...
Go 版使用 hystrix:
import "github.com/afex/hystrix-go/hystrix"
func init() {
hystrix.ConfigureCommand("inventory", hystrix.CommandConfig{
Timeout: 1000,
MaxConcurrentRequests: 100,
ErrorPercentThreshold: 25,
})
}
性能测试数据
测试环境:8 核 16G 云主机,MySQL 5.7,Redis 6.2
| 方案 | QPS | 99 线延迟 | CPU 使用率 |
|---|---|---|---|
| 同步调用 | 342 | 890ms | 78% |
| 异步队列 | 2104 | 210ms | 63% |
| 带缓存的同步 | 587 | 430ms | 69% |
生产环境避坑指南
上下文丢失三防措施
- 所有写操作增加 WAL 日志
- 重要字段采用 last-write-win 策略
- 定期扫描长时间未完成的 chain_id
循环调用检测
def detect_cycle(chain: list[Agent]) -> bool:
visited = set()
for agent in chain:
if id(agent) in visited:
return True
visited.add(id(agent))
return False
幂等性保障
- 每个请求携带唯一 request_id
- 关键操作记录执行状态
- 采用 CAS 方式更新结果
留给读者的思考题
- 如何在不重启服务的情况下,动态调整调用链顺序?比如根据流量特征把风控 Agent 从第一位移到第三位
- 当调用链跨多个微服务时,怎样实现类似 OpenTracing 的分布式追踪?特别是跨语言场景下的上下文传递
经过三个月的生产验证,这套方案成功将某金融系统的拒付处理流程从平均 2.3 秒缩短到 1.4 秒。最大的收获不是性能提升本身,而是发现良好的架构设计能让业务迭代速度提高 3 倍——当你需要新增反洗钱检查环节时,只需要编写新 Agent 并插入调用链,而不必修改任何现有代码。
正文完
