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 })
解决方案
核心修正点
- 修正
undefined判断逻辑:原代码判断chunk.answer != 'undefined'会误判字符串'undefined'为有效值,应改为严格判断chunk.answer !== undefined - 简化流转换流程:无需先转成
ReadableStream再转回IterableReadableStream,直接基于原始流生成目标流 - 正确返回流式响应:使用
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
相关产品推荐
相关产品推荐

