尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

R2DBC实战:Spring Boot响应式部门管理从零实现

R2DBC实战:Spring Boot响应式部门管理从零实现 最近在整理一个“由零开始写一个数字中枢”的系列先拿部门Department管理开刀。原因很直接再复杂的数字中枢落到数据库里都离不开组织、用户、权限这些基础主数据而部门往往是最先要建出来的那张表。这个基础篇我用的是 R2DBC而不是传统 JDBC 或 JPA所以整条链路都是响应式写法从 Repository 到 Service 再到 Controller到处都是 Mono 和 Flux。适合已经写过 Spring Boot CRUD、想了解 R2DBC 实际用法、或者准备把组织架构模块落地的同学。这篇文章不会只贴代码我会把环境依赖、表结构、映射方式、查询方法、接口返回、常见报错和后续扩展一起讲清楚。1. 先想清楚数字中枢里的部门管理要解决什么问题1.1 数字中枢的实体关系为什么从部门开始“数字中枢”这个名字听起来很大但如果真的落到代码里第一步不是设计微服务也不是引入一堆中间件而是先把基础主数据理清楚。部门就是最典型的基础主数据之一。用户会挂在部门下角色会关联部门范围数据权限会按部门隔离甚至审批流、工单、任务分配都会把部门作为归属条件。所以部门管理不是一个简单的 CRUD Demo它在整个系统里承担的是“口径统一”的作用。如果部门表字段设计得乱后面做用户、组织架构、权限都会跟着乱。比如部门编码可能需要唯一父部门用来表达层级关系排序值用来控制展示顺序创建时间、更新时间用来做数据追踪和同步。这些字段看起来基础但在响应式数据库的写法里每一个字段映射都需要考虑清楚。我建议在动手写 R2DBC 之前先画一个最小实体关系图部门、用户、角色、数据权限。第一版不需要做得很复杂部门表先把 id、name、code、parent_id、sort_order、created_at、updated_at 这几个字段稳住后面的模块都能复用这套习惯。1.2 为什么用 R2DBC 而不是 JDBC 或 JPA传统 Spring Boot 项目里大部分人用的是 Spring Data JPA 或 MyBatis。JPA 很成熟但它是基于 JDBC 的阻塞模型。一个请求查询数据库时对应的线程会等数据库返回这个等待过程对操作系统来说就是线程挂起。并发量上来之后要么增加线程池要么接受延迟变高。R2DBC 的全称是 Reactive Relational Database Connectivity它的核心不是“快”而是“不阻塞”。同样的线程在等待数据库期间可以继续处理其他请求的 IO等到数据库结果回来后再通过事件机制继续执行。这正好和 Spring WebFlux 的响应式链路搭配。这里放一张直观对比对比点JDBC / JPAR2DBC线程模型一个请求占用一个线程等待数据库返回事件驱动数据库等待期间线程可以处理其他任务返回类型List、Page、OptionalMono、Flux典型使用场景传统接口、内部管理系统高并发 IO、网关、实时数据服务学习成本熟悉 SQL 和 ORM 即可还需要理解响应式流和背压事务处理声明式事务成熟需要额外注意事务边界但有一点要先说明R2DBC 不是银弹。如果业务场景就是低并发、内部管理系统、报表导出JDBC 可能更简单直接。R2DBC 更适合以后要做高并发、多数据源、响应式链路或者想统一 Web 层和数据库层为非阻塞模型的场景。数字中枢这类系统大概率会涉及实时看板、在线审批、多端接入所以从部门管理开始就切换到 R2DBC后面扩展方向不会跑偏。2. 搭建 R2dbc 项目依赖、配置和表结构2.1 依赖怎么选我用的是 Spring Boot 的响应式技术栈所以先加上 WebFlux 和 Data R2DBC 两个核心 starter。下面是示例 pom 依赖dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-webflux/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-r2dbc/artifactId /dependency dependency groupIdio.r2dbc/groupId artifactIdr2dbc-h2/artifactId scoperuntime/scope /dependency dependency groupIdcom.h2database/groupId artifactIdh2/artifactId scoperuntime/scope /dependency /dependencies本地演示用 H2 内存库最方便不依赖外部数据库。生产环境换成 PostgreSQL 或 MySQL 时把r2dbc-h2和h2换成对应驱动即可比如org.postgresql:r2dbc-postgresql或io.asyncer:r2dbc-mysql。具体版本以你实际引入的 Spring Boot 版本为准。这里为什么要用spring-boot-starter-webflux因为如果 Controller 方法返回FluxDepartment和MonoDepartment用 WebFlux 是最自然的配合。R2DBC 的 Repository 返回的就是响应式流WebFlux 的响应式处理链路可以一直保持非阻塞。2.2 连接配置和连接池参数基础配置写进 application.ymlspring: r2dbc: url: r2dbc:h2:mem:///digitalhub;DB_CLOSE_DELAY-1 username: sa password: pool: initial-size: 2 max-size: 10 sql: init: mode: always schema-locations: classpath:schema.sql连接池参数不只是启动参数它直接影响并发能力。initial-size表示服务启动时先建立的连接数max-size是连接池最大连接数。对于内存库2 到 10 足够生产库不要照抄要根据最大并发、单查询耗时、数据库 CPU 连接数限制来调整。关键点R2DBC 的连接池不是越大越好。连接数太大数据库侧反而会变慢因为每个连接都要占用数据库内存和线程。排障时如果发现“连接池耗尽”先看是不是某处代码调用了block()阻塞了线程而不是盲目把max-size调大。2.3 部门表结构设计我这里在 classpath 下放一个 schema.sqlH2 启动时会自动执行CREATE TABLE IF NOT EXISTS department ( id BIGINT AUTO_INCREMENT PRIMARY KEY, name VARCHAR(128) NOT NULL, code VARCHAR(64) NOT NULL UNIQUE, parent_id BIGINT NULL, sort_order INT DEFAULT 0, created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP );这套字段是部门管理的最小可用集id自增主键R2DBC 里用Id标记。name部门名称不能为空。code部门编码唯一业务上用于导入、同步、权限匹配。parent_id父部门 ID顶级部门可以为空。sort_order排序值越小越靠前。created_at / updated_at创建和更新时间默认交给数据库更稳。如果用的是 PostgreSQLBIGINT AUTO_INCREMENT在 PostgreSQL 里不适用要换成BIGSERIAL。如果用的是 MySQL则要保证数据库版本的AUTO_INCREMENT行为和表引擎一致。这里给的是 H2 示例生产环境换库时schema 要按对应数据库方言调整。3. 写实体和 RepositoryR2dbc 的对象映射与查询方法3.1 实体用注解还是约定Spring Data R2DBC 的实体写法更像 Spring Data JDBC而不是 JPA。它没有复杂的一对多、多对多关系映射想表达关联关系时很多人在 Service 层手动处理。部门表比较简单实体代码如下import org.springframework.data.annotation.Id; import org.springframework.data.relational.core.mapping.Column; import org.springframework.data.relational.core.mapping.Table; import java.time.LocalDateTime; Table(department) public class Department { Id private Long id; private String name; private String code; Column(parent_id) private Long parentId; Column(sort_order) private Integer sortOrder; Column(created_at) private LocalDateTime createdAt; Column(updated_at) private LocalDateTime updatedAt; public Long getId() { return id; } public void setId(Long id) { this.id id; } public String getName() { return name; } public void setName(String name) { this.name name; } public String getCode() { return code; } public void setCode(String code) { this.code code; } public Long getParentId() { return parentId; } public void setParentId(Long parentId) { this.parentId parentId; } public Integer getSortOrder() { return sortOrder; } public void setSortOrder(Integer sortOrder) { this.sortOrder sortOrder; } public LocalDateTime getCreatedAt() { return createdAt; } public void setCreatedAt(LocalDateTime createdAt) { this.createdAt createdAt; } public LocalDateTime getUpdatedAt() { return updatedAt; } public void setUpdatedAt(LocalDateTime updatedAt) { this.updatedAt updatedAt; } }字段说明Table(department)指定表名。Id必须加在数据库主键字段上否则保存时可能报错。Column用来把 Java 字段映射到数据库列尤其下划线字段一定要写清楚。实体里没有用 Lombok是为了减少依赖。如果你项目里已经用了 Lombok可以换成Data但要注意Id的注解仍然需要保留。为什么不建议完全依赖约定映射因为一旦表结构里某个列名和 Java 字段对不上运行时查出来就是空值或者报 column not found。显式写Column虽然多几行代码但排查问题时能少一层不确定。3.2 Repository 接口怎么定义R2DBC 的 Repository 继承ReactiveCrudRepository方法返回类型是Mono或Fluximport org.springframework.data.repository.reactive.ReactiveCrudRepository; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; public interface DepartmentRepository extends ReactiveCrudRepositoryDepartment, Long { MonoDepartment findByCode(String code); FluxDepartment findByParentId(Long parentId); FluxDepartment findByParentIdOrderBySortOrderAsc(Long parentId); MonoBoolean existsByCode(String code); }使用时有几个原则如果查询可能返回多行返回类型必须用Flux不要用Mono。如果业务上约定只返回一行比如findByCode可以用Mono。判断是否存在用existsByCode返回MonoBoolean不要先findByCode再判断非空这样会多一次查询。排序可以直接写在方法名上比如OrderBySortOrderAsc适合简单的固定排序。R2DBC 之所以用响应式类型是为了让查询不会阻塞当前线程。如果自定义实现里有人调repository.findAll().block(),那就破坏了非阻塞链路。下面会在 Service 层示范正确写法。3.3 复杂查询用 DatabaseClient 还是 R2dbcEntityTemplateRepository 方法适合简单条件。如果查询条件多、动态拼接、或者要做聚合查询我一般用R2dbcEntityTemplate或DatabaseClient。R2dbcEntityTemplate适合把查询结果映射回实体的场景import org.springframework.data.r2dbc.core.R2dbcEntityTemplate; import org.springframework.data.relational.core.query.Criteria; import org.springframework.data.relational.core.query.Query; import org.springframework.stereotype.Repository; import reactor.core.publisher.Flux; Repository public class DepartmentSearchRepository { private final R2dbcEntityTemplate template; public DepartmentSearchRepository(R2dbcEntityTemplate template) { this.template template; } public FluxDepartment searchByName(String keyword) { return template.select(Department.class) .matching(Query.query(Criteria.where(name).like(% keyword %))) .all(); } }DatabaseClient适合直接写原生 SQL。它更灵活但映射需要自己处理适合报表、复杂 join、多表统计等场景import org.springframework.r2dbc.core.DatabaseClient; import org.springframework.stereotype.Repository; import reactor.core.publisher.Flux; Repository public class DepartmentSqlRepository { private final DatabaseClient databaseClient; public DepartmentSqlRepository(DatabaseClient databaseClient) { this.databaseClient databaseClient; } public FluxDepartment listByKeyword(String keyword) { String sql SELECT id, name, code, parent_id, sort_order, created_at, updated_at FROM department WHERE name LIKE :keyword ORDER BY sort_order ASC ; return databaseClient.sql(sql) .bind(keyword, % keyword %) .map((row, metadata) - { Department department new Department(); department.setId(row.get(id, Long.class)); department.setName(row.get(name, String.class)); department.setCode(row.get(code, String.class)); department.setParentId(row.get(parent_id, Long.class)); department.setSortOrder(row.get(sort_order, Integer.class)); department.setCreatedAt(row.get(created_at, LocalDateTime.class)); department.setUpdatedAt(row.get(updated_at, LocalDateTime.class)); return department; }) .all(); } }简单条件用 Repository 方法复杂动态查询用R2dbcEntityTemplate原生 SQL 和报表用DatabaseClient。不要一上来所有地方都用原生 SQL否则后面的维护成本会很高。4. 服务层和 Web 接口Mono/Flux 的正确打开方式4.1 服务层做校验、默认值和转换实体对象不应该直接暴露给前端。前端请求先到一个 Request 对象Service 层做转换、校验、默认值填充然后保存实体。这里用一个简单的 Service 示例import org.springframework.stereotype.Service; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.time.LocalDateTime; Service public class DepartmentService { private final DepartmentRepository departmentRepository; public DepartmentService(DepartmentRepository departmentRepository) { this.departmentRepository departmentRepository; } public FluxDepartment list(Long parentId) { if (parentId ! null) { return departmentRepository.findByParentIdOrderBySortOrderAsc(parentId); } return departmentRepository.findAll(); } public MonoDepartment detail(Long id) { return departmentRepository.findById(id); } public MonoDepartment create(CreateDepartmentRequest request) { Department department new Department(); department.setName(request.getName()); department.setCode(request.getCode()); department.setParentId(request.getParentId()); department.setSortOrder(request.getSortOrder() null ? 0 : request.getSortOrder()); department.setCreatedAt(LocalDateTime.now()); department.setUpdatedAt(LocalDateTime.now()); return departmentRepository.save(department); } public MonoDepartment update(Long id, UpdateDepartmentRequest request) { return departmentRepository.findById(id) .flatMap(department - { department.setName(request.getName()); department.setCode(request.getCode()); department.setParentId(request.getParentId()); department.setSortOrder(request.getSortOrder()); department.setUpdatedAt(LocalDateTime.now()); return departmentRepository.save(department); }); } public MonoVoid delete(Long id) { return departmentRepository.deleteById(id); } }注意update方法里的flatMap先找到实体再修改再保存。如果直接new Department再save很可能把已有数据覆盖掉。R2DBC 的save在有Id的情况下会尝试更新但没有值缺失风险。另外createdAt和updatedAt这里由应用层赋值。生产环境更推荐数据库默认值加触发器或者用数据库CURRENT_TIMESTAMP避免多个服务实例时间不一致。但是如果你选择了数据库默认值实体查询出来时字段会有值插入时应用层可以不设置。4.2 Controller 返回响应式类型会有什么差别Controller 示例import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; RestController RequestMapping(/api/departments) public class DepartmentController { private final DepartmentService departmentService; public DepartmentController(DepartmentService departmentService) { this.departmentService departmentService; } GetMapping public FluxDepartment list(RequestParam(required false) Long parentId) { return departmentService.list(parentId); } GetMapping(/{id}) public MonoDepartment detail(PathVariable Long id) { return departmentService.detail(id); } PostMapping public MonoDepartment create(RequestBody CreateDepartmentRequest request) { return departmentService.create(request); } PutMapping(/{id}) public MonoDepartment update(PathVariable Long id, RequestBody UpdateDepartmentRequest request) { return departmentService.update(id, request); } DeleteMapping(/{id}) public MonoVoid delete(PathVariable Long id) { return departmentService.delete(id); } }和传统 JPA 的区别是返回类型不再是ListDepartment而是FluxDepartment返回单条记录时是MonoDepartment删除返回MonoVoid表示一个没有结果的响应式信号。从前端调用角度看接口路径和 JSON 结构没有本质区别。区别在框架内部Flux是数据流可以边查询边推送给客户端Mono是 0 或 1 个结果。如果你在 Controller 里把Flux转成List还是要调用collectList()那会等所有数据到齐才返回响应式优势会打折扣。4.3 为什么不能在响应式链路里用 block这是一个非常常见的坑。新手写 R2DBC 时很容易在 Service 或 Controller 里写出这种代码Department department departmentRepository.findById(1L).block();如果运行在 WebFlux 的 reactor 线程上block()会让线程停下来等结果。这个问题在低并发下可能看不出异常但并发一旦上来线程池会被占满请求开始排队最后表现为接口超时、连接池耗尽、甚至死锁。正确做法是让Mono和Flux沿着调用链走完把最终结果交给 Spring WebFlux 的响应式链路处理。如果确实需要把Flux收集为MonoListT用collectList()而不是block()MonoListDepartment listMono departmentRepository.findAll().collectList();这条规则几乎适用于所有 R2DBC 代码。遇到“启动时报 block 操作不可用”或“请求卡住”优先排查是不是有block()出现在响应式链路里。5. 验证请求、监控资源与排查链路5.1 最小验证从启动到查询一条记录先把项目启动起来。启动成功后用 curl 做最小验证不需要打开任何前端页面。先查询列表空数据时返回空数组curl http://localhost:8080/api/departments再创建一条部门数据curl -X POST http://localhost:8080/api/departments \ -H Content-Type: application/json \ -d { name: 研发部, code: DEV, parentId: null, sortOrder: 1 }如果返回 JSON 里包含生成的 id 和 created_at说明链路已经通了。再查一次列表应该能看到刚才创建的部门。验证时我一般分三步先确认服务能启动日志里没有初始化错误。再跑一次查询确认表结构和实体映射没问题。然后跑一次插入确认主键生成、字段填充正常。不要一上来就写几十个接口、做批量导入。先把单条链路跑稳后面的功能都会轻松很多。5.2 常见报错和排查顺序R2DBC 项目刚开始跑起来时问题往往不是模型算法问题而是环境、依赖、映射问题。这里列一个我常用的排查表现象排查方向常见原因启动报错找不到数据库驱动检查 pom 依赖缺少r2dbc-h2或对应数据库驱动查询时表不存在检查 schema.sql 是否执行Spring SQL init 没配置或表名拼错字段报错sort_order not found检查实体 ColumnJava 属性名和数据库列名不一致使用block()时报错线程阻塞检查全链路调用响应式链路里不应出现 block请求很慢然后超时检查连接池和日志连接池太小或数据库查询没有索引插入时主键没有生成检查实体 Id 和表主键主键字段没有标记 Id或数据库方言配置不对排查顺序建议固定先看完整异常栈定位是启动阶段还是请求阶段。再看 SQL确认是不是表名、列名写错。再确认表结构是否真的创建成功。最后检查是不是把阻塞代码写进了响应式链路。不要一上来就怀疑 R2DBC 框架本身。多数问题都是依赖版本、schema 初始化、字段映射、block 调用这四类。6. 从部门管理延伸到批量、事务和后续扩展6.1 批量导入部门的 R2dbc 写法部门管理很快会遇到批量导入场景比如初始化组织架构或者从 Excel 导入部门。这时不要用 for 循环逐条save并block()应该用saveAllpublic FluxDepartment batchCreate(ListCreateDepartmentRequest requests) { ListDepartment departments requests.stream().map(request - { Department department new Department(); department.setName(request.getName()); department.setCode(request.getCode()); department.setParentId(request.getParentId()); department.setSortOrder(request.getSortOrder() null ? 0 : request.getSortOrder()); department.setCreatedAt(LocalDateTime.now()); department.setUpdatedAt(LocalDateTime.now()); return department; }).toList(); return departmentRepository.saveAll(departments); }saveAll会以响应式方式依次写入返回FluxDepartment。使用批量导入时有几个点不能忽略批量数据要提前校验比如部门编码是否重复、父部门是否存在。如果数据量大建议分批比如每 200 条一批避免单个事务过大。批量写入后要处理部分失败场景哪些成功、哪些失败、失败原因是什么不能只看最后返回的Flux是否正常。R2DBC 的saveAll不是银弹。它只负责写入不负责业务幂等。如果一条数据因为重复编码失败整个Flux会往 error 信号走你需要用onErrorContinue或分组处理来跳过单条错误。6.2 事务处理的关键点R2DBC 也支持事务但写法比 JPA 更需要注意边界。Spring 里可以配合 R2DBC 事务管理器使用Transactional不过更稳妥的方式是使用TransactionalOperator尤其是返回Flux的批量场景import org.springframework.transaction.reactive.TransactionalOperator; import reactor.core.publisher.Flux; public FluxDepartment createParentAndChildren(Department parent, ListDepartment children) { return departmentRepository.save(parent) .flatMap(savedParent - { children.forEach(child - child.setParentId(savedParent.getId())); return departmentRepository.saveAll(children); }) .as(transactionalOperator::transactional); }这里的as(transactionalOperator::transactional)会把整个响应式流包在一个事务里。流正常执行完成才提交任何一步出错都会回滚。需要注意如果方法返回的是Flux事务边界会随着Flux的完成信号结束。不要在事务里做阻塞调用也不要把一个长流无限推给客户端因为事务持有连接的时间会变长。实际业务中父部门和子部门一起创建是很常见的场景。如果你不在事务里处理可能出现父部门保存成功、子部门保存失败最后遗留不完整数据。不要等到线上出问题才补事务应该在第一个批量接口时就把事务边界想清楚。6.3 分页、树结构和软删除部门管理往后延伸有三个方向一定会碰到。第一个是分页。R2DBC 的 Repository 查询方法上可以加Pageable参数返回类型仍然一般是Flux。但获取总数往往需要单独 count。如果当前 Spring Data 版本对分页支持不稳定可以直接用DatabaseClient写LIMIT/OFFSET再用聚合查询拿总数。分页这个功能不要照搬 JPA 的Page对象R2DBC 的响应式模型更适合数据流式返回。第二个是树结构。部门表用parent_id可以表达一棵树但在查询整棵树时如果层级很深递归查询很容易变成 N 次查询。常见优化方案是冗余路径字段比如path字段存储/1/2/3或者引入 closure table也就是部门闭包表。第一种适合查询次数多、更新少的场景第二种适合频繁移动部门节点的场景。部门管理第一版可以先只做parent_id和一层查询后续再做整棵树的展开。第三个是软删除。部门通常不能物理删掉因为历史用户、历史审批流可能还挂着这个部门。常见做法是增加deleted字段默认 0删除时更新为 1。但要注意findAll()默认不会过滤软删除字段必须在 Repository 方法里手动加条件比如FluxDepartment findByDeletedFalseOrderBySortOrderAsc();如果你没有加这个条件软删除就只是一个摆设。部门数据一旦被标记删除查询、同步、导入都要统一排除这个逻辑要写入查询方法不能靠业务层每次手写条件。整体来看部门管理是数字中枢里最基础、也最适合练习 R2DBC 的模块。先把单条查询和插入跑稳再逐步补批量、事务、分页和软删除。响应式编程的核心不是把代码写成纯链式而是要让每个 IO 操作都不阻塞线程同时把事务边界、错误处理、数据一致性想清楚。等到后面接入用户和权限模块时你会发现把基础打稳比加新功能更重要。
返回列表