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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 20:14:48