如何将使用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
相关产品推荐
相关产品推荐

