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

Node项目中如何从langchain.js导入IterableReadableStream并处理流

问题:转换LangChain流并返回仅含answer字段的流式响应

需要在Node项目中从langchain.js导入IterableReadableStream类,将类型为IterableReadableStream<{answer:string,content:Document[]}>的流转换为IterableReadableStream<string>类型的可迭代流,仅提取其中的answer字段,最终将提取出的answer字段流发送至前端,而非返回包含content字段的原始流。

原代码实现如下:

const stream = await Package.retrievalChain.stream({
    input: textFeed,

});

const answerStream: ReadableStream<String> = new ReadableStream({
    async pull(controller) {
        try {
            for await (const chunk of stream) {
                if (chunk.answer != 'undefined')
                    controller.enqueue(chunk.answer);
            }
            controller.close();
        } catch (error: Error | any) {
            console.log(error.message)
        }

    },
});

console.log("---Readable Stream generated")
const iterableReadableStream = convertEventStreamToIterableReadableDataStream(answerStream);
try {
    for await (const chunk of iterableReadableStream) {
        process.stdout.write(`${chunk}`);
    }

} catch (error: Error | any) {
    console.log(error.message)
} finally {
}


// return new StreamingTextResponse(iterableReadableStream);
return new NextResponse('done', { status: 200 })

解决方案

核心修正点

  1. 修正undefined判断逻辑:原代码判断chunk.answer != 'undefined'会误判字符串'undefined'为有效值,应改为严格判断chunk.answer !== undefined
  2. 简化流转换流程:无需先转成ReadableStream再转回IterableReadableStream,直接基于原始流生成目标流
  3. 正确返回流式响应:使用StreamingTextResponse替代NextResponse('done'),确保流式内容发送到前端

修正后的完整代码

import { IterableReadableStream } from "@langchain/core/utils/stream";
// 确保已正确导入Package和Document类型

export async function handleStreamingResponse(textFeed) {
    // 获取LangChain返回的原始流
    const originalStream = await Package.retrievalChain.stream({
        input: textFeed,
    });

    // 创建仅提取answer字段的IterableReadableStream
    const answerStream = new IterableReadableStream(async function* () {
        try {
            for await (const chunk of originalStream) {
                // 过滤无效的answer内容
                if (chunk.answer !== undefined && chunk.answer.trim()) {
                    yield chunk.answer;
                }
            }
        } catch (error) {
            console.error('流处理异常:', error.message);
            throw error; // 抛出异常让前端捕获处理
        }
    });

    // 返回流式响应,前端可直接接收分块内容
    return new StreamingTextResponse(answerStream);
}

前端接收示例(可选)

前端可以通过fetchAPI接收流式响应:

async function fetchStream() {
    const response = await fetch('/your-api-endpoint', {
        method: 'POST',
        body: JSON.stringify({ textFeed: '你的输入内容' }),
        headers: { 'Content-Type': 'application/json' }
    });

    const reader = response.body.getReader();
    const decoder = new TextDecoder();

    while (true) {
        const { done, value } = await reader.read();
        if (done) break;
        // 处理每一块answer内容
        console.log(decoder.decode(value));
        // 比如插入到页面中
        document.getElementById('answer-container').textContent += decoder.decode(value);
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:37:02