从代码到落地:基于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工程+异常兜底”,只有把诸种细节处理妥帖,方能从“玩具”真正变为“可用的产品”。

更多推荐