1. 多Agent协作架构的核心挑战
在构建多Agent系统时,我们常常会遇到一个看似矛盾的现象:单个Agent表现良好,但将它们组合起来后系统性能却断崖式下降。这种现象背后隐藏着分布式系统的基本特性——错误传播。当一个Agent产生微小错误时,这个错误会作为输入传递给下一个Agent,经过多次传递和放大后,最终可能导致整个系统偏离预期目标。
1.1 错误传播的数学本质
让我们用一个简单的数学模型来说明这个问题。假设一个系统包含5个Agent,每个Agent在每一步操作中都有95%的准确率。经过5步操作后,系统整体准确率会下降到约28%(0.95^5≈0.28)。这解释了为什么在实际部署中,即使每个独立Agent表现良好,整个系统仍可能出现灾难性失败。
提示:在实际工程中,错误传播的影响往往比理论计算更严重,因为Agent之间的错误通常不是独立的,而是会相互增强。
1.2 三大核心挑战解析
1.2.1 状态同步难题
Agent之间的状态同步是多系统协作的首要挑战。我们需要决定:
- 采用共享内存(高效但难以扩展)
- 消息传递(灵活但复杂度高)
- 文件系统(简单但实时性差)
每种方案都有其适用场景和代价。例如,在需要高频通信的场景中,共享内存可能是最佳选择;而在分布式环境中,消息队列往往更合适。
1.2.2 任务分发机制
任务分发策略直接影响系统效率和资源利用率。我们面临两个主要选择:
- 集中式调度(Orchestrator模式)
- 分布式协商(Peer-to-Peer模式)
集中式调度易于实现和管理,但可能成为性能瓶颈;分布式协商更具弹性,但实现复杂度显著提高。
1.2.3 阻塞与死锁问题
多Agent系统中,阻塞问题尤为棘手。常见场景包括:
- 循环等待:Agent A等待Agent B,而Agent B又在等待Agent A
- 资源竞争:多个Agent争抢同一资源
- 优先级反转:低优先级任务持有高优先级任务所需资源
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 主流协作模式深度解析
2.1 Orchestrator/Worker模式实践
这是目前最成熟的多Agent协作架构,其核心是一个中央调度器(Orchestrator)和多个工作节点(Worker)。让我们看一个完整的Python实现示例:
python复制class Orchestrator:
def __init__(self, workers: List[Worker]):
self.workers = workers
self.task_queue = asyncio.Queue()
self.result_map = {}
async def dispatch(self, task: Task):
"""任务分解与分发"""
subtasks = self._decompose_task(task)
task_id = str(uuid.uuid4())
# 初始化结果存储
self.result_map[task_id] = {
'subtasks': len(subtasks),
'results': [],
'event': asyncio.Event()
}
# 分发子任务
for i, subtask in enumerate(subtasks):
worker = self._select_worker(subtask)
await self.task_queue.put({
'task_id': task_id,
'subtask_id': i,
'subtask': subtask,
'worker': worker
})
# 等待所有子任务完成
await self.result_map[task_id]['event'].wait()
return self._aggregate_results(task_id)
async def _worker_loop(self, worker: Worker):
"""Worker处理循环"""
while True:
task_item = await self.task_queue.get()
try:
result = await worker.execute(task_item['subtask'])
self._handle_result(
task_item['task_id'],
task_item['subtask_id'],
result
)
except Exception as e:
self._handle_failure(
task_item['task_id'],
task_item['subtask_id'],
str(e)
)
finally:
self.task_queue.task_done()
这种架构的优势在于:
- 故障隔离:单个Worker失败不会影响整个系统
- 弹性扩展:可以动态增加Worker数量
- 状态集中:Orchestrator掌握全局状态,便于监控和调试
然而,它也存在明显缺点:
- Orchestrator成为单点故障
- 随着Worker数量增加,Orchestrator可能成为性能瓶颈
- Worker之间无法直接通信,必须通过Orchestrator中转
2.2 Peer-to-Peer模式实战
在需要Agent之间直接交互的场景中,P2P模式展现出独特优势。以下是一个代码审查场景的实现:
python复制class CodeReviewSystem:
def __init__(self):
self.coder = Agent(
role="Senior Developer",
system_prompt="你是一名经验丰富的Python开发工程师,擅长编写高效、清晰的代码。"
)
self.reviewer = Agent(
role="Quality Engineer",
system_prompt="你是一名严格的代码审查专家,专注于发现代码中的潜在问题和优化点。"
)
async def review_cycle(self, requirement: str, max_rounds=3):
current_code = None
review_history = []
for round in range(max_rounds):
# 开发阶段
if current_code is None:
dev_prompt = f"""根据以下需求编写Python代码:
{requirement}
要求:
1. 包含完整的类型注解
2. 添加适当的docstring
3. 考虑异常处理
"""
else:
dev_prompt = f"""根据审查意见修改代码:
审查意见:
{review_history[-1]}
原始代码:
{current_code}
"""
current_code = await self.coder.generate(dev_prompt)
# 审查阶段
review_prompt = f"""审查以下Python代码:
{current_code}
原始需求:
{requirement}
请从以下角度进行审查:
1. 功能是否符合需求
2. 是否有潜在bug
3. 性能是否可以优化
4. 代码风格是否一致
"""
review_result = await self.reviewer.generate(review_prompt)
review_history.append(review_result)
# 检查是否通过审查
if "没有发现问题" in review_result or "无需修改" in review_result:
break
return {
"final_code": current_code,
"review_rounds": round + 1,
"review_history": review_history
}
注意:在实践中,P2P模式需要特别注意防止"互相讨好"现象。可以通过在prompt中明确要求批判性反馈,或引入第三个Agent作为仲裁者来解决这个问题。
2.3 Pipeline模式优化策略
线性流水线模式虽然简单,但在实际应用中需要特别注意错误处理和性能优化。下面是一个增强版的Pipeline实现:
python复制class ResilientPipeline:
def __init__(self, stages: List[PipelineStage]):
self.stages = stages
self.timeout = 30 # 默认超时时间(秒)
self.retry_count = 2 # 默认重试次数
async def execute(self, input_data):
intermediate_results = []
current_input = input_data
for stage in self.stages:
last_error = None
for attempt in range(self.retry_count + 1):
try:
stage_start = time.time()
result = await asyncio.wait_for(
stage.process(current_input),
timeout=self.timeout
)
stage_time = time.time() - stage_start
intermediate_results.append({
"stage": stage.name,
"status": "success",
"time": stage_time,
"attempt": attempt + 1
})
current_input = result
break
except Exception as e:
last_error = e
intermediate_results.append({
"stage": stage.name,
"status": "failed",
"error": str(e),
"attempt": attempt + 1
})
if attempt == self.retry_count:
raise PipelineError(
f"Stage {stage.name} failed after {self.retry_count} retries"
) from last_error
# 指数退避重试
await asyncio.sleep(min(2 ** attempt, 10))
return {
"final_output": current_input,
"pipeline_stats": intermediate_results
}
关键优化点包括:
- 超时控制:防止单个stage阻塞整个pipeline
- 指数退避重试:提高临时性故障的恢复能力
- 详细执行日志:便于问题诊断和性能分析
3. Agent通信机制技术选型
3.1 四种通信方式对比分析
| 通信方式 | 延迟 | 吞吐量 | 可靠性 | 实现复杂度 | 适用场景 |
|---|---|---|---|---|---|
| 共享文件系统 | 高 | 高 | 高 | 低 | 异步批处理、大数据交换 |
| 消息队列 | 中 | 高 | 高 | 中 | 分布式系统、松耦合架构 |
| 共享内存 | 低 | 极高 | 低 | 高 | 高性能计算、单机多进程 |
| 结构化API | 中-高 | 中 | 高 | 高 | 跨语言集成、企业级系统 |
3.2 文件系统通信的工程实践
文件系统通信虽然简单,但在实际工程中需要注意多个细节问题。以下是一个生产级别的文件通信实现:
python复制class FileSystemComm:
def __init__(self, base_dir: Path, lock_timeout=30):
self.base_dir = base_dir
self.lock_timeout = lock_timeout
self.base_dir.mkdir(exist_ok=True)
async def write_result(self, task_id: str, data: dict):
"""原子化写入结果文件"""
temp_file = self.base_dir / f".tmp.{task_id}.{os.getpid()}"
target_file = self.base_dir / f"result.{task_id}.json"
try:
# 写入临时文件
with open(temp_file, 'w') as f:
json.dump(data, f)
# 原子重命名
temp_file.replace(target_file)
finally:
if temp_file.exists():
temp_file.unlink()
async def wait_for_result(self, task_id: str, poll_interval=1):
"""等待结果文件出现"""
result_file = self.base_dir / f"result.{task_id}.json"
start_time = time.time()
while not result_file.exists():
if time.time() - start_time > self.lock_timeout:
raise TimeoutError(f"Timeout waiting for {task_id}")
await asyncio.sleep(poll_interval)
with open(result_file, 'r') as f:
return json.load(f)
async def lock_operation(self, lock_name: str):
"""基于文件的分布式锁"""
lock_file = self.base_dir / f"lock.{lock_name}"
start_time = time.time()
while True:
try:
# 尝试创建锁文件(原子操作)
fd = os.open(lock_file, os.O_CREAT | os.O_EXCL | os.O_WRONLY)
with os.fdopen(fd, 'w') as f:
f.write(str(os.getpid()))
yield
return
except FileExistsError:
if time.time() - start_time > self.lock_timeout:
raise TimeoutError(f"Lock {lock_name} timeout")
await asyncio.sleep(0.1)
finally:
if lock_file.exists():
try:
lock_file.unlink()
except:
pass
关键设计考虑:
- 原子写入:通过"写入临时文件+原子重命名"避免数据损坏
- 分布式锁:基于文件系统实现简单的进程间协调
- 超时处理:防止无限等待导致的系统挂起
3.3 消息队列集成方案
对于需要高吞吐量的场景,消息队列是更合适的选择。以下是使用RabbitMQ的集成示例:
python复制class MessageQueueComm:
def __init__(self, amqp_url: str):
self.connection = None
self.channel = None
self.amqp_url = amqp_url
async def connect(self):
"""建立RabbitMQ连接"""
self.connection = await aio_pika.connect_robust(self.amqp_url)
self.channel = await self.connection.channel()
async def publish_task(self, queue_name: str, task_data: dict):
"""发布任务到指定队列"""
if not self.channel:
await self.connect()
queue = await self.channel.declare_queue(
queue_name,
durable=True
)
message = aio_pika.Message(
body=json.dumps(task_data).encode(),
delivery_mode=aio_pika.DeliveryMode.PERSISTENT
)
await self.channel.default_exchange.publish(
message,
routing_key=queue_name
)
async def consume_tasks(self, queue_name: str, callback: Callable):
"""消费队列任务"""
if not self.channel:
await self.connect()
queue = await self.channel.declare_queue(
queue_name,
durable=True
)
async with queue.iterator() as queue_iter:
async for message in queue_iter:
try:
async with message.process():
task_data = json.loads(message.body.decode())
result = await callback(task_data)
# 如果需要返回结果
if message.reply_to:
reply_message = aio_pika.Message(
body=json.dumps(result).encode(),
correlation_id=message.correlation_id
)
await self.channel.default_exchange.publish(
reply_message,
routing_key=message.reply_to
)
except Exception as e:
logging.error(f"Task processing failed: {e}")
# 根据业务决定是否重新入队
4. 任务调度与资源管理
4.1 动态任务调度算法
高效的动态任务调度是多Agent系统的核心。以下是一个考虑多种因素的调度器实现:
python复制class DynamicScheduler:
def __init__(self, agents: List[Agent]):
self.agents = agents
self.agent_stats = {
agent.id: {
'last_used': 0,
'active_tasks': 0,
'success_rate': 1.0,
'avg_latency': 0
}
for agent in agents
}
self.lock = asyncio.Lock()
async def assign_task(self, task: Task) -> Agent:
"""基于多种因素选择最优Agent"""
async with self.lock:
available_agents = [
agent for agent in self.agents
if self._can_handle(agent, task)
and self.agent_stats[agent.id]['active_tasks'] < agent.max_concurrency
]
if not available_agents:
raise NoAvailableAgentError(task)
# 计算每个Agent的得分
scored_agents = []
current_time = time.time()
for agent in available_agents:
stats = self.agent_stats[agent.id]
# 空闲时间得分(优先使用空闲时间长的)
idle_score = current_time - stats['last_used']
# 成功率得分
success_score = stats['success_rate']
# 延迟得分(优先选择平均延迟低的)
latency_score = 1 / (stats['avg_latency'] + 0.1)
# 负载均衡得分(优先选择当前任务少的)
load_score = 1 / (stats['active_tasks'] + 1)
# 综合得分(可调整权重)
total_score = (
0.3 * idle_score +
0.3 * success_score +
0.2 * latency_score +
0.2 * load_score
)
scored_agents.append((total_score, agent))
# 选择得分最高的Agent
_, best_agent = max(scored_agents, key=lambda x: x[0])
# 更新Agent状态
self.agent_stats[best_agent.id]['active_tasks'] += 1
self.agent_stats[best_agent.id]['last_used'] = current_time
return best_agent
async def update_agent_stats(self, agent_id: str, success: bool, latency: float):
"""更新Agent性能指标"""
async with self.lock:
stats = self.agent_stats[agent_id]
stats['active_tasks'] = max(0, stats['active_tasks'] - 1)
# 更新成功率(指数移动平均)
stats['success_rate'] = 0.9 * stats['success_rate'] + 0.1 * float(success)
# 更新延迟(指数移动平均)
stats['avg_latency'] = 0.8 * stats['avg_latency'] + 0.2 * latency
该调度器考虑以下因素:
- Agent空闲时间:避免某些Agent长期闲置
- 历史成功率:优先选择更可靠的Agent
- 平均延迟:优化系统响应时间
- 当前负载:实现负载均衡
4.2 任务依赖图执行引擎
复杂任务通常包含多个依赖步骤,需要有向无环图(DAG)调度。以下是DAG调度器的实现:
python复制class DAGExecutor:
def __init__(self, scheduler: DynamicScheduler):
self.scheduler = scheduler
self.task_graph = nx.DiGraph()
def add_task(self, task: Task, dependencies: List[Task] = None):
"""添加任务及其依赖关系"""
self.task_graph.add_node(task.id, task=task)
if dependencies:
for dep in dependencies:
self.task_graph.add_edge(dep.id, task.id)
if not nx.is_directed_acyclic_graph(self.task_graph):
raise ValueError("Task graph contains cycles")
async def execute(self):
"""执行DAG任务"""
completed = {}
in_progress = set()
# 拓扑排序确定执行顺序
execution_order = list(nx.topological_sort(self.task_graph))
async def run_task(task_id):
# 等待所有前置任务完成
for pred in self.task_graph.predecessors(task_id):
await asyncio.shield(completed[pred]['event'].wait())
task = self.task_graph.nodes[task_id]['task']
try:
# 选择Agent执行任务
agent = await self.scheduler.assign_task(task)
in_progress.add(task_id)
start_time = time.time()
result = await agent.execute(task)
latency = time.time() - start_time
await self.scheduler.update_agent_stats(
agent.id, True, latency
)
return {
'status': 'completed',
'result': result,
'latency': latency
}
except Exception as e:
await self.scheduler.update_agent_stats(
agent.id, False, 0
)
return {
'status': 'failed',
'error': str(e)
}
finally:
in_progress.remove(task_id)
# 为每个任务创建完成事件
for task_id in execution_order:
completed[task_id] = {
'event': asyncio.Event(),
'result': None
}
# 并行执行可并行任务
pending_tasks = {
task_id: asyncio.create_task(run_task(task_id))
for task_id in execution_order
}
# 等待所有任务完成并设置事件
for task_id, task_future in pending_tasks.items():
result = await task_future
completed[task_id]['result'] = result
completed[task_id]['event'].set()
return {
task_id: completed[task_id]['result']
for task_id in execution_order
}
关键特性:
- 基于拓扑排序的任务调度
- 自动处理任务依赖关系
- 并行执行独立任务
- 完善的错误处理和状态跟踪
5. 错误处理与系统健壮性
5.1 断路器模式实现
断路器模式是防止级联失败的关键技术。以下是Python实现:
python复制class CircuitBreaker:
def __init__(self, max_failures=3, reset_timeout=60):
self.max_failures = max_failures
self.reset_timeout = reset_timeout
self.failure_count = 0
self.last_failure_time = 0
self.state = 'closed' # closed, open, half-open
self.lock = threading.Lock()
async def execute(self, coro):
"""在断路器保护下执行协程"""
if self.state == 'open':
# 检查是否应该尝试恢复
if time.time() - self.last_failure_time > self.reset_timeout:
with self.lock:
if self.state == 'open':
self.state = 'half-open'
else:
raise CircuitOpenError("Circuit is open")
try:
result = await coro
if self.state == 'half-open':
with self.lock:
if self.state == 'half-open':
self.state = 'closed'
self.failure_count = 0
return result
except Exception as e:
with self.lock:
self.failure_count += 1
self.last_failure_time = time.time()
if self.failure_count >= self.max_failures:
self.state = 'open'
raise
def __call__(self, coro_func):
"""装饰器用法"""
async def wrapper(*args, **kwargs):
return await self.execute(coro_func(*args, **kwargs))
return wrapper
使用示例:
python复制@CircuitBreaker(max_failures=3, reset_timeout=60)
async def call_agent_b(task):
return await agent_b.execute(task)
try:
result = await call_agent_b(important_task)
except CircuitOpenError:
# 断路器打开时的降级处理
result = fallback_operation()
5.2 超时与重试策略
合理的超时和重试策略可以显著提高系统稳定性:
python复制class RetryPolicy:
def __init__(self,
max_retries=3,
initial_delay=0.1,
max_delay=10,
backoff_factor=2,
timeout=30):
self.max_retries = max_retries
self.initial_delay = initial_delay
self.max_delay = max_delay
self.backoff_factor = backoff_factor
self.timeout = timeout
async def execute(self, coro_func, *args, **kwargs):
"""执行带重试策略的操作"""
last_error = None
current_delay = self.initial_delay
for attempt in range(self.max_retries + 1):
try:
# 带超时的执行
return await asyncio.wait_for(
coro_func(*args, **kwargs),
timeout=self.timeout
)
except Exception as e:
last_error = e
if attempt == self.max_retries:
break
# 指数退避
await asyncio.sleep(current_delay)
current_delay = min(
current_delay * self.backoff_factor,
self.max_delay
)
raise MaxRetriesExceededError(
f"Operation failed after {self.max_retries} retries"
) from last_error
最佳实践建议:
- 为不同操作设置不同的超时时间(如CPU密集型操作和I/O操作)
- 对于非幂等操作要谨慎使用重试
- 结合断路器模式使用,避免无限制重试失败操作
6. 性能监控与调试技巧
6.1 分布式追踪实现
在多Agent系统中,端到端的请求追踪至关重要。以下是一个轻量级实现:
python复制class TraceContext:
def __init__(self, trace_id=None, parent_span_id=None):
self.trace_id = trace_id or str(uuid.uuid4())
self.span_id = str(uuid.uuid4())
self.parent_span_id = parent_span_id
self.start_time = time.time()
self.tags = {}
def child_context(self):
"""创建子追踪上下文"""
return TraceContext(
trace_id=self.trace_id,
parent_span_id=self.span_id
)
def to_dict(self):
return {
'trace_id': self.trace_id,
'span_id': self.span_id,
'parent_span_id': self.parent_span_id,
'start_time': self.start_time,
'duration': time.time() - self.start_time,
'tags': self.tags
}
class Tracer:
def __init__(self, exporter=None):
self.exporter = exporter or LoggingExporter()
async def trace(self, operation_name, context=None):
"""追踪一个操作的执行"""
context = context or TraceContext()
context.tags['operation'] = operation_name
try:
start_time = time.time()
yield context
except Exception as e:
context.tags['error'] = str(e)
raise
finally:
context.tags['duration'] = time.time() - start_time
await self.exporter.export(context.to_dict())
使用示例:
python复制async def process_task(task):
async with tracer.trace("process_task") as ctx:
ctx.tags['task_id'] = task.id
result = await agent.process(task)
ctx.tags['result_size'] = len(result)
return result
6.2 关键性能指标收集
以下是要监控的核心指标:
- 每个Agent的请求处理时间(P50, P90, P99)
- 任务队列长度
- 错误率和错误类型分布
- 资源利用率(CPU、内存、网络)
- 消息延迟(生产到消费的时间差)
Prometheus监控示例:
python复制from prometheus_client import Summary, Counter, Gauge
# 定义指标
REQUEST_TIME = Summary(
'agent_request_processing_seconds',
'Time spent processing request',
['agent_type']
)
REQUEST_COUNT = Counter(
'agent_requests_total',
'Total requests processed',
['agent_type', 'status']
)
QUEUE_SIZE = Gauge(
'task_queue_size',
'Current number of tasks in queue'
)
# 在关键位置记录指标
@REQUEST_TIME.time()
async def process_with_metrics(agent, task):
try:
result = await agent.process(task)
REQUEST_COUNT.labels(
agent_type=agent.type,
status='success'
).inc()
return result
except Exception as e:
REQUEST_COUNT.labels(
agent_type=agent.type,
status='failed'
).inc()
raise
7. 安全设计与访问控制
7.1 Agent认证机制
在分布式环境中,必须确保只有授权的Agent可以参与系统:
python复制class AgentAuthenticator:
def __init__(self, shared_secret):
self.shared_secret = shared_secret
def generate_token(self, agent_id, timestamp=None):
"""生成认证token"""
timestamp = timestamp or int(time.time())
msg = f"{agent_id}:{timestamp}".encode()
signature = hmac.new(
self.shared_secret.encode(),
msg,
hashlib.sha256
).hexdigest()
return f"{agent_id}:{timestamp}:{signature}"
def validate_token(self, token, max_age=60):
"""验证token有效性"""
try:
agent_id, timestamp_str, signature = token.split(':')
timestamp = int(timestamp_str)
# 检查时间有效性
if abs(time.time() - timestamp) > max_age:
return False
# 验证签名
expected = self.generate_token(agent_id, timestamp)
return hmac.compare_digest(token, expected)
except:
return False
7.2 最小权限原则实施
为每个Agent定义明确的权限边界:
python复制class PolicyEngine:
def __init__(self):
self.policies = {}
def add_policy(self, agent_type, resource, actions):
"""添加访问控制策略"""
if agent_type not in self.policies:
self.policies[agent_type] = {}
self.policies[agent_type][resource] = actions
def check_permission(self, agent_type, resource, action):
"""检查是否允许访问"""
return (
agent_type in self.policies and
resource in self.policies[agent_type] and
action in self.policies[agent_type][resource]
)
# 示例策略配置
policy_engine = PolicyEngine()
policy_engine.add_policy(
agent_type='data_processor',
resource='customer_data',
actions=['read', 'analyze']
)
policy_engine.add_policy(
agent_type='report_generator',
resource='customer_data',
actions=['aggregate']
)
8. 测试策略与质量保证
8.1 单元测试框架
针对Agent的核心功能进行单元测试:
python复制class AgentTestCase(unittest.TestCase):
def setUp(self):
self.agent = Agent(system_prompt="你是一个有帮助的助手")
async def test_simple_response(self):
response = await self.agent.generate("你好")
self.assertIsInstance(response, str)
self.assertGreater(len(response), 0)
async def test_error_handling(self):
with self.assertRaises(AgentError):
await self.agent.generate("" * 10000) # 触发输入过长错误
8.2 集成测试方案
验证多个Agent的协作行为:
python复制class IntegrationTest(unittest.IsolatedAsyncioTestCase):
async def asyncSetUp(self):
self.orchestrator = Orchestrator()
self.worker1 = Worker(capabilities=['analysis'])
self.worker2 = Worker(capabilities=['formatting'])
await self.orchestrator.register_worker(self.worker1)
await self.orchestrator.register_worker(self.worker2)
async def test_pipeline_execution(self):
task = Task(
id="test1",
payload="Analyze and format this data",
steps=['analysis', 'formatting']
)
result = await self.orchestrator.process(task)
self.assertIn('analysis', result)
self.assertIn('formatted', result)
async def test_failure_recovery(self):
faulty_worker = Worker(capabilities=['analysis'], fail_prob=0.9)
await self.orchestrator.register_worker(faulty_worker)
task = Task(
id="test2",
payload="Analyze this data",
steps=['analysis']
)
# 应该自动重试直到成功
result = await self.orchestrator.process(task)
self.assertIn('analysis', result)
9. 部署架构与扩展策略
9.1 容器化部署方案
使用Docker部署多Agent系统的示例配置:
dockerfile复制# Agent基础镜像
FROM python:3.9-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
# 不同类型的Agent使用不同的entrypoint
CMD ["python", "worker.py"]
使用Docker Compose编排:
yaml复制version: '3'
services:
orchestrator:
build: .
command: python orchestrator.py
environment:
- REDIS_HOST=redis
ports:
- "8000:8000"
depends_on:
- redis
worker_analysis:
build: .
command: python worker.py --type analysis
environment:
- REDIS_HOST=redis
scale: 3
worker_formatting:
build: .
command: python worker.py --type formatting
environment:
- REDIS_HOST=redis
scale: 2
redis:
image: redis:alpine
ports:
- "6379:6379"
9.2 自动扩展策略
基于负载的自动扩展实现:
python复制class AutoScaler:
def __init__(self, min_workers=1, max_workers=10, scale_up_threshold=0.8):
self.min_workers = min_workers
self.max_workers = max_workers
self.scale_up_threshold = scale_up_threshold
self.current_workers = min_workers
async def monitor_and_scale(self, queue):
"""监控队列并调整worker数量"""
while True:
queue_size = await queue.size()
utilization = queue_size / (self.current_workers * 5) # 假设每个worker可处理5个任务
if utilization > self.scale_up_threshold and self.current_workers < self.max_workers:
await self.scale_out()
elif utilization < 0.2 and self.current_workers > self.min_workers:
await self.scale_in()
await asyncio.sleep(10)
async def scale_out(self):
"""增加worker实例"""
# 实际实现可能调用Kubernetes API或Docker API
self.current_workers += 1
logging.info(f"Scaling out to {self.current_workers} workers")
async def scale_in(self):
"""减少worker实例"""
self.current_workers -= 1
logging.info(f"Scaling in to {self.current_workers} workers")
10. 成本优化与性能权衡
10.1 Token使用优化
在多Agent系统中,LLM调用的token消耗是主要成本来源。优化策略包括:
- 上下文压缩:定期总结对话历史而非完整保留
- 分层处理:简单任务使用小模型,复杂任务才用大模型
- 缓存机制:缓存常见问题的响应
python复制class ConversationCache:
def __init__(self, max_size=1000):
self.cache = {}
self.max_size = max_size
self.lock = asyncio.Lock()
async def get_response(self, conversation_hash):
"""获取缓存的响应"""
async with self.lock:
return self.cache.get(conversation_hash)
async def store_response(self, conversation_hash, response):
"""存储响应到缓存"""
async with self.lock:
if len(self.cache) >= self.max_size:
# 简单的LRU淘汰
self.cache.pop(next(iter(self.cache)))
self.cache[conversation_hash] = response
def hash_conversation(messages):
"""计算对话内容的哈希值"""
return hashlib.md5(
json.dumps(messages, sort_keys=True).encode()
).hexdigest()
class CachedAgent:
def __init__(self, agent, cache):
self.agent = agent
self.cache = cache
async def generate(self, messages):
# 检查缓存
conv_hash = hash_conversation(messages)
cached = await self.cache.get_response(conv_hash)
if cached:
return cached
# 调用真实Agent
response = await self.agent.generate(messages)
# 存储到缓存
await self.cache.store_response(conv_hash, response)
return response
10.2 混合精度执行
根据任务需求动态选择模型精度:
python复制class HybridAgent:
def __init__(self, fast_model, accurate_model, router_model):
self.fast_model = fast_model # 小模型,响应快
self.accurate_model = accurate_model # 大模型,质量高
self.router_model = router_model # 路由决策模型
