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

Chainlit中异步生成器实现流式响应失败问题排查与解决

Chainlit异步流式输出失效问题排查与解决

问题说明

编写的Chainlit异步RAG应用无法实现逐token流式输出,最终会一次性收到完整答案,但同步版本逻辑可正常流式输出。已知:

  • rag.astream(message.content) 是异步Python生成器
  • adapter.chat 为支持流式的 langchain_community.chat_models.BedrockChat

异步代码示例

import chainlit as cl

from langchain.prompts import PromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough, RunnableParallel
from langchain_community.document_transformers import LongContextReorder

from lib.aws import AWSAdapter


def format_docs(docs):
    """Format the retrieved documents."""
    return "\n\n".join(doc.page_content for doc in docs)


@cl.on_chat_start
async def on_chat_start():
    adapter = AWSAdapter()

    vectordb = adapter.opensearch_client
    retriever = vectordb.as_retriever()

    prompt = PromptTemplate(
        template= \
            """Answer the question based only on the following context:
            {context}
            
            Question: {question}
            """
            ,
        input_variables=["context", "question"],
    )

    reordering = LongContextReorder()

    chain = (prompt |
             adapter.chat |
             StrOutputParser())

    rag = RunnableParallel(
        {"context": retriever | reordering.transform_documents | format_docs,
         "question": RunnablePassthrough()}
    ).assign(answer=chain)

    cl.user_session.set("rag", rag)


@cl.on_message
async def main(message: cl.Message):
    rag = cl.user_session.get("rag")

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

    async for chunk in rag.astream(message.content):
        for key in chunk:
            if key == "answer":
                await msg.stream_token(chunk[key])

    await msg.send()

同步工作代码示例

output = {}
curr_key = None
for chunk in rag.stream("Question ..."):
    for key in chunk:
        if key not in output:
            output[key] = chunk[key]
        else:
            output[key] += chunk[key]
            
        if key == 'answer':
            print(chunk[key], end="", flush=True)
            
        curr_key = key

解决方案

1. 移除冗余的消息发送调用

原代码末尾的await msg.send()会将已通过stream_token逐步更新的消息再次完整发送,覆盖流式效果。修改后的消息处理逻辑:

@cl.on_message
async def main(message: cl.Message):
    rag = cl.user_session.get("rag")

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

    async for chunk in rag.astream(message.content):
        for key in chunk:
            if key == "answer":
                await msg.stream_token(chunk[key])
    # 移除此处多余的await msg.send()调用

2. 确保BedrockChat显式启用流式模式

检查AWSAdapter中BedrockChat的初始化代码,必须设置streaming=True,异步模式依赖该参数触发流式输出:

# AWSAdapter类中chat初始化示例
from langchain_community.chat_models import BedrockChat

class AWSAdapter:
    def __init__(self):
        self.chat = BedrockChat(
            model_id="your-target-model-id",  # 例如 anthropic.claude-v2
            streaming=True,  # 必须开启该参数
            # 其他配置:credentials、region_name等
        )
        # 其他初始化逻辑

3. 优化同步逻辑的异步兼容性(可选)

如果retriever和LongContextReorder.transform_documents是同步方法,会阻塞异步流程,需将其包装为异步逻辑:

# 包装同步函数为异步
async def async_format_docs(docs):
    return format_docs(docs)

async def async_reorder_docs(docs):
    reordering = LongContextReorder()
    return reordering.transform_documents(docs)

# 重构rag链使用异步包装后的逻辑
rag = RunnableParallel(
    {"context": retriever | RunnableLambda(async_reorder_docs) | RunnableLambda(async_format_docs),
     "question": RunnablePassthrough()}
).assign(answer=chain)

完成上述修改后,重启Chainlit应用即可实现逐token流式输出。

内容的提问来源于stack exchange,提问作者Vasil Yordanov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 02:26:08