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

Angular无法订阅Spring WebFlux流的问题求助

Spring Boot WebFlux SSE服务Angular客户端订阅异常解决方案

问题描述

  • 基于Spring Boot WebFlux实现SSE(服务器推送事件)数据流服务,核心代码如下:
    1. 数据流生成方法:
      public Flux<Data> getStream() {
          Flux<Data> initialData = Flux.just(this.data);
          Flux<Data> periodicData = Flux.interval(Duration.ofSeconds(60))
                  .map(tick -> this.data);
          return Flux.concat(initialData, periodicData);
      }
      
    2. 请求处理器:
      public Mono<ServerResponse> streamData(ServerRequest request) {
          return ServerResponse.ok()
                  .contentType(MediaType.TEXT_EVENT_STREAM)
                  .body(DataService.getStream(), Data.class);
      }
      
    3. 路由配置:
      RouterFunction<?> routes (RequestHandler requestHandler) {
          return RouterFunctions
                  .route(RequestPredicates
                          .GET("/data"),
                          requestHandler::streamData);
      }
      
  • Postman测试可正常订阅数据流,但Angular客户端订阅时触发异常,客户端代码:
    getStream() {
        return this.http.get<Data>(this.apiUrl+'/data', {observe : 'events'}).subscribe(
          (event : HttpEvent<Data>) => {
            if (event.type === HttpEventType.Response){
              console.log('Message : ', event.body);
            }
          }
        );
      }
    

问题根源

  1. 事件类型判断错误:HttpEventType.Response仅在整个响应完成时触发,而SSE是持续分块推送的数据流,每个推送事件对应的类型是HttpEventType.ServerSentEvent,原代码无法捕获实时推送的数据。
  2. 响应类型配置缺失:HttpClient默认解析逻辑不匹配SSE的文本流格式,需要明确指定响应类型以正确解析数据。

修复方案

方案1:调整HttpClient配置适配SSE

修改Angular客户端代码,指定响应类型并监听正确的事件类型:

import { HttpClient, HttpEvent, HttpEventType } from '@angular/common/http';
import { Data } from './data.model'; // 导入你的Data模型定义

// ...

getStream() {
  return this.http.get<Data>(`${this.apiUrl}/data`, {
    observe: 'events',
    responseType: 'text' as 'json', // 强制指定响应类型为文本,兼容TypeScript类型检查
    reportProgress: true
  }).subscribe((event: HttpEvent<Data>) => {
    // 捕获SSE实时推送事件
    if (event.type === HttpEventType.ServerSentEvent) {
      const data = JSON.parse(event.data) as Data;
      console.log('Received stream data:', data);
    }
    // 监听连接关闭事件
    else if (event.type === HttpEventType.Response) {
      console.log('SSE connection closed');
    }
  }, error => {
    console.error('SSE subscription error:', error);
  });
}

方案2:使用原生EventSource(推荐)

浏览器原生EventSource API专门用于处理SSE数据流,实现更简洁稳定:

import { Data } from './data.model';

// ...

getStream() {
  const eventSource = new EventSource(`${this.apiUrl}/data`);

  // 监听SSE推送的消息
  eventSource.onmessage = (event) => {
    const data = JSON.parse(event.data) as Data;
    console.log('Received stream data:', data);
  };

  // 监听连接建立事件
  eventSource.onopen = () => {
    console.log('SSE connection established');
  };

  // 监听连接异常
  eventSource.onerror = (error) => {
    console.error('SSE connection error:', error);
    eventSource.close(); // 出现异常时关闭连接,可按需添加重连逻辑
  };

  // 组件销毁时关闭连接,避免内存泄漏
  // 在组件的ngOnDestroy生命周期中调用:eventSource.close();
}

服务端补充配置(若遇跨域问题)

如果Angular客户端和Spring Boot服务端不在同一域名,需配置CORS允许SSE请求:

@Bean
public CorsWebFilter corsWebFilter() {
    CorsConfiguration config = new CorsConfiguration();
    config.setAllowedOrigins(List.of("http://your-angular-domain:port")); // 替换为你的Angular应用地址
    config.setAllowedMethods(List.of("GET"));
    config.setAllowedHeaders(List.of("*"));
    config.setAllowCredentials(true);
    config.setMaxAge(3600L);

    UrlBasedCorsConfigurationSource source = new UrlBasedCorsConfigurationSource();
    source.registerCorsConfiguration("/**", config);

    return new CorsWebFilter(source);
}

内容的提问来源于stack exchange,提问作者Farfetch'd

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 20:10:30