springboot3+r2dbc——响应式编程实践

2023-11-06 20:11

本文主要是介绍springboot3+r2dbc——响应式编程实践,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

在这里插入图片描述

Spring boot3已经M1了,最近群佬们也开始蠢蠢欲动的开始整活Reactive+Spring Boot3,跟着大家的步伐,我也来整一篇工程入门,我们将用java17+Spring Boot3+r2dbc+Reactive栈来讲述,欢迎大家来讨论。(关于响应式,请大家异步到之前的文章里,有详细介绍。)

r2dbc

Reactor还有基于其之上的Spring WebFlux框架。包括vert.xrxjava等等reactive技术。我们实际上在应用层已经有很多优秀的响应式处理框架。

但是有一个问题就是所有的框架都需要获取底层的数据,而基本上关系型数据库的底层读写都还是同步的。

为了解决这个问题,出现了两个标准,一个是oracle提出的 ADBC (Asynchronous Database Access API),另一个就是Pivotal提出的R2DBC (Reactive Relational Database Connectivity)。

R2DBC是基于Reactive Streams标准来设计的。通过使用R2DBC,你可以使用reactive API来操作数据。

同时R2DBC只是一个开放的标准,而各个具体的数据库连接实现,需要实现这个标准。

今天我们以r2dbc-h2为例,讲解一下r2dbcSpring webFlux中的使用。

工程依赖

以下是 pom.xml清单

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"><modelVersion>4.0.0</modelVersion><parent><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-parent</artifactId><version>3.0.0-M1</version><relativePath/> <!-- lookup parent from repository --></parent><groupId>wang.datahub</groupId><artifactId>springboot3demo</artifactId><version>0.0.1-SNAPSHOT</version><name>springboot3demo</name><description>Demo project for Spring Boot</description><properties><java.version>17</java.version></properties><dependencies><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-data-r2dbc</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-data-redis-reactive</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-data-rest</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-groovy-templates</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-hateoas</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-webflux</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-configuration-processor</artifactId><optional>true</optional></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-devtools</artifactId></dependency><dependency><groupId>io.r2dbc</groupId><artifactId>r2dbc-h2</artifactId></dependency><dependency><groupId>com.h2database</groupId><artifactId>h2</artifactId></dependency><dependency><groupId>mysql</groupId><artifactId>mysql-connector-java</artifactId><scope>runtime</scope></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-test</artifactId><scope>test</scope></dependency><dependency><groupId>io.projectreactor</groupId><artifactId>reactor-test</artifactId><scope>test</scope></dependency><dependency><groupId>io.projectreactor</groupId><artifactId>reactor-test</artifactId>
<!--			<version>3.4.14</version>-->
<!--			<scope>compile</scope>--></dependency></dependencies><build><plugins><plugin><groupId>org.springframework.boot</groupId><artifactId>spring-boot-maven-plugin</artifactId></plugin></plugins></build><repositories><repository><id>spring-milestones</id><name>Spring Milestones</name><url>https://repo.spring.io/milestone</url><snapshots><enabled>false</enabled></snapshots></repository><repository><id>spring-snapshots</id><name>Spring Snapshots</name><url>https://repo.spring.io/snapshot</url><releases><enabled>false</enabled></releases></repository></repositories><pluginRepositories><pluginRepository><id>spring-milestones</id><name>Spring Milestones</name><url>https://repo.spring.io/milestone</url><snapshots><enabled>false</enabled></snapshots></pluginRepository><pluginRepository><id>spring-snapshots</id><name>Spring Snapshots</name><url>https://repo.spring.io/snapshot</url><releases><enabled>false</enabled></releases></pluginRepository></pluginRepositories></project>

配置文件

这里我们只配置了r2dbc链接信息

r2dbc:url: r2dbc:h2:mem:///test?options=DB_CLOSE_DELAY=-1;DB_CLOSE_ON_EXIT=FALSE

配置类

用于配置默认链接,创建初始化数据

package wang.datahub.springboot3demo.config;import io.netty.util.internal.StringUtil;
import io.r2dbc.spi.ConnectionFactories;
import io.r2dbc.spi.ConnectionFactory;
import io.r2dbc.spi.ConnectionFactoryOptions;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import reactor.core.publisher.Flux;
import static io.r2dbc.spi.ConnectionFactoryOptions.*;@Configuration
@ConfigurationProperties(prefix = "r2dbc")
public class DBConfig {private String url;private String user;private String password;public String getUrl() {return url;}public void setUrl(String url) {this.url = url;}public String getUser() {return user;}public void setUser(String user) {this.user = user;}public String getPassword() {return password;}public void setPassword(String password) {this.password = password;}@Beanpublic ConnectionFactory connectionFactory() {System.out.println("url ==> "+url);ConnectionFactoryOptions baseOptions = ConnectionFactoryOptions.parse(url);ConnectionFactoryOptions.Builder ob = ConnectionFactoryOptions.builder().from(baseOptions);if (!StringUtil.isNullOrEmpty(user)) {ob = ob.option(USER, user);}if (!StringUtil.isNullOrEmpty(password)) {ob = ob.option(PASSWORD, password);}return ConnectionFactories.get(ob.build());}@Beanpublic CommandLineRunner initDatabase(ConnectionFactory cf) {return (args) ->Flux.from(cf.create()).flatMap(c ->Flux.from(c.createBatch().add("drop table if exists Users").add("create table Users(" +"id IDENTITY(1,1)," +"firstname varchar(80) not null," +"lastname varchar(80) not null)").add("insert into Users(firstname,lastname)" +"values('Jacky','Li')").add("insert into Users(firstname,lastname)" +"values('Doudou','Li')").add("insert into Users(firstname,lastname)" +"values('Maimai','Li')").execute()).doFinally((st) -> c.close())).log().blockLast();}}

bean

创建用户bean

package wang.datahub.springboot3demo.bean;import org.springframework.data.annotation.Id;public class Users {@Idprivate Long id;private String firstname;private String lastname;public Users(){}public Users(Long id, String firstname, String lastname) {this.id = id;this.firstname = firstname;this.lastname = lastname;}public Long getId() {return id;}public void setId(Long id) {this.id = id;}public String getFirstname() {return firstname;}public void setFirstname(String firstname) {this.firstname = firstname;}public String getLastname() {return lastname;}public void setLastname(String lastname) {this.lastname = lastname;}@Overridepublic String toString() {return "User{" +"id=" + id +", firstname='" + firstname + '\'' +", lastname='" + lastname + '\'' +'}';}
}

DAO

dao代码清单如下,包含查询列表、按id查询,以及创建用户等操作

package wang.datahub.springboot3demo.dao;import io.r2dbc.spi.Connection;
import io.r2dbc.spi.ConnectionFactory;
import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
import org.springframework.data.relational.core.query.Query;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import wang.datahub.springboot3demo.bean.Users;import static org.springframework.data.r2dbc.query.Criteria.where;
import static org.springframework.data.relational.core.query.Query.query;@Component
public class UsersDao {private ConnectionFactory connectionFactory;private R2dbcEntityTemplate template;public UsersDao(ConnectionFactory connectionFactory) {this.connectionFactory = connectionFactory;this.template = new R2dbcEntityTemplate(connectionFactory);}public Mono<Users> findById(long id) {return this.template.selectOne(query(where("id").is(id)),Users.class);//        return Mono.from(connectionFactory.create())
//                .flatMap(c -> Mono.from(c.createStatement("select id,firstname,lastname from Users where id = $1")
//                                .bind("$1", id)
//                                .execute())
//                        .doFinally((st) -> close(c)))
//                .map(result -> result.map((row, meta) ->
//                        new Users(row.get("id", Long.class),
//                                row.get("firstname", String.class),
//                                row.get("lastname", String.class))))
//                .flatMap( p -> Mono.from(p));}public Flux<Users> findAll() {return this.template.select(Users.class).all();
//        return Mono.from(connectionFactory.create())
//                .flatMap((c) -> Mono.from(c.createStatement("select id,firstname,lastname from users")
//                                .execute())
//                        .doFinally((st) -> close(c)))
//                .flatMapMany(result -> Flux.from(result.map((row, meta) -> {
//                    Users acc = new Users();
//                    acc.setId(row.get("id", Long.class));
//                    acc.setFirstname(row.get("firstname", String.class));
//                    acc.setLastname(row.get("lastname", String.class));
//                    return acc;
//                })));}public Mono<Users> createAccount(Users account) {return Mono.from(connectionFactory.create()).flatMap(c -> Mono.from(c.beginTransaction()).then(Mono.from(c.createStatement("insert into Users(firstname,lastname) values($1,$2)").bind("$1", account.getFirstname()).bind("$2", account.getLastname()).returnGeneratedValues("id").execute())).map(result -> result.map((row, meta) ->new Users(row.get("id", Long.class),account.getFirstname(),account.getLastname()))).flatMap(pub -> Mono.from(pub)).delayUntil(r -> c.commitTransaction()).doFinally((st) -> c.close()));}private <T> Mono<T> close(Connection connection) {return Mono.from(connection.close()).then(Mono.empty());}
}

controller

controller代码清单如下,包含了查询列表、按id查询,以及创建用户等操作

package wang.datahub.springboot3demo.controller;import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Controller;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import wang.datahub.springboot3demo.bean.Users;
import wang.datahub.springboot3demo.dao.UsersDao;@RestController
public class UsersController {@Autowiredprivate final UsersDao usersDao;public UsersController(UsersDao usersDao) {this.usersDao = usersDao;}@GetMapping("/users/{id}")public Mono<ResponseEntity<Users>> getUsers(@PathVariable("id") Long id) {return usersDao.findById(id).map(acc -> new ResponseEntity<>(acc, HttpStatus.OK)).switchIfEmpty(Mono.just(new ResponseEntity<>(null, HttpStatus.NOT_FOUND)));}@GetMapping("/users")public Flux<Users> getAllAccounts() {return usersDao.findAll();}@PostMapping("/createUser")public Mono<ResponseEntity<Users>> createUser(@RequestBody Users user) {return usersDao.createAccount(user).map(acc -> new ResponseEntity<>(acc, HttpStatus.CREATED)).log();}
}

启动类清单:

package wang.datahub.springboot3demo;import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import wang.datahub.springboot3demo.config.DBConfig;@SpringBootApplication
@EnableConfigurationProperties(DBConfig.class)
public class WebFluxR2dbcApp {public static void main(String[] args) {SpringApplication.run(WebFluxR2dbcApp.class, args);}
}

好了,致此我们整个 Demo 就实现完成了

参考链接:

https://zhuanlan.zhihu.com/p/299069835

这篇关于springboot3+r2dbc——响应式编程实践的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/358809

相关文章

防止Linux rm命令误操作的多场景防护方案与实践

《防止Linuxrm命令误操作的多场景防护方案与实践》在Linux系统中,rm命令是删除文件和目录的高效工具,但一旦误操作,如执行rm-rf/或rm-rf/*,极易导致系统数据灾难,本文针对不同场景... 目录引言理解 rm 命令及误操作风险rm 命令基础常见误操作案例防护方案使用 rm编程 别名及安全删除

C++统计函数执行时间的最佳实践

《C++统计函数执行时间的最佳实践》在软件开发过程中,性能分析是优化程序的重要环节,了解函数的执行时间分布对于识别性能瓶颈至关重要,本文将分享一个C++函数执行时间统计工具,希望对大家有所帮助... 目录前言工具特性核心设计1. 数据结构设计2. 单例模式管理器3. RAII自动计时使用方法基本用法高级用法

PHP应用中处理限流和API节流的最佳实践

《PHP应用中处理限流和API节流的最佳实践》限流和API节流对于确保Web应用程序的可靠性、安全性和可扩展性至关重要,本文将详细介绍PHP应用中处理限流和API节流的最佳实践,下面就来和小编一起学习... 目录限流的重要性在 php 中实施限流的最佳实践使用集中式存储进行状态管理(如 Redis)采用滑动

ShardingProxy读写分离之原理、配置与实践过程

《ShardingProxy读写分离之原理、配置与实践过程》ShardingProxy是ApacheShardingSphere的数据库中间件,通过三层架构实现读写分离,解决高并发场景下数据库性能瓶... 目录一、ShardingProxy技术定位与读写分离核心价值1.1 技术定位1.2 读写分离核心价值二

深入浅出Spring中的@Autowired自动注入的工作原理及实践应用

《深入浅出Spring中的@Autowired自动注入的工作原理及实践应用》在Spring框架的学习旅程中,@Autowired无疑是一个高频出现却又让初学者头疼的注解,它看似简单,却蕴含着Sprin... 目录深入浅出Spring中的@Autowired:自动注入的奥秘什么是依赖注入?@Autowired

MySQL分库分表的实践示例

《MySQL分库分表的实践示例》MySQL分库分表适用于数据量大或并发压力高的场景,核心技术包括水平/垂直分片和分库,需应对分布式事务、跨库查询等挑战,通过中间件和解决方案实现,最佳实践为合理策略、备... 目录一、分库分表的触发条件1.1 数据量阈值1.2 并发压力二、分库分表的核心技术模块2.1 水平分

SpringBoot通过main方法启动web项目实践

《SpringBoot通过main方法启动web项目实践》SpringBoot通过SpringApplication.run()启动Web项目,自动推断应用类型,加载初始化器与监听器,配置Spring... 目录1. 启动入口:SpringApplication.run()2. SpringApplicat

Python异步编程之await与asyncio基本用法详解

《Python异步编程之await与asyncio基本用法详解》在Python中,await和asyncio是异步编程的核心工具,用于高效处理I/O密集型任务(如网络请求、文件读写、数据库操作等),接... 目录一、核心概念二、使用场景三、基本用法1. 定义协程2. 运行协程3. 并发执行多个任务四、关键

AOP编程的基本概念与idea编辑器的配合体验过程

《AOP编程的基本概念与idea编辑器的配合体验过程》文章简要介绍了AOP基础概念,包括Before/Around通知、PointCut切入点、Advice通知体、JoinPoint连接点等,说明它们... 目录BeforeAroundAdvise — 通知PointCut — 切入点Acpect — 切面

SpringBoot3匹配Mybatis3的错误与解决方案

《SpringBoot3匹配Mybatis3的错误与解决方案》文章指出SpringBoot3与MyBatis3兼容性问题,因未更新MyBatis-Plus依赖至SpringBoot3专用坐标,导致类冲... 目录SpringBoot3匹配MyBATis3的错误与解决mybatis在SpringBoot3如果