You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Google ADK流式响应时出现重复消息的技术求助

问题:Google ADK Sequential Workflow流式响应重复输出

环境配置

  • 使用Google ADK搭配Sequential Workflow编排器
  • 配置两个智能体:
    • sql_agent:生成SQL查询语句
    • execute_sql_agent:执行SQL查询
  • 通过/run_sse端点将响应流式传输到自定义Chainlit UI

问题现象

收到重复的流式消息,每个智能体的结果被多次输出,示例:

sql_agent: final SQL result...
execute_sql_agent: final SQL result...
sql_agent: final SQL result...
execute_sql_agent: final SQL result...

相关Chainlit代码

@cl.on_message
async def on_message(message: cl.Message):
    user_id = cl.user_session.get("user_id")
    session_id = cl.user_session.get("session_id")

    if not session_id:
        await cl.Message(
            content="❌ No active session found. Please restart the chat."
        ).send()
        return

    payload = {
        "app_name": APP_NAME,
        "user_id": user_id,
        "session_id": session_id,
        "streaming": True,
        "new_message": {
            "role": "user",
            "parts": [{"text": message.content}],
        },
    }

    msg = cl.Message(content="")
    await msg.send()

    try:
        async with aiohttp.ClientSession() as session:
            async with session.post(
                f"{API_BASE_URL}/run_sse",
                headers={"Content-Type": "application/json"},
                json=payload,
            ) as response:

                if response.status != 200:
                    text = await response.text()
                    await cl.Message(content=f"❌ API Error: {text}").send()
                    return

                async for line in response.content:
                    decoded = line.decode("utf-8").strip()

                    print("######## LOOP START ########")
                    print("SSE:", decoded)
                    print("######## LOOP END ##########")

                    # Skip empty lines or keepalive events
                    if not decoded or not decoded.startswith("data:"):
                        continue

                    data = decoded.replace("data:", "").strip()

                    if data == "[DONE]":
                        break

                    try:
                        event = json.loads(data)
                        content = event.get("content", {})
                        parts = content.get("parts", [])

                        if parts and "text" in parts[0]:
                            token = parts[0]["text"]
                            await msg.stream_token(token)

                    except json.JSONDecodeError:
                        continue

    except Exception as e:
        await cl.Message(
            content=f"⚠️ Connection error: {e}"
        ).send()

排查与解决思路

1. 确认后端SSE输出是否重复

先查看代码中的print日志,如果日志里确实多次出现相同的sql_agent和execute_sql_agent输出,说明问题出在Google ADK的工作流配置上:

  • 检查Sequential Workflow是否被重复触发,比如session复用导致工作流多次执行
  • 验证智能体的输出推送配置,是否被设置为重复发送
  • 排查工作流的循环逻辑,是否存在意外的重复执行分支

2. 前端添加去重逻辑

如果后端暂时无法调整,可在Chainlit代码中加入去重机制,记录已处理的智能体输出:

@cl.on_message
async def on_message(message: cl.Message):
    user_id = cl.user_session.get("user_id")
    session_id = cl.user_session.get("session_id")

    if not session_id:
        await cl.Message(
            content="❌ No active session found. Please restart the chat."
        ).send()
        return

    payload = {
        "app_name": APP_NAME,
        "user_id": user_id,
        "session_id": session_id,
        "streaming": True,
        "new_message": {
            "role": "user",
            "parts": [{"text": message.content}],
        },
    }

    msg = cl.Message(content="")
    await msg.send()
    # 新增:记录已处理的输出标识
    processed_outputs = set()

    try:
        async with aiohttp.ClientSession() as session:
            async with session.post(
                f"{API_BASE_URL}/run_sse",
                headers={"Content-Type": "application/json"},
                json=payload,
            ) as response:

                if response.status != 200:
                    text = await response.text()
                    await cl.Message(content=f"❌ API Error: {text}").send()
                    return

                async for line in response.content:
                    decoded = line.decode("utf-8").strip()

                    print("######## LOOP START ########")
                    print("SSE:", decoded)
                    print("######## LOOP END ##########")

                    if not decoded or not decoded.startswith("data:"):
                        continue

                    data = decoded.replace("data:", "").strip()

                    if data == "[DONE]":
                        break

                    try:
                        event = json.loads(data)
                        content = event.get("content", {})
                        parts = content.get("parts", [])

                        if parts and "text" in parts[0]:
                            token = parts[0]["text"]
                            # 根据内容前缀标识去重
                            if "sql_agent:" in token:
                                key = "sql_agent"
                            elif "execute_sql_agent:" in token:
                                key = "execute_sql_agent"
                            else:
                                key = token

                            if key not in processed_outputs:
                                await msg.stream_token(token)
                                processed_outputs.add(key)

                    except json.JSONDecodeError:
                        continue

    except Exception as e:
        await cl.Message(
            content=f"⚠️ Connection error: {e}"
        ).send()

3. 检查Session复用问题

确认session_id是否被重复使用,导致同一工作流实例被多次触发。可以尝试在每次请求时生成新的session_id,或者检查后端的session管理逻辑,确保每个用户请求对应唯一的工作流实例。

内容的提问来源于stack exchange,提问作者Meenakshi Kharwade

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 00:05:16