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

Spring Boot WebFlux+Angular:如何消费并渲染响应式API数据?

问题解决:响应式加载用户数据+实时更新

核心问题分析

  1. 一次性加载而非流式显示:后端同时引入spring-boot-starter-web和spring-boot-starter-webflux,导致WebFlux退化为Servlet栈,流式Flux被缓冲为完整响应后一次性返回;前端用HttpClient.get<UserInfo[]>会等待整个响应完成才触发回调。
  2. 无法实时更新新增用户:后端没有实现数据变更的推送机制,前端也没有保持长连接监听新数据。

后端调整

1. 移除冲突依赖

修改build.gradle,删除spring-boot-starter-web(WebFlux已包含响应式Web支持,两者共存会切换到非响应式Servlet环境):

dependencies {
    // 移除这行:implementation 'org.springframework.boot:spring-boot-starter-web'
    implementation 'org.springframework.boot:spring-boot-starter-actuator'
    implementation 'org.springframework.boot:spring-boot-starter-data-r2dbc'
    implementation 'org.springframework.boot:spring-boot-starter-webflux'
    compileOnly 'org.projectlombok:lombok'
    developmentOnly 'org.springframework.boot:spring-boot-devtools'
    runtimeOnly 'org.postgresql:postgresql'
    runtimeOnly 'org.postgresql:r2dbc-postgresql'
    annotationProcessor 'org.projectlombok:lombok'
    testImplementation 'org.springframework.boot:spring-boot-starter-test'
    testImplementation 'io.projectreactor:reactor-test'
}

2. 配置流式响应(SSE)

修改控制器,指定响应类型为text/event-stream,确保数据流式传输:

import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping(value = "/api/userinfo")
@CrossOrigin(origins = "*")
public class WebController {
    private final UserInfoService userInfoService;

    // 构造注入替代@Autowired
    public WebController(UserInfoService userInfoService) {
        this.userInfoService = userInfoService;
    }

    @GetMapping(value = "/users", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<UserInfo> getAll() {
        return userInfoService.getAll();
    }
}

3. 实现实时数据推送

修改服务类,用EmitterProcessor实现用户新增的推送逻辑:

import reactor.core.publisher.EmitterProcessor;
import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import org.springframework.stereotype.Service;

@Service
public class UserInfoService {
    private final UserInfoRepository repo;
    private final EmitterProcessor<UserInfo> userEventProcessor;
    private final FluxSink<UserInfo> userEventSink;

    public UserInfoService(UserInfoRepository repo) {
        this.repo = repo;
        this.userEventProcessor = EmitterProcessor.create(false);
        this.userEventSink = userEventProcessor.sink();
    }

    // 先返回所有已有用户,再持续推送新增用户
    public Flux<UserInfo> getAll() {
        return repo.findAll().concatWith(userEventProcessor);
    }

    // 新增用户时推送事件
    public Mono<UserInfo> createUser(UserInfo user) {
        return repo.save(user)
                .doOnSuccess(savedUser -> userEventSink.next(savedUser));
    }
}

前端调整

1. 修改服务类,支持流式接收SSE

更新HttpServiceService,用EventSource处理SSE流式数据:

import { Injectable } from '@angular/core';
import { Observable } from 'rxjs';
import { UserInfo } from '../model/userinfo';

@Injectable({
  providedIn: 'root'
})
export class HttpServiceService {
  private baseUrl = "http://localhost:9095/api/userinfo";

  constructor() {}

  getUsers(): Observable<UserInfo> {
    const eventSource = new EventSource(`${this.baseUrl}/users`);
    return new Observable(observer => {
      // 接收单条用户数据
      eventSource.onmessage = (event) => {
        const user = JSON.parse(event.data) as UserInfo;
        observer.next(user);
      };

      // 处理错误
      eventSource.onerror = (error) => {
        observer.error(error);
        eventSource.close();
      };

      // 销毁时关闭连接
      return () => eventSource.close();
    });
  }

  // 新增用户的方法(对应后端createUser接口)
  createUser(user: UserInfo): Observable<UserInfo> {
    // 实现POST请求,此处省略具体代码
  }
}

2. 组件逐步加载数据

修改DashboardComponent,逐步添加用户到列表而非一次性赋值:

import { Component, OnInit } from '@angular/core';
import { UserInfo } from 'src/app/model/userinfo';
import { HttpServiceService } from 'src/app/services/http-service.service';

@Component({
  selector: 'app-dashboard',
  templateUrl: './dashboard.component.html',
  styleUrls: ['./dashboard.component.scss']
})
export class DashboardComponent implements OnInit {
  count = 0;
  usersList: UserInfo[] = [];

  constructor(private service: HttpServiceService) {}

  ngOnInit(): void {
    this.service.getUsers().subscribe({
      next: (user) => {
        this.usersList.push(user);
        this.count = this.usersList.length;
      },
      error: (err) => console.error('加载用户失败:', err),
      complete: () => console.log('初始用户加载完成,持续监听新增用户')
    });
  }
}

额外注意事项

  • 确保application.yml中使用R2DBC配置而非JDBC:
    spring:
      r2dbc:
        url: r2dbc:postgresql://localhost:5432/your_db_name
        username: your_username
        password: your_password
    
  • EventSource默认会自动重连,可根据业务需求调整重连逻辑;
  • 10万条数据流式加载时,可在后端添加delayElements(Duration.ofMillis(10))模拟分批,避免前端渲染压力过大。

内容的提问来源于stack exchange,提问作者Thomson Mathew

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 17:09:46