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

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();
    }
});

额外注意事项

  1. CORS配置:如果前端和API跨域,确保CORS策略允许text/plain或text/event-stream类型的响应。
  2. 响应缓冲:如果API启用了UseResponseBuffering,可能需要针对这个Action禁用缓冲,避免数据被缓存后一次性发送。
  3. 资源释放:前端取消请求时(比如用户关闭聊天框),要确保通过takeUntil或AbortSignal终止请求,避免内存泄漏。

备注:内容来源于stack exchange,提问作者Gayantha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 12:00:27