1. 项目背景与核心思路
去年在重构一个遗留系统时,我遇到了一个典型的技术债务问题:项目中有个3000多行的Python模块,负责处理各种业务逻辑的编排和调度。这个模块不仅维护困难,每次修改都像在拆炸弹,而且执行效率也随着业务增长越来越低。
经过分析,我发现这个模块的核心功能其实可以抽象为一种"智能工作流引擎"的需求——根据输入动态决定执行路径,处理异常,并维护执行上下文。这让我联想到当时刚发布的Claude Code框架,它提供的Agent机制正好能解决这类问题。
但直接引入Claude Code意味着:
- 需要额外部署服务
- 引入新的技术栈
- 增加系统复杂度
于是我开始思考:能否用Python原生特性实现一个轻量级的Agent核心?经过两周的探索,最终用40行代码实现了核心调度逻辑,替代了原来3000行的"意大利面条式"代码。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 关键技术实现解析
2.1 状态机设计
Agent的核心是一个状态机,我使用Python的enum和dataclass实现:
python复制from enum import Enum, auto
from dataclasses import dataclass
from typing import Any, Callable, Optional
class AgentState(Enum):
IDLE = auto()
PROCESSING = auto()
WAITING = auto()
COMPLETED = auto()
FAILED = auto()
@dataclass
class AgentContext:
state: AgentState = AgentState.IDLE
memory: dict[str, Any] = field(default_factory=dict)
current_task: Optional[str] = None
这个设计有几点精妙之处:
- 使用enum明确限定状态值,避免魔法字符串
- dataclass自动生成__init__方法,简化对象创建
- 类型注解让IDE能提供更好的代码补全
2.2 消息调度器
Agent的核心是一个不到20行的消息处理器:
python复制class PyAgent:
def __init__(self):
self._handlers = {}
self.ctx = AgentContext()
def register_handler(self, msg_type: str, handler: Callable):
self._handlers[msg_type] = handler
async def process(self, message: dict):
if not (handler := self._handlers.get(message['type'])):
raise ValueError(f"No handler for {message['type']}")
self.ctx.current_task = message.get('task_id')
try:
self.ctx.state = AgentState.PROCESSING
return await handler(message, self.ctx)
except Exception as e:
self.ctx.state = AgentState.FAILED
raise
finally:
self.ctx.state = AgentState.COMPLETED
关键设计点:
- 使用字典实现的消息路由,避免大量if-else
- 协程支持(async/await)确保高并发能力
- 自动状态管理,业务代码无需关心状态转换
2.3 工具集成层
为了兼容现有系统,我实现了一个工具适配层:
python复制class ToolManager:
def __init__(self):
self._tools = {}
def register_tool(self, name: str, tool: Callable):
self._tools[name] = tool
async def execute(self, tool_name: str, params: dict):
if tool := self._tools.get(tool_name):
return await tool(**params)
raise ValueError(f"Unknown tool: {tool_name}")
这个设计实现了:
- 统一工具调用接口
- 自动参数解包
- 异步执行支持
3. 完整实现与集成
将上述组件组合起来,就得到了完整的微型Agent框架:
python复制from contextlib import asynccontextmanager
class MiniAgent:
def __init__(self):
self.dispatcher = PyAgent()
self.tools = ToolManager()
# 注册内置处理器
self.dispatcher.register_handler('tool_call', self._handle_tool_call)
self.dispatcher.register_handler('user_input', self._handle_user_input)
async def _handle_tool_call(self, msg, ctx):
result = await self.tools.execute(msg['tool'], msg.get('params', {}))
ctx.memory[msg.get('store_as', 'last_result')] = result
return {'status': 'success', 'result': result}
async def _handle_user_input(self, msg, ctx):
ctx.memory['last_input'] = msg['content']
return {'status': 'received'}
@asynccontextmanager
async def session(self):
try:
self.dispatcher.ctx.state = AgentState.IDLE
yield self
finally:
self.dispatcher.ctx.state = AgentState.COMPLETED
使用示例:
python复制agent = MiniAgent()
# 注册业务工具
@agent.tools.register_tool('analyze_data')
async def analyze(data_source: str):
# 实现具体业务逻辑
return {"insight": "..."}
# 使用Agent
async with agent.session() as sess:
result = await sess.dispatcher.process({
'type': 'tool_call',
'tool': 'analyze_data',
'params': {'data_source': 'sales.csv'},
'store_as': 'analysis_result'
})
4. 性能优化技巧
在实际使用中,我总结了几个关键优化点:
4.1 内存管理
对于长时间运行的Agent,需要注意内存泄漏问题:
python复制class AgentContext:
def __init__(self):
self._memory = {}
self._max_memory = 1000 # 限制内存条目数
@property
def memory(self):
return self._memory
def cleanup(self):
if len(self._memory) > self._max_memory:
# LRU清理
self._memory = dict(sorted(
self._memory.items(),
key=lambda x: x[1]['_timestamp']
)[-self._max_memory:])
4.2 超时控制
为工具调用添加超时保护:
python复制from asyncio import TimeoutError, wait_for
class ToolManager:
async def execute(self, tool_name: str, params: dict, timeout=30):
try:
return await wait_for(
self._tools[tool_name](**params),
timeout=timeout
)
except TimeoutError:
log.warning(f"Tool {tool_name} timeout")
raise
4.3 批处理优化
对于高频小消息,可以实现批处理:
python复制class PyAgent:
def __init__(self):
self._batch = []
self._batch_size = 10
self._batch_interval = 0.1 # 秒
async def _batch_processor(self):
while True:
if len(self._batch) >= self._batch_size:
await self._process_batch()
await asyncio.sleep(self._batch_interval)
async def _process_batch(self):
tasks = [self.process(msg) for msg in self._batch]
self._batch.clear()
return await asyncio.gather(*tasks)
5. 实际应用案例
在我们的订单处理系统中,原本有这样一个复杂流程:
python复制def process_order(order):
if not validate(order):
return error("Invalid order")
inventory = check_inventory(order.items)
if not inventory['all_available']:
if order['allow_partial']:
order = adjust_order(order, inventory)
else:
return error("Out of stock")
payment = process_payment(order)
if not payment['success']:
retry = handle_payment_failure(order, payment)
if not retry:
return error("Payment failed")
# 还有10多个其他步骤...
改用Agent实现后:
python复制@agent.tools.register_tool('process_order')
async def process_order(order: dict, ctx):
await ctx.agent.process({'type': 'validate', 'order': order})
await ctx.agent.process({'type': 'check_inventory', 'order': order})
await ctx.agent.process({'type': 'process_payment', 'order': order})
# 其他步骤...
优势对比:
- 代码量从300+行减少到40行核心逻辑
- 每个步骤成为独立可测试的组件
- 状态管理变得明确且可追踪
- 可以动态添加/修改处理流程
6. 扩展与演进
这个微型框架后来发展出了几个重要扩展:
6.1 持久化支持
python复制class PersistentAgent(MiniAgent):
def __init__(self, storage):
super().__init__()
self.storage = storage
self.dispatcher.register_handler('persist', self._handle_persist)
async def _handle_persist(self, msg, ctx):
await self.storage.save(
key=msg['key'],
value=msg['value']
)
6.2 分布式支持
python复制class DistributedAgent(MiniAgent):
def __init__(self, rpc_client):
super().__init__()
self.rpc = rpc_client
self.dispatcher.register_handler('rpc_call', self._handle_rpc)
async def _handle_rpc(self, msg, ctx):
return await self.rpc.call(
msg['method'],
msg.get('args', [])
)
6.3 监控集成
python复制class MonitoredAgent(MiniAgent):
def __init__(self, metrics):
super().__init__()
self.metrics = metrics
async def process(self, message):
with self.metrics.timer('agent.process'):
return await super().process(message)
7. 经验总结
在实现和使用这个微型Agent框架的过程中,我总结了以下几点关键经验:
- 状态隔离:每个会话应该有独立的上下文,避免状态污染
- 错误处理:在框架层统一处理异常,业务代码只需关注成功路径
- 性能监控:关键路径添加指标收集,便于后期优化
- 文档驱动:为每个工具自动生成使用文档,降低使用门槛
一个典型的错误处理改进示例:
python复制class RobustAgent(MiniAgent):
async def process(self, message):
try:
return await super().process(message)
except Exception as e:
self.ctx.state = AgentState.FAILED
self.ctx.last_error = str(e)
# 自动重试逻辑
if self._should_retry(e):
await asyncio.sleep(1)
return await super().process(message)
raise
这个40行的Python实现虽然简单,但通过合理的抽象和扩展,成功替代了我们系统中多个复杂的业务流程模块。它不仅减少了代码量,还提高了系统的可维护性和扩展性。最重要的是,它证明了有时候简单的设计反而能解决复杂的问题。
