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

如何通过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)协议实现,浏览器原生支持,无需手动处理流拆分:

  1. 后端仅需修改返回的Content-Type:把APPLICATION_STREAM_JSON_VALUE替换为MediaType.TEXT_EVENT_STREAM_VALUE
  2. 前端直接用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 08:15:03