如何通过Flux与Angular Observables实现流式数据逐项消费
问题原因
你当前的写法存在两个核心问题:
- Angular HttpClient 默认会等待完整响应接收完毕后,再统一解析返回结果,不会逐块读取流式返回的内容,所以只会在所有数据返回后触发一次next回调
- 你设置的后端返回格式是
APPLICATION_STREAM_JSON_VALUE,每条数据是独立的分段内容,前端默认的JSON解析逻辑无法识别分段的JSON结构,必须手动逐块读取拆分
修复步骤
1. 后端适配(必做优化)
在返回的每条数据末尾增加换行符作为分隔符,避免流传输过程中数据块拆分异常:
// import reactor.core.publisher.Flux; import org.springframework.http.MediaType; @GetMapping(value = "/data/stream", produces = MediaType.APPLICATION_STREAM_JSON_VALUE) public Flux<String> streamDataArray() { String [] array = new String[20]; for(int i=0;i<array.length; i++){ array[i] = i+1+"" ; } return Flux.fromArray(array).map(s->{ try { Thread.sleep(250L); } catch (InterruptedException interruptedException) { interruptedException.printStackTrace(); } // 增加换行作为单条数据分隔符 return s + "\n"; }); }
2. 前端Angular代码修改
调整StreamingService的流请求配置
开启进度上报、逐块读取响应内容:
import { Injectable } from '@angular/core'; import { HttpClient, HttpEvent, HttpEventType } from '@angular/common/http'; import { Observable, from } from 'rxjs'; import { filter, map, mergeMap } from 'rxjs/operators'; @Injectable({ providedIn: 'root' }) export class StreamingService { constructor(private httpClient: HttpClient) { } getStream(): Observable<string> { return this.httpClient.get('/api/data/stream', { observe: 'events', responseType: 'text', // 必须开启进度上报,才能拿到逐块返回的响应片段 reportProgress: true }).pipe( // 过滤出有效响应事件 filter((event: HttpEvent<string>) => event.type === HttpEventType.DownloadProgress || event.type === HttpEventType.Response ), map(event => { if (event.type === HttpEventType.DownloadProgress) { // Angular 14+ 支持partialText直接获取已接收的文本片段 // @ts-ignore return event.partialText || ''; } if (event.type === HttpEventType.Response) { return event.body || ''; } return ''; }), // 按分隔符拆分出单条数据 mergeMap(fullText => { const items = fullText.split('\n').filter(item => item.trim()); return from(items); }) ); } }
提示:如果你使用的是Angular 14以下版本不支持
partialText,可以将responseType设置为'arraybuffer',自己通过TextDecoder逐块解码二进制内容。
组件订阅逻辑无需大幅调整
直接订阅即可拿到逐条返回的单条数据:
getStream() { this.streamingService.getStream() .subscribe({ next: (item: string) => { // 每250ms会触发一次,拿到单条数据 console.log('收到数据:', item); // 可直接push到展示数组实现实时渲染 }, error: (msg) => { console.error(JSON.stringify(msg)); }, complete: () => { console.log('Observable completed'); } }); }
可选优化方案:用SSE实现更稳定的流传输
如果你的场景就是服务端单向推送数据,更推荐使用标准SSE(Server-Sent Events)协议实现,浏览器原生支持,无需手动处理流拆分:
- 后端仅需修改返回的Content-Type:把
APPLICATION_STREAM_JSON_VALUE替换为MediaType.TEXT_EVENT_STREAM_VALUE - 前端直接用
EventSource对接即可,代码更简洁稳定:
getStreamBySse(): Observable<string> { return new Observable(observer => { const sse = new EventSource('/api/data/stream'); sse.onmessage = (event) => observer.next(event.data); sse.onerror = (err) => { observer.error(err); sse.close(); }; return () => sse.close(); }); }
内容的提问来源于stack exchange,提问作者Mehdi
相关产品推荐
相关产品推荐

