在电商精细化运营的大背景下,评论数据已经成为企业洞察用户需求、优化产品策略、提升营销转化的核心资产。但传统的人工分析效率低、维度单一,而单纯的 LLM 调用又难以适配复杂的业务分析场景。基于此,我主导设计并实现了一套多 Agent 协同的电商评论数据分析系统,通过 Agent 分工协作、异步工程化、实时通信、流式交互等技术,将评论分析从 “单点查询” 升级为 “全链路智能洞察”。本文将从技术架构、核心实现、工程思考三个维度,拆解这套系统的设计与落地过程。

一、系统核心定位与技术挑战

1. 业务核心诉求

电商评论分析的核心痛点在于:数据量大(单数据集可达 10 万 + 条)、分析维度多(情感倾向、主题挖掘、竞品对比、营销提炼等)、交互要求高(实时反馈、流式响应、进度可视化)。单纯的 “LLM + 数据库” 模式无法满足:

  • 多维度分析需要专业分工(如情感分析 Agent、营销文案 Agent);
  • 大数量级数据处理需要异步化,避免接口阻塞;
  • 用户需要实时感知分析进度,而非被动等待;
  • 不同场景(闲聊、精准分析、摘要生成)需要差异化的交互方式。

2. 核心技术挑战

  • 如何设计 Agent 分工与调度机制,避免 “大而全” 的单 Agent 性能瓶颈;
  • 异步任务与实时通信的协同(WebSocket + 异步 DB 操作);
  • 流式响应(SSE)与 LLM 输出的高效结合;
  • 异常处理与工程化保障(任务状态持久化、错误重试、日志监控);
  • 数据处理性能(10 万 + 条评论的 DataFrame 构建与分析)。

二、核心技术架构与实现细节

整个系统的 API 层基于 FastAPI 构建,核心分为任务管理、多 Agent 调度、实时通信、流式交互、数据持久化五大模块,以下是关键模块的设计与实现思考:

1. 异步任务管理:从 “同步阻塞” 到 “后台异步 + 状态追踪”

传统的接口调用如果直接处理 10 万 + 条评论的分析,会因超时、阻塞导致体验极差。因此我设计了 “任务 ID + 异步后台 + 状态持久化” 的模式:

核心实现
@router.post("/run", summary="启动多Agent工作流")
async def run_agents(
    request: AgentRunRequest,
    background_tasks: BackgroundTasks,
    db: AsyncSession = Depends(get_db),
):
    # 1. 任务初始化:生成唯一task_id,状态置为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. 后台任务解耦:核心逻辑丢到background_tasks,接口立即返回
    background_tasks.add_task(_do_agent_run, task.id, request, db)
    return {"task_id": task.id, "status": "pending", "message": "多Agent工作流已启动"}
设计思考
  • 任务状态持久化:通过AgentTaskModel将任务状态(pending/running/done/failed)、结果、错误信息全量存储,解决了 “任务中断后无法追溯” 的问题;
  • 异步会话隔离:核心分析逻辑_do_agent_run中重新创建AsyncSessionLocal,避免 FastAPI 的 Depends 注入的 Session 因请求结束被释放,保证长时任务的 DB 连接稳定性;
  • 进度回调机制:设计progress_callback函数,将每个 Agent 的执行进度(如 “情感分析 Agent:50%”)实时推送到 WebSocket,让用户感知分析过程:
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

2. 多 Agent 调度:分工协作,而非 “单 Agent 全包揽”

多 Agent 的核心价值是 “专业分工”,我设计了AgentDispatcher调度器,将不同分析维度拆解为独立子 Agent(情感分析、主题挖掘、营销文案、数据摘要等),核心逻辑如下:

核心实现
# 执行多Agent工作流
result_data = await dispatcher.dispatch(
    request.dataset_id,
    df,  # 评论数据DataFrame
    request.user_input,
    competitor_ids=request.competitor_ids,
    progress_callback=progress_callback,  # 进度回调
)
设计思考
  • Agent 职责边界:每个子 Agent 只负责单一维度(如MarketingAgent仅处理营销文案生成,SentimentAgent仅做情感分析),避免单 Agent 逻辑臃肿,同时便于后续扩展(如新增 “差评根因分析 Agent”);
  • 数据传递效率:将评论数据封装为 Pandas DataFrame,统一数据格式,避免各 Agent 重复解析 DB 数据;
  • 意图识别短路逻辑:在对话接口中,先通过SchedulingAgent识别用户意图 —— 如果是闲聊(SmallTalk),直接走闲聊 Agent,无需加载数据集,提升响应速度:
intent_result = await dispatcher.scheduling_agent._identify_intent(request.message)
intent_type = intent_result.get("intent_type", IntentType.FULL_ANALYSIS)
# 闲聊短路:无需数据集,直接流式输出
if intent_type == IntentType.SMALL_TALK:
    async for chunk in llm_client.astream(...):
        yield f"data: {json.dumps({'content': chunk})}\n\n"
    return

3. 实时通信与流式交互:从 “被动等待” 到 “实时反馈”

用户体验的核心是 “感知可控”,因此系统同时实现了WebSocket(任务进度) 和SSE(流式对话) 两种实时交互方式:

(1)WebSocket:任务进度实时推送
@router.websocket("/status/{task_id}")
async def agent_status_ws(task_id: str, websocket: WebSocket):
    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)

设计思考:通过active_connections字典维护 WebSocket 连接映射,任务执行过程中通过progress_callback推送进度,解决了 “用户不知道分析到哪一步” 的痛点;同时增加心跳检测(ping/pong),避免长连接被网关断开。

(2)SSE:LLM 响应流式输出

传统的 “一次性返回结果” 会让用户等待数秒,而 SSE(Server-Sent Events)可以将 LLM 的 token 输出实时推送给前端,实现 “边想边说” 的交互体验:

def _sse_stream(generator):
    """将异步生成器包装为 SSE StreamingResponse"""
    return StreamingResponse(generator, media_type="text/event-stream")

@router.post("/chat", summary="主Agent对话(SSE流式)")
async def agent_chat(request: AgentChatRequest, db: AsyncSession = Depends(get_db)):
    async def generate():
        # 流式输出LLM响应
        async for chunk in llm_client.astream(messages, temperature=0.5, max_tokens=800):
            yield f"data: {json.dumps({'content': chunk})}\n\n"
        yield "data: [DONE]\n\n"  # 结束标记
    return _sse_stream(generate())

设计思考

  • 异步生成器(async generator)保证了流式输出的效率,避免内存堆积;
  • 统一的_sse_stream包装函数,让所有流式接口复用逻辑,降低代码冗余;
  • 增加[DONE]结束标记,前端可精准感知响应完成,避免无限等待。

4. 营销文案生成:从 “纯 LLM 生成” 到 “数据驱动 + 规则约束”

营销文案生成是核心业务场景,单纯的 LLM 生成容易脱离评论数据,因此我设计了 “数据提取 + 规则约束 + LLM 生成” 的三层逻辑:

核心实现
# 1. 数据提取:筛选4分以上好评,提取高频词
good_reviews = [r for r in reviews if r.rating and 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)]

# 2. 规则约束:明确生成格式(情感种草/功能卖点/数据背书)
prompt = f"""根据以下用户好评高频词,生成三种风格的标题文案。
好评高频词(Top30):{', '.join(top_words)}
请生成:
1. 情感种草风文案(3条)
2. 功能卖点风文案(3条)
3. 数据背书风文案(3条)
...
请以JSON格式返回:{{...}}"""

# 3. LLM生成+容错:解析失败返回兜底文案
try:
    response = await llm_client.achat(...)
    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": ["品质生活,从这里开始"...]}  # 兜底

设计思考

  • 数据驱动:基于真实好评的高频词生成文案,避免 LLM “凭空创作”;
  • 格式约束:通过 prompt 明确输出 JSON 格式,降低前端解析成本;
  • 容错机制:LLM 响应解析失败时返回兜底文案,保证接口可用性。

三、工程化思考与避坑指南

在落地这套系统的过程中,我踩了很多坑,也形成了一些工程化的思考:

1. 异步 DB 操作的 “坑”:会话隔离与事务管理

FastAPI 的Depends(get_db)注入的 AsyncSession 是请求级别的,而后台任务的执行时间远长于请求生命周期 —— 如果直接复用该 Session,会导致 “Session closed” 错误。因此核心解决思路是:

  • 后台任务中重新创建AsyncSessionLocal,保证会话独立;
  • 关键操作(如更新任务状态)单独提交事务,避免一个操作失败导致全量回滚。

2. 性能优化:数据处理的 “取舍”

10 万 + 条评论构建 DataFrame 时,最初直接遍历所有字段导致内存溢出,优化策略:

  • 只加载分析所需字段(如content/rating/clean_content),丢弃冗余字段;
  • 好评分析时限制最多取 500 条,平衡分析效果与性能;
  • LLM 输入截断(如摘要上下文限制 2000 字),避免 token 超限。

3. 异常处理:“不放过错误,也不暴露错误”

  • 全链路日志:通过logger.error(f"...", exc_info=True)记录异常栈,便于定位问题;
  • 错误持久化:将任务错误信息存储到AgentTaskModelerror_message字段,支持用户查询失败原因;
  • 前端友好:接口返回错误时,将技术错误转为用户可理解的提示(如 “数据集不存在” 而非 “SQLAlchemy NoResultFound”)。

4. 扩展性设计:Agent 的 “插拔式” 架构

AgentDispatcher中维护了 Agent 的注册表,新增 Agent 只需:

  1. 实现execute方法的 Agent 类;
  2. 在 Dispatcher 中注册该 Agent;
  3. 扩展意图识别规则。这种设计让系统可以快速适配新的分析场景(如新增 “用户画像 Agent”)。

四、总结与未来规划

这套多 Agent 电商评论分析系统的核心价值,在于将 “LLM 的智能” 与 “工程化的稳定” 结合 —— 既通过 Agent 分工实现了专业的多维度分析,又通过异步、流式、实时通信保证了用户体验,同时通过工程化手段保障了系统的稳定性。

未来优化方向

  1. 竞品数据整合:目前competitor_data为 TODO,后续将接入竞品评论数据,实现跨数据集对比分析;
  2. Agent 能力增强:引入 RAG(检索增强生成),将评论数据嵌入向量库,提升分析的精准度;
  3. 性能监控:接入 Prometheus 监控任务执行时长、LLM token 用量、DB 查询耗时,针对性优化;
  4. 多模态扩展:支持评论中的图片 / 视频分析,结合 OCR 提取图文信息。

从技术角度来说,多 Agent 系统的核心不是 “堆功能”,而是 “分治 + 协同”—— 把复杂的业务拆分为小而专的 Agent,再通过调度和通信让它们协同工作,同时兼顾工程化的稳定性和用户体验。这也是我在这次开发中最深的体会:AI 系统的落地,一半是算法智能,一半是工程能力。

更多推荐