ReactiveX结合OpenAI异步流式接口使用to_iterable()导致进程挂起的问题咨询
ReactiveX结合OpenAI异步流式接口使用to_iterable()导致进程挂起的问题咨询
大家好,我最近在做ReactiveX和OpenAI异步流式接口结合的开发时遇到了一个奇怪的问题,想过来请教一下各位。
先说说正常工作的情况:当我用ReactiveX的subscribe(print)来处理OpenAI返回的流时,一切都很顺畅——所有print命令都能按预期执行,内容打印正常,进程也能顺利退出,代码如下:
stream = await client.create_completion(...) stream.subscribe(print) # works perfectly
只要是用响应式的方式处理这个流,就没什么问题。
但当我尝试用pipe(ops.to_iterable()).run()把流转换成可迭代对象时,进程就会无限挂起,完全停不下来,代码是这样的:
stream = await client.create_completion(...) stream.pipe(ops.to_iterable()).run() # hangs :(
我实在搞不懂为什么会这样,明明都是处理同一个流,换了个操作就出问题了。有没有大佬能帮我看看,我到底哪里操作错了?
提前感谢大家的帮助!
以下是我的完整代码:
from openai import AsyncOpenAI, AsyncStream from openai.types.chat import ChatCompletionChunk from openai.types.chat.chat_completion_chunk import Choice import reactivex as rx from app.ai.bridge.chat.chat_completion_types import ChatRequest, ToolConfig from app.ai.bridge.chat.drivers.openai_chat_driver_reactive_mappers import ( map_message_to_openai, map_toolconfig_to_openai, ) import asyncio class OpenaiClientReactive: def __init__(self, openai_client: AsyncOpenAI) -> None: self.api = openai_client async def create_completion( self, chat_request: ChatRequest, tool_config: ToolConfig | None = None ) -> rx.Observable[Choice]: stream: rx.subject.ReplaySubject[Choice] = rx.subject.ReplaySubject() async def do_stream() -> None: async_stream: AsyncStream[ChatCompletionChunk] = await self._stream_openai( chat_request, tool_config ) max_variants_expected = chat_request.options.num_variants num_indexes_completed = 0 async for chunk in async_stream: for choice in chunk.choices: if choice.finish_reason: num_indexes_completed += 1 stream.on_next(choice) if max_variants_expected == num_indexes_completed: # If all indexes are complete, we can complete the stream print("All indexes complete") break await async_stream.close() stream.on_completed() asyncio.create_task(do_stream()) return stream async def _stream_openai( self, chat_request: ChatRequest, tool_config: ToolConfig | None = None, ) -> AsyncStream[ChatCompletionChunk]: mapped_messages = mapped_messages = [ map_message_to_openai(message) for message in chat_request.context.messages ] if tool_config: return await self.api.chat.completions.create( # TODO: Move this to database driven configuration, since it's an LLM. model="gpt-3.5-turbo", messages=mapped_messages, stream=True, n=chat_request.options.num_variants, tools=map_toolconfig_to_openai(tool_config), ) else: return await self.api.chat.completions.create( # TODO: Move this to database driven configuration, since it's an LLM. model="gpt-3.5-turbo", messages=mapped_messages, stream=True, n=chat_request.options.num_variants, )
备注:内容来源于stack exchange,提问作者Monarch Wadia
相关产品推荐
相关产品推荐

