ASP.NET Core Web API流式响应到Angular前端,如何实现逐段返回内容而非一次性获取全部
ASP.NET Core Web API流式响应到Angular前端,如何实现逐段返回内容而非一次性获取全部
咱们先搞清楚问题根源:你用yield return返回IEnumerable<string>时,ASP.NET Core默认的JSON格式化器会把所有yield的结果先收集到内存里,等整个序列生成完再一次性序列化成JSON数组返回给前端,所以Angular拿到的是完整的数组字符串,而不是逐段的流。要实现真正的流式传输,得从API和前端两方面调整,让数据生成一段就发送一段。
下面给你两种可行的解决方案:
方案一:直接写入响应流(轻量实现)
这种方式不需要额外依赖,直接通过ASP.NET Core的HttpResponse逐段写数据,前端监听进度事件处理每一段内容。
1. 修改ASP.NET Core API代码
把原来返回IEnumerable<string>的Action改成直接操作响应流,禁用缓存并确保每段数据写完就立即发送给客户端:
[HttpPost("getCopilotResponseV3")] public async Task GetCopilotResponsev3([FromBody] BLCopilotRequestV2 req, HttpResponse response) { // 设置响应头:禁用缓存、指定纯文本类型 response.Headers.ContentType = new MediaTypeHeaderValue("text/plain"); response.Headers.CacheControl = new CacheControlHeaderValue { NoCache = true, NoStore = true, MustRevalidate = true }; try { var streamResponse = _bLOpenAIService.GetOpenAIResponseStream(req); foreach (var update in streamResponse) { foreach (var cnt in update.ContentUpdate) { // 写入当前文本段 await response.WriteAsync(cnt.Text); // 强制把缓冲区内容推送给客户端,不等待后续数据 await response.Body.FlushAsync(); // 可选:给前端一点处理时间,避免发送过快导致UI卡顿 await Task.Delay(10); } } } catch (Exception ex) { _logger.LogError(ex, "Error streaming OpenAI response."); // 给前端发送错误提示 await response.WriteAsync($"[ERROR]: {ex.Message}"); } }
2. 修改Angular代码
开启进度报告,监听DownloadProgress事件获取逐段返回的文本:
Angular服务调整
requestFromBLCopilotAIV3(request: BLCopilotRequestV2) { return this.http.post( `${environment.getEndpoint()}Common/getCopilotResponseV3`, request, { responseType: 'text', observe: 'events', reportProgress: true // 必须开启,才能收到进度事件 } ); }
组件订阅逻辑调整
let accumulatedText = ''; // 记录已接收的所有文本,避免重复处理 this.copilotService.requestFromBLCopilotAIV3({ messageList: cpMessageList, indexName: this.copilotComponent.defaultValue, semanticConfigName: this.copilotComponent.mapName } as BLCopilotRequestV2) .pipe(takeUntil(this.stopSignal)) .subscribe({ next: (event) => { if (event.type === HttpEventType.DownloadProgress && event.partialText) { // 提取本次新增的文本段 const newText = event.partialText.substring(accumulatedText.length); accumulatedText = event.partialText; console.log('收到新内容:', newText); // 这里更新UI,比如把新文本追加到聊天框 this.chatResponse += newText; this.changeDetector.detectChanges(); // 强制更新UI } else if (event.type === HttpEventType.Response) { // 处理响应末尾的剩余文本 const remainingText = event.body?.substring(accumulatedText.length) || ''; if (remainingText) { this.chatResponse += remainingText; } console.log('流传输完成'); } }, error: (err) => { console.error('流传输错误:', err); this.isLoading = false; this.responseStarted = false; this.changeDetector.detectChanges(); }, complete: () => { this.isLoading = false; this.changeDetector.detectChanges(); } });
方案二:使用Server-Sent Events(SSE,标准流式协议)
SSE是专门为服务器向客户端流式发送文本数据设计的协议,兼容性好,且自带消息分割、重连机制,适合聊天类场景。
1. 修改ASP.NET Core API代码
按照SSE格式发送数据(每条消息以data: 开头,\n\n结尾):
[HttpPost("getCopilotResponseV3")] public async Task GetCopilotResponsev3([FromBody] BLCopilotRequestV2 req, HttpResponse response) { // 设置SSE专属响应头 response.Headers.ContentType = new MediaTypeHeaderValue("text/event-stream"); response.Headers.CacheControl = new CacheControlHeaderValue { NoCache = true, NoStore = true }; response.Headers.Connection = "keep-alive"; try { var streamResponse = _bLOpenAIService.GetOpenAIResponseStream(req); foreach (var update in streamResponse) { foreach (var cnt in update.ContentUpdate) { // 按照SSE格式写入文本段 await response.WriteAsync($"data: {cnt.Text}\n\n"); await response.Body.FlushAsync(); await Task.Delay(10); } } // 发送结束标记 await response.WriteAsync("data: [DONE]\n\n"); await response.Body.FlushAsync(); } catch (Exception ex) { _logger.LogError(ex, "Error streaming SSE response."); await response.WriteAsync($"data: [ERROR] {ex.Message}\n\n"); await response.Body.FlushAsync(); } }
2. 修改Angular代码
用fetch API结合ReadableStream处理SSE消息,比原生EventSource更灵活(支持POST请求):
Angular服务调整
requestFromBLCopilotAIV3(request: BLCopilotRequestV2): Observable<string> { return new Observable(observer => { const fetchOptions: RequestInit = { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(request), signal: this.stopSignal // 绑定取消信号,避免内存泄漏 }; fetch(`${environment.getEndpoint()}Common/getCopilotResponseV3`, fetchOptions) .then(response => { if (!response.ok) throw new Error(`请求失败:${response.status}`); const reader = response.body?.getReader(); if (!reader) { observer.error('无响应内容可读取'); observer.complete(); return; } const decoder = new TextDecoder(); let buffer = ''; // 递归读取每一段数据 const readChunk = () => { reader.read().then(({ done, value }) => { if (done) { observer.complete(); return; } // 解码二进制数据为文本,拼接缓冲区 buffer += decoder.decode(value, { stream: true }); // 按SSE消息分隔符分割内容 const messages = buffer.split('\n\n'); buffer = messages.pop() || ''; // 保留未完成的消息片段 // 处理每条SSE消息 messages.forEach(msg => { if (msg.startsWith('data: ')) { const text = msg.substring(6).trim(); if (text === '[DONE]') { observer.complete(); } else if (text.startsWith('[ERROR]')) { observer.error(text.substring(8)); } else { observer.next(text); } } }); readChunk(); }).catch(err => { observer.error(err); observer.complete(); }); }; readChunk(); }).catch(err => { observer.error(err); observer.complete(); }); }); }
组件订阅逻辑调整
this.copilotService.requestFromBLCopilotAIV3({ messageList: cpMessageList, indexName: this.copilotComponent.defaultValue, semanticConfigName: this.copilotComponent.mapName } as BLCopilotRequestV2) .pipe(takeUntil(this.stopSignal)) .subscribe({ next: (chunk) => { console.log('收到流式内容:', chunk); this.chatResponse += chunk; this.changeDetector.detectChanges(); }, error: (err) => { console.error('流式请求出错:', err); this.isLoading = false; this.responseStarted = false; this.changeDetector.detectChanges(); }, complete: () => { this.isLoading = false; this.changeDetector.detectChanges(); } });
额外注意事项
- CORS配置:如果前端和API跨域,确保CORS策略允许
text/plain或text/event-stream类型的响应。 - 响应缓冲:如果API启用了
UseResponseBuffering,可能需要针对这个Action禁用缓冲,避免数据被缓存后一次性发送。 - 资源释放:前端取消请求时(比如用户关闭聊天框),要确保通过
takeUntil或AbortSignal终止请求,避免内存泄漏。
备注:内容来源于stack exchange,提问作者Gayantha
相关产品推荐
相关产品推荐

