Upstash context.call返回ReadableStream时前端仅收到workflowRunId的问题排查求助
Upstash context.call返回ReadableStream时前端仅收到workflowRunId的问题排查求助
看起来你遇到的核心问题是:后端通过Upstash Workflow的context.call能正确拿到LangGraph的流式自动补全数据,但前端调用后端接口时,只能收到workflowRunId,无法获取实际的自动补全内容。我来帮你拆解可能的原因和解决方案:
核心原因分析
Upstash Workflow的serve函数默认行为会将处理函数的返回值包装为包含workflow元数据(比如workflowRunId)的响应,而不是直接透传你返回的ReadableStream。另外,LangGraph返回的是JSON Lines(ndjson)格式的流式数据,需要前后端都正确处理这种格式才能拿到有效内容。
解决方案分步实施
1. 后端修改:手动转发流式数据给前端
你需要在后端消费context.call返回的流,然后创建一个新的流直接转发给前端,同时设置正确的响应头来标识流式内容。修改你的后端代码如下:
import { serve } from '@upstash/workflow/nextjs'; export const { POST } = serve<{ noteText: string }>( async (context) => { const { noteText } = context.requestPayload; const { status, headers, body } = await context.call<ReadableStream>('langgraph-request', { url: `${process.env.LANGGRAPH_RUN_URL}/runs/stream`, method: 'POST', headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${process.env.LANGSMITH_API_KEY}`, }, body: { assistant_id: 'autocomplete_agent', input: { note_text: noteText }, stream_mode: ['updates'], }, }); // 消费LangGraph的流式响应,转发给前端 const reader = body.getReader(); const forwardStream = new ReadableStream({ async start(controller) { try { while (true) { const { done, value } = await reader.read(); if (done) break; // 将每个数据块转发给前端 controller.enqueue(value); } controller.close(); } catch (err) { controller.error(err); } finally { reader.releaseLock(); } }, }); // 返回包含正确响应头的流式结果 return { headers: { // 保留LangGraph的响应头,覆盖或添加流式相关头 ...headers, 'Content-Type': 'application/x-ndjson', // 明确标识JSON Lines格式 'Transfer-Encoding': 'chunked', 'Cache-Control': 'no-cache', 'Content-Length': undefined, // 移除Content-Length,因为是流式 }, body: forwardStream, }; }, { failureFunction: ({ failStatus, failResponse, failHeaders }) => { console.error('Autocomplete workflow failed:', { status: failStatus, response: failResponse, headers: failHeaders, }); }, } );
2. 前端修改:正确解析JSON Lines流式数据
LangGraph返回的是每行一个JSON对象的流式数据(JSON Lines),你需要在前端分割每行并解析,提取自动补全内容。修改你的前端代码:
const getSuggestion = debounce(async (noteText: string, cb: (suggestion: string) => void) => { try { const response = await fetch('/api/workflow/autocomplete/', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ noteText }), }); if (!response.ok || !response.body) { throw new Error(`HTTP error! status: ${response.status}`); } const reader = response.body.getReader(); const decoder = new TextDecoder(); let partialLine = ''; // 处理跨块的不完整行 let suggestion = ''; while (true) { const { done, value } = await reader.read(); if (done) { // 处理最后一行未完成的内容 if (partialLine.trim()) processLine(partialLine); break; } const chunk = decoder.decode(value, { stream: true }); const lines = chunk.split('\n'); // 第一行可能和上一个块的末尾拼接 lines[0] = partialLine + lines[0]; partialLine = lines.pop() || ''; for (const line of lines) { if (line.trim()) processLine(line); } } cb(suggestion); // 最终更新 // 解析单个JSON行,提取自动补全内容 function processLine(line: string) { try { const data = JSON.parse(line); if (data.event === 'updates' && data.data?.autocompleteNote?.suggested_completion) { suggestion = data.data.autocompleteNote.suggested_completion; cb(suggestion); // 实时更新UI } } catch (e) { console.warn('Failed to parse streaming chunk:', e); } } } catch (err) { console.error('Autocomplete failed:', err); cb(''); } }, 300);
3. 额外验证点
- 确认你的Next.js版本支持流式响应(Next.js 13+的Pages Router或App Router都支持)
- 检查Upstash Workflow的权限配置,确保它能正确转发流式数据
- 在浏览器开发者工具的"网络"标签中,查看后端接口的响应,确认是否有分块的JSON Lines数据返回
内容来源于stack exchange
相关产品推荐
相关产品推荐

