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

如何将使用RESTFul的Spring Boot微服务迁移至Spring WebFlux响应式架构

Spring WebFlux迁移性能无提升解决方案

核心问题原因

你当前的改造仅在接口返回值层面包装了Flux,底层的Service、Repository逻辑依然是基于阻塞式Spring JDBC实现,请求处理过程中会占用WebFlux的事件循环线程等待数据库返回,完全没有发挥WebFlux非阻塞的特性,因此压测不会有性能提升。

完整改造步骤

1. 依赖替换

移除原有的spring-boot-starter-web、spring-jdbc依赖,引入WebFlux和响应式数据库访问依赖:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<dependency>
    <groupId>com.oracle.database.r2dbc</groupId>
    <artifactId>oracle-r2dbc</artifactId>
    <scope>runtime</scope>
</dependency>

同时修改配置文件,把原来的JDBC数据源配置改成R2DBC配置:

spring:
  r2dbc:
    url: r2dbc:oracle://<host>:<port>/<service-name>
    username: <username>
    password: <password>
    pool:
      max-size: 20 # 根据实际业务调整连接池大小

2. Repository层改造

替换阻塞的JdbcTemplate为R2dbcDatabaseClient,方法返回值改为响应式类型Flux<PersonSaverItem>:

public Flux<PersonSaverItem> getPublishItemsBycountryId(Integer countryId, PublishItemState state, int limit, int offset) {
    StringBuilder sql = new StringBuilder("select * from person where country_id = :countryId ");
    if (state != null) {
        sql.append("and state = :state ");
    }
    sql.append("limit :limit offset :offset");
    
    DatabaseClient.GenericExecuteSpec spec = databaseClient.sql(sql.toString())
            .bind("countryId", countryId)
            .bind("limit", limit)
            .bind("offset", offset);
    
    if (state != null) {
        spec = spec.bind("state", state.name()); // 按数据库实际存储的枚举类型调整转换逻辑
    }
    return spec.mapProperties(PersonSaverItem.class).all();
}

3. Service层改造

调整方法返回值为响应式类型,使用响应式算子处理异常,不要使用阻塞式try-catch:

@Override
public Flux<PersonSaverItem> getPublishItemListBy(Integer countryId, PublishItemState state, int limit, int offset) {
    return personSaverRepository.getPublishItemsBycountryId(countryId, state, limit, offset)
            .onErrorResume(e -> {
                // 自定义异常处理逻辑,比如日志打印、业务异常转换
                log.error("查询人员列表失败 countryId:{}", countryId, e);
                return Flux.error(new BusinessException("数据查询失败"));
            });
}

4. Controller层改造

直接返回Service层的Flux结果即可,不需要手动包装集合:

@GetMapping(value = "/world/{countryId}/persons", produces = MediaType.APPLICATION_JSON_VALUE)
public Flux<PersonSaverItem> getPublishPersonsByCountry(@PathVariable("countryId") Integer countryId,
                                                        @RequestParam(value = "state", required = false) PublishItemState state,
                                                        @RequestParam(value = "limit", defaultValue = "50", required = false) int limit,
                                                        @RequestParam(value = "offset", defaultValue = "0", required = false) int offset) {
    return personSaverService.getPublishItemListBy(countryId, state, limit, offset);
}

注意事项

  • 全链路非阻塞是WebFlux性能提升的核心,整个调用链中不能出现Thread.sleep、同步IO、阻塞集合操作等逻辑,否则依然会卡住事件循环线程,无法提升并发能力。
  • 如果存在暂时无法改造的阻塞遗留代码,可以将阻塞逻辑放到Schedulers.boundedElastic()调度器中执行,避免占用事件循环线程:
    // 临时兼容方案,性能弱于全链路非阻塞
    return Mono.fromCallable(() -> personSaverService.getPublishItemListBy(countryId, state, limit, offset))
            .subscribeOn(Schedulers.boundedElastic())
            .flatMapIterable(PersonSaverList::getpersonSaverList);
    
  • 压测前需要根据服务器配置调整Netty工作线程数、R2DBC连接池大小等参数,避免资源瓶颈导致性能无法体现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 18:06:04