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
相关产品推荐
相关产品推荐

