Spring Boot 4.x响应式数据访问:R2DBC实战


Spring Boot 4.x响应式数据访问:R2DBC实战

技术原理

R2DBC(Reactive Relational Database Connectivity)是Spring Boot 4.x中响应式数据访问的标准。与传统JDBC不同,R2DBC完全非阻塞,能够在少量线程的情况下处理大量并发数据库操作,非常适合微服务和高并发场景。

核心优势

  1. 非阻塞I/O: 基于Reactor和Netty
  2. 高并发: 单个连接可以处理多个并发请求
  3. 资源效率: 相比JDBC减少了线程上下文切换
  4. 流式处理: 支持响应式流的数据处理

实践案例:构建响应式数据访问层

项目配置

<!-- pom.xml --> <dependencies> <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> <!-- R2DBC驱动 --> <dependency> <groupId>org.postgresql</groupId> <artifactId>r2dbc-postgresql</artifactId> </dependency> <!-- 连接池 --> <dependency> <groupId>io.r2dbc</groupId> <artifactId>r2dbc-pool</artifactId> </dependency> <!-- Lombok --> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <scope>provided</scope> </dependency> </dependencies>

数据库配置

# application.yml spring: r2dbc: url: r2dbc:postgresql://localhost:5432/mydb username: user password: password pool: initial-size: 5 max-size: 20 max-idle-time: 30m max-life-time: 1h max-create-connection-time: 5s acquire-timeout: 5s max-validation-time: 3s data: r2dbc: repositories: enabled: true logging: level: org.springframework.data.r2dbc: DEBUG io.r2dbc.postgresql.QUERY: DEBUG

数据模型

package com.example.demo.model; import org.springframework.data.annotation.Id; import org.springframework.data.relational.core.mapping.Table; import lombok.Data; import lombok.Builder; import lombok.AllArgsConstructor; import lombok.NoArgsConstructor; import java.time.LocalDateTime; @Data @Builder @AllArgsConstructor @NoArgsConstructor @Table("users") public class User { @Id private Long id; private String username; private String email; private String password; private LocalDateTime createdAt; private LocalDateTime updatedAt; private Boolean active; }

Repository层

package com.example.demo.repository; import com.example.demo.model.User; import org.springframework.data.r2dbc.repository.Query; import org.springframework.data.r2dbc.repository.R2dbcRepository; import org.springframework.stereotype.Repository; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Repository public interface UserRepository extends R2dbcRepository<User, Long> { // 根据用户名查找 Mono<User> findByUsername(String username); // 根据邮箱查找 Mono<User> findByEmail(String email); // 查找活跃用户 Flux<User> findByActiveTrue(); // 自定义查询 @Query("SELECT * FROM users WHERE username LIKE :pattern") Flux<User> findByUsernamePattern(String pattern); // 批量更新 @Query("UPDATE users SET active = :active WHERE id IN (:ids)") Mono<Integer> updateActiveStatusByIds(Boolean active, Long... ids); }

Service层

package com.example.demo.service; import com.example.demo.model.User; import com.example.demo.repository.UserRepository; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.time.LocalDateTime; @Service public class UserService { private final UserRepository userRepository; public UserService(UserRepository userRepository) { this.userRepository = userRepository; } // 创建用户 @Transactional public Mono<User> createUser(User user) { return userRepository.save(user.toBuilder() .createdAt(LocalDateTime.now()) .updatedAt(LocalDateTime.now()) .active(true) .build()); } // 更新用户 @Transactional public Mono<User> updateUser(Long id, User userDetails) { return userRepository.findById(id) .switchIfEmpty(Mono.error(new RuntimeException("User not found"))) .map(existingUser -> existingUser.toBuilder() .username(userDetails.getUsername()) .email(userDetails.getEmail()) .updatedAt(LocalDateTime.now()) .build()) .flatMap(userRepository::save); } // 删除用户 @Transactional public Mono<Void> deleteUser(Long id) { return userRepository.deleteById(id); } // 批量创建用户 @Transactional public Flux<User> createUsers(Flux<User> users) { return users.flatMap(user -> userRepository.save(user.toBuilder() .createdAt(LocalDateTime.now()) .updatedAt(LocalDateTime.now()) .active(true) .build())); } // 复杂查询 public Flux<User> searchUsers(String keyword) { return userRepository.findByUsernamePattern("%" + keyword + "%") .filter(User::getActive); } }

高级特性

响应式事务

package com.example.demo.service; import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.reactive.TransactionalOperator; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @Service public class TransactionalService { private final UserRepository userRepository; private final TransactionalOperator rxtx; public TransactionalService(UserRepository userRepository, TransactionalOperator rxtx) { this.userRepository = userRepository; this.rxtx = rxtx; } // 声明式事务 @Transactional public Mono<User> createUserWithAudit(User user) { return userRepository.save(user) .flatMap(savedUser -> { // 创建审计日志 return createAuditLog("CREATE", user.getId()) .thenReturn(savedUser); }); } // 编程式事务 public Mono<User> updateUserWithAudit(Long id, User userDetails) { return rxtx.execute(status -> userRepository.findById(id) .switchIfEmpty(Mono.error(new RuntimeException("User not found"))) .flatMap(existingUser -> userRepository.save( existingUser.toBuilder() .username(userDetails.getUsername()) .updatedAt(LocalDateTime.now()) .build() )) .flatMap(updatedUser -> createAuditLog("UPDATE", updatedUser.getId()) .thenReturn(updatedUser) ) ); } private Mono<Void> createAuditLog(String action, Long userId) { // 审计日志逻辑 return Mono.empty(); } }

自定义转换器

package com.example.demo.config; import io.r2dbc.spi.ConnectionFactory; import org.springframework.context.annotation.Configuration; import org.springframework.data.convert.CustomConversions; import org.springframework.data.r2dbc.convert.R2dbcCustomConversions; import org.springframework.data.r2dbc.dialect.DialectResolver; import org.springframework.data.r2dbc.dialect.R2dbcDialect; import org.springframework.lang.NonNull; import java.util.Arrays; @Configuration public class R2dbcConfig { @NonNull public R2dbcCustomConversions r2dbcCustomConversions(ConnectionFactory connectionFactory) { R2dbcDialect dialect = DialectResolver.getDialect(connectionFactory); return new R2dbcCustomConversions( dialect.getSimpleTypeHolder(), Arrays.asList( new LocalDateTimeConverter(), new JsonNodeConverter() ) ); } }

批量操作

@Service public class BatchOperationService { private final UserRepository userRepository; private final DatabaseClient databaseClient; public BatchOperationService(UserRepository userRepository, DatabaseClient databaseClient) { this.userRepository = userRepository; this.databaseClient = databaseClient; } // 批量插入 public Flux<User> batchInsert(Flux<User> users) { return databaseClient.insert() .into("users") .nullValue("id") .value("username", user -> user.getUsername()) .value("email", user -> user.getEmail()) .value("password", user -> user.getPassword()) .value("created_at", user -> LocalDateTime.now()) .value("updated_at", user -> LocalDateTime.now()) .value("active", user -> true) .then() .thenMany(Flux.empty()); } // 批量更新 public Mono<Integer> batchUpdateStatus(Flux<Long> ids, Boolean status) { return databaseClient.update() .table("users") .set("active", status) .set("updated_at", LocalDateTime.now()) .matching( Criteria.where("id").in( ids.collectList().block(Arrays.asList()) ) ) .fetch() .rowsUpdated(); } }

Controller层

package com.example.demo.controller; import com.example.demo.model.User; import com.example.demo.service.UserService; import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @RestController @RequestMapping("/api/users") public class UserController { private final UserService userService; public UserController(UserService userService) { this.userService = userService; } @GetMapping("/{id}") public Mono<User> getUser(@PathVariable Long id) { return userService.findById(id); } @GetMapping public Flux<User> getAllUsers() { return userService.findAll(); } @PostMapping public Mono<User> createUser(@RequestBody User user) { return userService.createUser(user); } @PutMapping("/{id}") public Mono<User> updateUser(@PathVariable Long id, @RequestBody User user) { return userService.updateUser(id, user); } @DeleteMapping("/{id}") public Mono<Void> deleteUser(@PathVariable Long id) { return userService.deleteUser(id); } @GetMapping("/search") public Flux<User> searchUsers(@RequestParam String keyword) { return userService.searchUsers(keyword); } }

性能优化

连接池调优

spring: r2dbc: pool: max-size: 50 # 最大连接数 initial-size: 10 # 初始连接数 max-idle-time: 30m # 最大空闲时间 max-life-time: 1h # 连接最大生命周期 max-create-connection-time: 5s # 创建连接超时 acquire-timeout: 5s # 获取连接超时 max-validation-time: 3s # 验证超时

背压处理

@Service public class StreamService { private final UserRepository userRepository; public StreamService(UserRepository userRepository) { this.userRepository = userRepository; } // 流式处理大数据集 public Flux<User> streamAllUsers() { return userRepository.findAll() .limitRate(100) // 限制请求速率 .buffer(50) // 批量处理 .flatMap(Flux::fromIterable); } }

测试

单元测试

@SpringBootTest class UserRepositoryTest { @Autowired private UserRepository userRepository; @Test void testFindByUsername() { User user = User.builder() .username("testuser") .email("test@example.com") .password("password") .build(); StepVerifier.create(userRepository.save(user)) .expectNextMatches(saved -> saved.getUsername().equals("testuser")) .verifyComplete(); StepVerifier.create(userRepository.findByUsername("testuser")) .expectNextMatches(found -> found.getEmail().equals("test@example.com")) .verifyComplete(); } }

总结

R2DBC是Spring Boot 4.x响应式编程的重要组成部分,提供了高效的非阻塞数据访问能力。通过合理的配置和设计,可以构建高性能、可扩展的数据访问层。建议在高并发场景下优先考虑R2DBC,但需要注意响应式编程的学习曲线和调试复杂度。


作者与出处
整理: 灏天文库整理
本站整理收录,版权归原作者/开源协议所有;欢迎通过原文链接访问源仓库。
发布者: 作者: 灏天学者_SSSV45的小龙虾 转发
评论区 (0)
U