创新实训:多agent系统总体架构的设计与落地
从代码到落地:基于LLM的多Agent工作流系统设计与实现可以很自然、妥帖地联系到最近所做的尝试:用大语言模型(LLM)来构建一套多Agent协同工作流系统,应用于电商评论数据分析、营销文案生成、差评预警诸种场景,在VS Code中梳理代码时从单功能接口出发,直至完整的Agent框架,过程中有若干值得系统总结的要点,因此本文打算从技术架构、核心设计、代码实现三个维度对此进行清晰阐述。
一、需求与架构思考:为什么要做多Agent?
最初的需求十分明确:分析电商评论数据,自动生成营销文案,监控差评率并预警。若直接写一个大函数把所有逻辑都塞进去,短期确实能跑通,但是扩展性极差,稍后要加入“竞品分析Agent”“报告生成Agent”时,代码就会变得极其臃肿。因此核心的思考点十分明确:
1.单一职责:把不同业务能力合理、明确地拆分成独立的Agent(评论分析Agent、文案生成Agent、预警Agent等),各Agent各就其位,各司其职。
2 .异步执行:由于Agent执行分析10万条评论的过程必然耗时,因此首先要做异步化处理,以免阻塞接口。
3.状态可追踪:即用户要能清楚地看到Agent的执行进度、结果及失败原因,第三宜用WebSocket实时推送执行状态以改善用户体验。
4.数据持久化:所有任务、报告、预警记录都应做数据持久化处理,便于以后查询。
二、核心模块设计与代码实现
1. 数据模型层:ORM 设计(agent_task_model.py)
要系统、有层次地讨论核心模块设计及其实现。首先从数据模型层入手,即ORM设计(agent_task_model. py),先解决“数据存什么、怎么存”的基本问题,自然、合理地设计任务、报告、预警三个核心模型,并用SQLAlchemy做ORM映射:
# agent_task_model.py 核心代码
from sqlalchemy import Column, String, Integer, Float, DateTime, Text, ForeignKey
from app.core.database import Base
class AgentTaskModel(Base):
"""多Agent工作流任务:记录任务状态、输入、结果、错误信息"""
__tablename__ = "agent_tasks"
id = Column(String(36), primary_key=True, default=lambda: str(uuid.uuid4()))
dataset_id = Column(String(36), ForeignKey("datasets.id"), nullable=False, index=True)
user_input = Column(Text) # 用户输入的任务指令
status = Column(String(20), default="pending") # pending/running/done/failed
result_json = Column(Text) # 任务结果(JSON字符串)
error_message = Column(Text) # 失败时的错误信息
created_at = Column(DateTime, default=datetime.utcnow)
updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
class ReportModel(Base):
"""报告记录:存储自动生成的周报/月报/竞品分析报告"""
__tablename__ = "reports"
id = Column(String(36), primary_key=True, default=lambda: str(uuid.uuid4()))
dataset_id = Column(String(36), ForeignKey("datasets.id"), nullable=False, index=True)
template_type = Column(String(50)) # weekly/monthly/competitor
title = Column(String(200))
status = Column(String(20), default="pending")
file_path_docx = Column(String(500)) # 生成的文档路径
file_path_pdf = Column(String(500))
class Alert(Base):
"""差评率预警:监控差评率超过阈值时记录并告警"""
__tablename__ = "alerts"
id = Column(String(36), primary_key=True, default=lambda: str(uuid.uuid4()))
dataset_id = Column(String(36), ForeignKey("datasets.id"), nullable=False, index=True)
alert_type = Column(String(50), default="negative_rate")
threshold = Column(Float) # 预警阈值
current_value = Column(Float) # 当前差评率
message = Column(Text) # 预警信息
triggered_at = Column(DateTime, default=datetime.utcnow)
设计思考:
- 所有表都用UUID作为主键,由此自然解决分布式环境下自增ID的种种问题
-“AgentTaskModel”是整个多Agent工作流生命周期的核心记录体
- 状态字段(status)统一用字符串枚举(pending/running/done/failed),有利于前端展示及后端逻辑判断
- 结果用“Text”类型存储JSON字符串,能很好地容纳多Agent返回的各类结构化、甚至格式各异的结果。
2.核心业务层:多 Agent 调度与异步执行(agent.py)
这是系统的核心,我把 Agent 相关的接口都集中在agent.py里,主要包含:主要包含:任务启动、异步执行、状态监听、结果查询、营销文案生成等功能。
(1)异步任务启动:避免阻塞
由于用户调用`/run`接口启动多Agent工作流时不宜阻塞,因此先创建任务记录,将核心逻辑放到后台任务(BackgroundTasks)中执行,马上返回`task_id`给用户,由用户用`task_id`去监听任务状态。
# agent.py /run 接口核心代码
@router.post("/run", summary="启动多Agent工作流")
async def run_agents(
request: AgentRunRequest,
background_tasks: BackgroundTasks,
db: AsyncSession = Depends(get_db),
):
"""启动多Agent协同工作流,返回task_id监听状态"""
# 1. 创建任务记录,初始状态pending
task = AgentTaskModel(
dataset_id=request.dataset_id,
user_input=request.user_input,
status="pending",
)
db.add(task)
await db.commit()
await db.refresh(task)
# 2. 后台执行Agent工作流(不阻塞接口)
background_tasks.add_task(_do_agent_run, str(task.id), request, db)
return {"task_id": str(task.id), "status": "pending", "message": "多Agent工作流已启动"}
设计思考:
用FastAPI的`BackgroundTasks`来做轻量级异步处理(如需更复杂的场景,可自然、合理地替换为Celery)。具体做法是:任务创建之后立刻提交数据库,因此用户可以即时查询任务状态,同时把核心执行逻辑抽离成`_do_agent_run`函数,这样更利于维护、测试。
(2)多Agent执行核心逻辑
`_do_agent_run`函数实质上是多Agent调度的正式入口:
1.把任务状态改为`running`
2.从数据库取出评论数据并转为DataFrame以方便分析
3.定义进度回调函数(用WebSocket推送进度)
4.调用`AgentDispatcher`调度诸个Agent执行
5.执行结束后更新任务结果及错误信息。以下是该函数的Python实现片段。
# agent.py _do_agent_run 核心代码
async def _do_agent_run(task_id: str, request: AgentRunRequest, db: AsyncSession):
from app.core.database import AsyncSessionLocal
async with AsyncSessionLocal() as session:
try:
# 1. 更新任务状态为running
await session.execute(
update(AgentTaskModel).where(AgentTaskModel.id == task_id).values(status="running")
)
await session.commit()
# 2. 读取评论数据并转为DataFrame
result = await session.execute(
select(Review).where(Review.dataset_id == request.dataset_id, Review.is_noise == 0)
)
reviews = result.scalars().all()
df = pd.DataFrame([{
"id": r.id,
"content": r.content,
"clean_content": r.clean_content or r.content,
"rating": r.rating,
"review_time": r.review_time,
} for r in reviews])
# 3. 进度回调:通过WebSocket推送Agent执行进度
async def progress_callback(agent_name: str, status: str, progress: int):
if task_id in active_connections:
try:
await active_connections[task_id].send_json({
"agent": agent_name,
"status": status,
"progress": progress,
})
except Exception:
pass
# 4. 调度多Agent执行
result_data = await dispatcher.dispatch(
request.dataset_id,
df,
request.user_input,
progress_callback=progress_callback,
)
# 5. 更新任务结果,状态改为done
await session.execute(
update(AgentTaskModel).where(AgentTaskModel.id == task_id).values(
status="done",
result_json=json.dumps(result_data, ensure_ascii=False, default=str),
)
)
await session.commit()
# 6. 通过WebSocket通知客户端完成
if task_id in active_connections:
await active_connections[task_id].send_json({"status": "done", "result": result_data})
except Exception as e:
# 异常处理:更新错误信息,状态改为failed
logger.error(f"Agent工作流失败: {e}", exc_info=True)
await session.execute(
update(AgentTaskModel).where(AgentTaskModel.id == task_id).values(
status="failed", error_message=str(e)
)
)
await session.commit()
设计思考:
- 单独创建数据库会话(AsyncSessionLocal),因而能自然地避开原请求中db会话被释放的问题
- 将评论数据转为DataFrame,有利于后续Agent做系统、规范的数据分析(如统计好评关键词、差评维度)
- 进度回调函数把Agent执行与WebSocket推送解耦,故WebSocket连接中断时不会影响Agent执行
- 异常捕获后及时、妥帖地更新任务状态及错误信息,对用户问题排查极有帮助
(3)WebSocket 实时状态推送
由于要让用户能及时、可靠地看到Agent的执行进度,故本节自然、妥帖地介绍了用WebSocket实时推送Agent状态的方法,即实现WebSocket接口`/status/{task_id}`,并用`active_connections`字典来管理当前所有连接:
# agent.py WebSocket 核心代码
active_connections: dict[str, WebSocket] = {}
@router.websocket("/status/{task_id}")
async def agent_status_ws(task_id: str, websocket: WebSocket):
"""WebSocket 实时推送 Agent 执行状态"""
await websocket.accept()
active_connections[task_id] = websocket
try:
while True:
data = await websocket.receive_text()
if data == "ping":
await websocket.send_text("pong") # 心跳检测
except WebSocketDisconnect:
pass
finally:
active_connections.pop(task_id, None) # 连接断开时清理
设计思考:
- 心跳检测(ping/pong)用来防止连接被网关/浏览器断开
- 连接断开时主动清理`active_connections`以避免内存泄漏
- Agent执行过程中通过`progress_callback` 函数实时、友好的推送进度信息,例如“评论分析Agent:已完成50%”。
(4)营销文案生成:LLM 的落地应用
除了多 Agent 工作流,我还单独实现了/marketing-copy接口,基于好评数据生成营销文案,这是 LLM 的直接落地场景:
# agent.py 营销文案生成核心代码
@router.post("/marketing-copy", summary="营销文案生成")
async def generate_marketing_copy(
request: MarketingCopyRequest,
db: AsyncSession = Depends(get_db),
):
"""根据好评数据生成三种风格营销文案"""
# 1. 读取好评数据
result = await db.execute(
select(Review).where(
Review.dataset_id == request.dataset_id,
Review.is_noise == 0,
)
)
reviews = result.scalars().all()
if not reviews:
raise HTTPException(404, "无评论数据")
# 2. 提取好评高频词
good_reviews = [r for r in reviews if r.rating is not None and cast(float, r.rating) >= 4]
word_counter = Counter()
for r in good_reviews[:500]:
for word in jieba.cut(r.content):
if len(word) >= 2:
word_counter[word] += 1
top_words = [w for w, _ in word_counter.most_common(30)]
# 3. 构造LLM Prompt
prompt = f"""你是专业电商营销文案师。根据以下用户好评高频词,生成三种风格的标题文案。
好评高频词(Top30):{', '.join(top_words)}
请生成:
1. 情感种草风文案(3条)
2. 功能卖点风文案(3条)
3. 数据背书风文案(3条)
4. 卖点提炼清单(3-5条)
5. 差评应对话术(关切型/解释型/补偿型各1条)
请以JSON格式返回:
{{
"emotional_copy": ["文案1", "文案2", "文案3"],
"feature_copy": ["文案1", "文案2", "文案3"],
"data_copy": ["文案1", "文案2", "文案3"],
"key_points": ["卖点1", "卖点2", "卖点3"],
"response_templates": {{
"caring": "关切型话术",
"explanatory": "解释型话术",
"compensatory": "补偿型话术"
}}
}}"""
# 4. 调用LLM生成文案
try:
response = await llm_client.achat([
{"role": "system", "content": "你是专业电商营销文案师,请以JSON格式返回。"},
{"role": "user", "content": prompt}
], temperature=0.7, max_tokens=1500)
# 提取JSON结果(防止LLM返回多余文本)
import re
json_match = re.search(r'\{.*\}', response, re.DOTALL)
if json_match:
return json.loads(json_match.group())
except Exception as e:
logger.error(f"营销文案生成失败: {e}")
# 兜底返回默认文案
return {"emotional_copy": ["品质生活,从这里开始", ...]}
三、框架设计的核心思考
1. 扩展性
- AgentDispatcher 调度器:所有 Agent 都通过
dispatcher.dispatch执行,后续新增 Agent(如竞品分析 Agent)只需在 Dispatcher 里注册,无需修改核心接口; - 数据模型分层:任务、报告、预警分开存储,后续加新功能(如自动生成 PPT)只需新增模型,不影响现有表结构。
2. 可靠性
- 异步执行 + 状态持久化:即使服务重启,任务状态也能从数据库恢复;
- 异常捕获:每个关键步骤都有异常处理,避免单个 Agent 失败导致整个工作流崩溃;
- 兜底逻辑:LLM 调用失败时返回默认文案,WebSocket 连接断开不影响 Agent 执行。
3. 可观测性
- WebSocket 实时推送:用户能看到执行进度;
- 详细的错误日志:记录 Agent 失败的具体原因;
- Token 用量统计:掌握 LLM 使用成本。
四、后续优化方向
1. Agent解耦:改为基于消息队列(RabbitMQ)的分布式调度,支持多实例部署
2.缓存层:加入缓存层Redis缓存高频查询的任务状态及好评高频词
3.可视化:增加任务结果的可视化报表(差评率趋势图)
4.权限控制:对dataset_id做权限控制,防止越权查询
5.LLM 微调:用自有评论数据对LLM做微调以提高文案生成、评论分析的准确性。
五、总结
这套基于 LLM 的多 Agent 工作流系统,核心是“把复杂业务拆成小Agent,用异步+状态持久化保证可靠性,用WebSocket提升体验”。并且我觉得很重要的值得考虑的一点是:技术落地必须兼顾“当下能用”和“未来可扩展”,因此单一职责、异步化、可观测性是核心原则。具体实现上,SQLAlchemy的ORM设计让数据层十分清晰,FastAPI的异步特性与Agent耗时执行的场景高度契合,WebSocket直接、巧妙地解决实时反馈问题。最后关于LLM的落地,我认为:真正的落地不在“调用API”,而在“数据预处理+Prompt工程+异常兜底”,只有把诸种细节处理妥帖,方能从“玩具”真正变为“可用的产品”。
更多推荐



所有评论(0)