Spring WebFlux R2DBC插入超256条记录报请求队列超限异常
256这个阈值刚好对应Reactor框架flatMap操作符的默认并发上限(内置常量Queues.SMALL_BUFFER_SIZE = 256)。之前修改Flux.create搭配OverflowStrategy.BUFFER的方案没有解决问题,是因为该溢出策略仅控制Flux生产端的缓存逻辑,流到flatMap环节时,操作符还是会默认同时订阅最多256个内部的save发布者,把请求全部推给下游R2DBC层。当插入数据量超过256时,短时间内堆积的数据库请求会超出R2DBC连接池的请求队列上限,直接抛出Cannot exchange messages because the request queue limit is exceeded异常。
以下两种方案均完全保留响应式非阻塞特性,可根据业务场景选择:
方案1:限制flatMap并发度(改动量最小)
给flatMap显式传入并发度参数,将同时执行的数据库写入操作数量限制在R2DBC连接池可承载的范围内,从根源上避免请求瞬间打满连接池队列。原代码里的Flux.fromIterable本身没有问题,不需要额外包装一层Flux.create。
实现代码:
if (employeeList == null || employeeList.isEmpty()) return Flux.empty(); return Flux.fromIterable(employeeList) .doOnNext(employee -> { employee.setParentResourceId(parentResourceId); employee.setCreatedBy(userId); }) .map(employeeMapper::entityFromModel) // 并发度根据R2DBC连接池配置的最大连接数调整,通常设为连接数的1~2倍即可,比如最大连接为10时设10~20 .flatMap(employeeRepository::save, 16) .map(employeeMapper::fromEntity);
如果对吞吐量要求不高,也可以直接把flatMap替换为concatMap,该操作符会严格按顺序串行执行所有写入操作,并发度固定为1,完全不会出现请求堆积的问题,只是写入效率会低于配置了合理并发度的flatMap。
方案2:使用R2DBC批量写入(性能最优)
循环单条插入本身IO开销很高,数据量较大时优先选择非阻塞批量写入API,一次数据库请求完成整批数据写入,既不会触发队列超限问题,写入性能也比单条循环插入高1~2个量级。
基于Spring Data R2DBC的实现参考:
if (employeeList == null || employeeList.isEmpty()) return Flux.empty(); // 提前统一赋值 employeeList.forEach(employee -> { employee.setParentResourceId(parentResourceId); employee.setCreatedBy(userId); }); List<EmployeeEntity> entities = employeeList.stream() .map(employeeMapper::entityFromModel) .toList(); // 执行批量插入 return databaseClient.sql("INSERT INTO employee (parent_resource_id, created_by, /* 其余业务字段 */) " + "VALUES (:parentResourceId, :createdBy, /* 其余字段占位符 */)") .filter(statement -> { for (EmployeeEntity entity : entities) { statement.bind("parentResourceId", entity.getParentResourceId()) .bind("createdBy", entity.getCreatedBy()) // 绑定其余业务字段 .add(); } return statement; }) .map((row, metadata) -> { // 映射返回结果为Employee模型 return employeeMapper.fromEntity(/* 从row读取字段组装实体 */); }) .all();
- 不要为了快速解决问题在响应式链路中调用
block()、同步批量写入等阻塞方法,会直接破坏WebFlux的非阻塞线程模型 flatMap的并发值不要设置得远高于R2DBC连接池的最大连接数,否则还是会出现请求排队超限的问题- 不要在
doOnNext等钩子操作里写阻塞逻辑,会拖慢整个响应式链路的执行效率
内容的提问来源于stack exchange,提问作者Sooraj V

