第13章 开发响应式Web应用

随着计算机软件行业的发展,不仅诞生了各种各样的编程语言,也产生了很多编程范式,从一开始的命令式编程,到后面的面向对象编程及函数式编程,现在,响应式编程也流行起来了。但与前几种范式不同,响应式编程并非不建立在特定的语言基础上,很多语言,比如Java、Kotlin、Scala等都可以进行响应式编程,尤其是在Web开发上的应用变得越来越流行。本章将会带读者了解响应式编程的核心特点,并介绍一个适配Kotlin且原生支持响应式开发的Web框架——Spring 5。最后还将介绍如何动手去实现一个简单响应式Web应用。

13.1 响应式编程的关键:非阻塞异步编程模型

很多人都听说过响应式编程,也使用过这种范式进行过开发,但还是有很多人没有真正理解它背后的思想。下面我们就来看看为什么会诞生响应式编程,以及它到底是为了解决什么问题。

假设一个用户要购物下单,我们需要先获取商品详情和用户的地址,然后根据这些信息进行下单操作。一开始你可能会这么去实现:

data class Goods(val id: Long, val name: String, val stock: Int)
data class Address(val userId: Long, val location: String)
fun getGoodsFromDB(goodsId: Long): Goods {    //获取商品详情
    Thread.sleep(1000)                        //模拟IO操作
    return Goods(goodsId, "深入Kotlin",10)
}

fun getAddressFromDB(userId: Long): Address { //获取地址详情
    Thread.sleep(1000)                               //模拟IO操作
    return  Address(userId, "杭州")
}

fun doOrder(goods: Goods, address: Address): Long {  //进行下单操作
    Thread.sleep(1000)                               //模拟IO操作
    return 1L
}

fun order(goodsId: Long, userId: Long) {
    val goods = getGoodsFromDB(goodsId)
    val address = getAddressFromDB(userId)
    doOrder(goods, address)
}

这是我们通常的做法,很简单也很好理解。但是这种做法有一个缺点:它是一种同步阻塞的方式,在第11章中我们提到了同步阻塞的劣势。这段程序虽然简单,然而更好的方式是将获取商品信息、获取地址这两个没有关联的操作设计成并行执行,这样就可以拥有更快的响应速度。假设我们每次IO操作耗时是100ms,那么上面这段代码的执行时间起码是300ms,其实理论上可以将时间控制在200ms左右。下面我们来看看具体如何去做。

13.1.1 使用CompletableFuture实现异步非阻塞

其实在第11章中也讲过改进这种问题的方案,那就是让这些IO操作能够并行执行且整个过程都是非阻塞的。若要实现异步非阻塞,在Kotlin中有两种方式,一种是利用Java标准库中的CompletableFuture,另一种则是通过协程来实现。这里我们使用CompletableFuture来改进上面的代码:

fun getGoodsFromDB(goodsId: Long): CompletableFuture<Goods> {  
    //返回的是CompletableFuture<Goods>
    return CompletableFuture.supplyAsync {
        Thread.sleep(1000) //模拟IO操作
        Goods(goodsId, "深入Kotlin", 10)
    }
}

fun getAddressFromDB(userId: Long): CompletableFuture<Address> { 
    //返回的是CompletableFuture<Address>
    return CompletableFuture.supplyAsync {
        Thread.sleep(1000) //模拟IO操作
        Address(userId, "杭州")
    }
}

fun doOrder(goods: Goods, address: Address): CompletableFuture<Long> {
    return CompletableFuture.supplyAsync {
        Thread.sleep(1000) //模拟IO操作
        1L
    }
}

fun main(args: Array<String>) {
    val goodsF = getGoodsFromDB(1)
    val addressF = getAddressFromDB(1)
    CompletableFuture.allOf(goodsF, addressF).thenApply { //保证前两个IO操作返回结果后再执行
        Stream.of(goodsF, addressF).map { it.join() }.collect(Collectors.toList<Any>()) //这里需要借助Stream来获取结果
    }.thenApply { 
        doOrder(it[0] as Goods, it[1] as Address) 
    }.join()
}

在Java 8之后我们确实可以使用CompletableFuture来写异步非阻塞代码,但是我们发现对CompletableFuture的操作却不怎么直观。比如上面的合并操作,还需要借助Stream来得到结果,这让我们开发变得烦琐,不容易理解。

13.1.2 使用RxKotlin进行响应式编程

一种更直观的实现异步非阻塞程序的解决方案是利用RxJava,它同样适用于SE 8之前的Java版本。你很可能知道Rx系列的类库,它的一个主要作用就是提供统一的接口来帮助我们更方便地处理异步数据流。其中RxJava提供了对Java的支持,但这里我们将会使用RxKotlin,它的实现基于RxJava,但是增加了一些Kotlin独有的特性。下面我们就利用RxKotlin来实现一个新的版本:

val threadCount = Runtime.getRuntime().availableProcessors()
val threadPoolExecutor = Executors.newFixedThreadPool(threadCount) //线程池
val scheduler = Schedulers.from(threadPoolExecutor)                //调度器

fun getGoodsFromDB(goodsId: Long): Observable<Goods> { //返回的是一个Observable<T>类型数据
    return Observable.defer {
        Thread.sleep(1000)                                        //模拟IO操作
        Observable.just(Goods(goodsId, "深入Kotlin", 10))
    }
} 
fun getAddressFromDB(userId: Long): Observable<Address> {
    return Observable.defer {
        Thread.sleep(1000)                                        //模拟IO操作
        Observable.just(Address(userId, "杭州"))
    }
}

fun rxOrder(goodsId: Long, userId: Long) {
    var goods: Goods? = null
    var address: Address? = null
    val goodsF = getGoodsFromDB(1).subscribeOn(scheduler)  //方法再指定调度器中执行
    val addressF = getAddressFromDB(1).subscribeOn(scheduler) //方法再指定调度器中执行
    
    Observable.merge(goodsF, addressF).subscribeBy( //合并两个Observable
        onNext = { when(it) {
            is Goods -> goods = it
            is Address -> address = it
        } }, 
        onComplete = {  //全部执行后
            doOrder(goods!!, address!!)  
        }
    )
}

看上去以上程序代码比较多,但使用RxKotlin带来了以下优势:

• 将异步编程变得优雅、直观,不用对每个异步请求都执行一个回调,同时还可以组合多个异步任务;

• 不需要书写多线程代码,只需指定相应的策略便可使用多线程的功能;

• Java 6及以上的版本都可用。

这让我们在实现需求的同时,又保持了代码的简洁和优雅。响应式编程除了异步编程模型这个特点外,还有另一个特点,那就是数据流处理。简单来说就是将数据处理的过程变得像流水线一样,比如A=>B=>C=>……,后继者不需要阻塞等待结果,而是由前一个处理者将结果通知它。

举个例子,假设下单之后需要给商家及消费者推送消息,那么用流处理如下表示:

doOrder()
    .map(doNotifyCustomer)
    .map(doNotifyShop)
    .map(doOther)
    ...

我们可以对原始数据进行处理,生成一个新的数据然后传递给下一个处理者。这些处理过程都是异步非阻塞的。可以看出,流式调用相对于回调的方式实在是优雅得太多了,不再需要编写大量的嵌套回调函数,从而使代码更加简洁易懂。

13.1.3 响应式Web编程框架

通过上面的这些例子可以发现,应用响应式编程拥有诸多优势。应用一些第三方响应式的类库则能帮助我们更快地进行响应式程序开发,而不必关心异步、线程等细节,只着重于业务逻辑的处理。

既然响应式编程有这么多优点,那么为什么在以前的Java Web生态中应用得却不那么广泛呢?原因有以下几点:

1)传统的Servlet容器,比如Tomcat是同步阻塞的模型(Servlet 3.1 Async IO之前)​;

2)主流的Java Web框架对响应式的支持不是很好,比如Spring MVC、Spring Boot等,当然也有支持响应式编程的Web框架,比如Vertx、Play!Framework;

3)一些主流的第三方类库的实现是同步阻塞的,比如连接MySQL的驱动包,所以很难使整个系统真正做到异步非阻塞。

这些原因使得响应式编程在Java Web领域使用得不是很多。当然,如果用Play!进行Web开发的话,你自然就会进行响应式编程,因为它就是全面支持响应式编程。幸运的是,随着Spring 5的发布,这种局面将会被打破,因为Spring 5开始全面拥抱响应式编程,而且适配Kotlin。下面我们就来看看Spring 5到底给我们带来了些什么。

13.2 Spring 5:响应式Web框架

Spring的大名每个Java程序员都应该听过,但Spring 5或许很多人并不是很了解,它是2017年9月才发布的,引入的一些崭新的特性,带来的不仅仅是技术上的改变,更多的是开发思维上的变化。下面我们就来看看Spring 5的这些新特性。

13.2.1 支持响应式编程

前面一节我们已经说了响应式编程的好处,但在Spring 5版本以前它并不是原生支持响应式编程的,主要原因是底层Web容器的限制。Tomcat等容器在Servlet 3.1支持Async IO之前,并不能做到真正的异步非阻塞,而集成一些支持异步非阻塞的容器,比如Netty,又相对比较复杂。然而,在Spring 5发布后,你可以轻松选择自己所需的Web容器,比如Tomcat或者Netty等,这给Spring支持响应式编程提供了底层基础。而传统的Spring MVC并不原生支持响应式编程,所以Spring 5引入了一个全新的Web框架,那就是Spring Webflux。

Spring Webflux主要帮助我们在框架层面实现响应式编程,它不再使用传统基于Servlet实现的HttpServletRequest和HttpServletResponse,而是采用全新的ServerRequest和ServerResponse。同时Spring Webflux还要求请求的返回数据类型为Flux,这是一种响应式的数据流类型,比如我们在上一节中提到的Observable类型。

不过需要注意的一点是,Spring 5并没有使用RxJava 2作为程序的响应式类库,默认集成的是Reactor库。那么这是基于什么考量呢?其实了解响应式编程的读者应该对两个库都比较熟悉,我们不能说谁好谁坏。但RxJava库早于Reactor库诞生,所以RxJava一开始是处于响应式编程的探索阶段,当时Java并没有提出相应的响应式编程规范,所以RxJava 2受限于兼容RxJava遗留的历史包袱,有些方面使用起来并不是很方便。而Reactor则完全是基于响应式流规范设计和实现的类库,同时JDK的最低版本是JDK 8,所以可以使用JDK8提供的流操作。如果你想写更加简洁、更加函数式的代码,Reactor或许是个更好的选择。下面我们就来看一下Spring Webflux中最基础的两个数据类型。

1.Mono

在传统的Spring MVC里,请求的返回直接是一个对象,比如查询一个用户,返回的是一个User对象或者一个null。而在Spring WebFlux则是使用了Mono,它代表的是0~1个元素,比如它的返回类型为Mono<User>,代表返回流中只有一个数据或者为空数据。

2.Flux

在业务开发中,我们除了返回一个简单的对象外,有时还会返回集合对象,比如查询一批用户,那么返回值为List<User>。而在Spring WebFlux中则使用了Flux,它代表的是0~N个元素,比如它的返回类型为Flux<User>,代表返回流中有0~N个数据。Spring 5除了引入Spring Webflux来提供响应式编程特性以外,它还有另一个特性让Kotlin开发者非常兴奋,那就是适配Kotlin。

13.2.2 适配Kotlin

为什么这里我们会讲Spring 5而不是其他的支持响应式的Web框架呢?除了它受众面比较广以外,另一个原因是它全面适配Kotlin。Kotlin虽然一直在安卓开发中被广泛采用,但在Web开发中却少见其身影,一个很重要的原因就是没有一个好的Web框架适配它。虽然在Spring 5之前已经有了Ktor、Javalin等框架支持Kotlin,但由于相对比较小众,并没有被广泛应用。...

13.2.3 函数式路由

路由配置是一个Web框架的特色,Spring从最早的XML配置到后来的注解配置,现在也支持了函数式路由。当然,现在大多数人还是使用注解来作为Spring的路由配置。我们不去探讨注解配置好还是函数式路由配置好,而是着重介绍一下函数式路由能实现以前用注解无法实现的功能。

我们知道,用注解来配置路由虽然很简单,也很直接,但是随着微服务及模块化程序开发趋势的发展,路由分模块化统一管理是一个需求,但用传统的注解方式却很难做到。而Spring 5最新支持的函数式路由却可以实现这个功能,而且结合Kotlin DSL,语法也非常简洁。我们来看一个简单的例子。

假设现在有2个handler,分别是UserHandler以及CustomerHandler,里面都有3个方法。若是用注解的方式,我们会这么做:

@Component
class UserHandler {

    @RequestMapping(value = "user/getUser", method = [RequestMethod.GET], produces = [MediaType.APPLICATION_JSON_VALUE])   //注解配置路由
    fun getUser() {}

    @RequestMapping(value = "user/addUser", method = [RequestMethod.POST], produces = [MediaType.APPLICATION_JSON_VALUE])
    fun addUser() {}

    @RequestMapping(value = "user/updateUser", method = [RequestMethod.PUT], produces = [MediaType.APPLICATION_JSON_VALUE])
    fun updateUser() {}
    
}

@Component
class CustomerHandler {

    @RequestMapping(value = "customer/getCustomer", method = [RequestMethod.GET], produces = [MediaType.APPLICATION_JSON_VALUE])
    fun getCustomer() {}

    @RequestMapping(value = "customer/addCustomer", method = [RequestMethod.POST], produces = [MediaType.APPLICATION_JSON_VALUE])
    fun addCustomer() {}

    @RequestMapping(value = "customer/updateCustomer", method = [RequestMethod.PUT], produces = [MediaType.APPLICATION_JSON_VALUE])
    fun updateCustomer() {}
    
}

我们发现这种方式虽然直接方便,但是如果Handler里面方法一多,路由信息与方法掺杂在一起,会导致整个类变得臃肿,不易维护。所以我们希望有一种方式既能保持声明路由的简洁性和功能性,比如支持REST请求、指定请求及返回的数据类型等,同时又方便统一管理。Spring 5中的函数式路由能帮我们解决这个问题。下面我们就来看一下改造后的代码:

import org.springframework.http.MediaType

@Component
class UserHandler {  //类中没有路由信息
    fun getUser() {}
    fun addUser() {}
    fun updateUser() {}   
}

@Component
class CustomerHandler {
    fun getCustomer() {}
    fun addCustomer() {}
    fun updateCustomer() {}
}

@Configuration
class Routes(userHandler: UserHandler, customerHandler: CustomerHandler) { 
    //定义路由类统一管理
    @Bean
    fun userRouter() = router {     //不同类的路由分开管理
        "user".nest {
            GET("/getUser").nest {  //支持REST请求
                accept(APPLICATION_JSON, userHandler::getUser)
            }
            POST("/addUser").nest {
                accept(APPLICATION_JSON, userHandler::addUser)
            }
            PUT("/updateUser").nest {
                accept(APPLICATION_JSON, userHandler::updateUser)
            }
        }
    }
    @Bean
    fun customerRouter() = router {
        "customer".nest {
            GET("/getCustomer").nest {
                accept(APPLICATION_JSON, userHandler::getCustomer)
            }
            POST("/addCustomer").nest {
                accept(APPLICATION_JSON, userHandler::addCustomer)
            }
            PUT("/updateCustomer").nest {
                accept(APPLICATION_JSON, userHandler::updateCustomer)
            }
        }
    }
}

乍一看这种方式似乎并没有简单多少,甚至感觉代码更多了。但仔细思考一下,其实这是一个更合理的方式,它帮助我们将配置与业务逻辑分离,而且统一管理,功能点上也没有很大的缺失。同时这种方式也更符合函数式编程的风格,结合Kotlin DSL使代码变得更加精简优雅,可读性也更好。

13.2.4 异步数据库驱动

如果一个请求在执行过程中有一部分是同步阻塞的,那么整个应用就不能算异步非阻塞。而我们知道,在实际的业务场景中与数据库打交道是无法避免的,也就是说如果想要实现整个系统保证异步非阻塞的架构,那么数据库操作也必须是异步非阻塞的,程序与数据库通信的驱动需要支持异步非阻塞,比如现在Spring支持的MongoDB、Redis等。但我们在很多场景用的是MySQL,由于我们使用的JDBC驱动是同步阻塞的,所以我们将无法达到全异步非阻塞的架构。那么如果我们需要使用MySQL,并且还要保证整个系统是异步非阻塞的架构,就需要一个支持异步非阻塞操作的数据库驱动。

其实在Scala上已经有了这么一个驱动:postgresql-async(项目地址:https://github.com/mauricio/postgresql-async)​,全异步,基于Netty实现,同时支持MySQL和Postgresql。一些开源项目和公司也已经在实际中使用它了,比如Quill(官网地址:http://getquill.io/)​,该项目的github地址:https://github.com/mauricio/postgresql-async。有兴趣的读者可以去看看。但是不幸的是,这个项目的作者声明已经不再维护了。

因为受限于这个项目实现使用了很多Scala才有的数据类型,比如Future(与Java中的Future不一样)​,所以我们无法在Java以及Kotlin的环境中使用它。但幸运的是有个Kotlin的社区人员将这个项目用Kotlin重写了一遍,项目叫作jasync-sql,基于Java8的CompletableFuture,完全适配Java及Kotlin,与Spring 5最新的webflux也可以结合得很好。当然这只是一个小众项目,没有经历过大量的测试以及实践的考验,仅供学习,不推荐大家在一些大型项目中使用。这个项目的github地址:https://github.com/jasync-sql/jasync-sql。

这里给大家简单演示一下,如何使用jasync-sql与Spring Webflux相结合:

//创建一个数据库连接
Connection connection = new MySQLConnection(
    new Configuration(
        "root",
        "localhost",
        3306,
        "123456",
        "test"
    )
);
//执行连接
CompletableFuture<Connection> connectFuture = connection.connect()

//执行数据库操作
CompletableFuture<QueryResult> queryResult = connection.sendPreparedStatement ("select * from user");

val result: Mono<QueryResult> = Mono.fromFuture(queryResult)

其实书写方式跟我们以前用JDBC写数据库操作很类似,也是先创建连接,然后执行数据库操作。不同的在于返回的数据类型,不是简单的QueryResult,而是一个CompletableFuture<T>类型。它不同于Future,不仅仅是异步执行的,而且获取值的时候也是非阻塞的,同时CompletableFuture<T>类型的值转化为Webflux所要求的数据类型很容易,比如使用Mono.fromFuture就可以将一个CompletableFuture类型的值转换为Mono或者Flux,所以,我们可以使用jasync-sql来构建基于MySQL和Spring Webflux响应式应用。

13.3 Spring 5响应式编程实战

本节将会带大家基于Spring WebFlux+Kotlin+MySQL,实现一个简化的股票行情实时推送功能。本示例需要对Spring以及gradle有基本的了解。

实现这种需求有很多方式,比如Ajax轮询、长轮询、WebSocket等。但这个例子中我们将使用另外一种方式,那就是Server Sent Event。它虽然不能像WebSocket一样实现双工通信,只能由服务器不断地向客户端发送消息,但它也有自己的优势,比如基于Http协议,会自动断开重连等。所以这里我们就通过这种方式来模拟实现股票行情的实时推送,因为查看股票行情实时行情,往往只需要服务端向客户端不断推送消息即可。

这里我们使用gradle来构建我们的项目,同时我们将加入13.2节所讲的MySQL异步数据库驱动,来保证我们在使用MySQL作为存储DB进行数据库操时也是异步非阻塞的。最终这个项目的目录结构大致如下:

|____main
| |____kotlin
| | |____Application.kt
| | |____handler
| | | |____StockHandler.kt
| | |____JasyncPool.kt
| | |____models
| | | |____StockQuotation.kt
| | |____Routes.kt
| | |____service
| | | |____StockService.kt
| |____resources
| | |____application.properties
| | |____static
| | | |____index.js
| | |____templates
| | | |____index.mustache

build.gradle.kts配置文件:

import org.jetbrains.kotlin.gradle.tasks.KotlinCompile
import org.jetbrains.kotlin.gradle.plugin.KotlinPluginWrapper

group = "spring-kotlin-jasync-sql"
version = "1.0-SNAPSHOT"

val springBootVersion: String by extra

plugins {
    application
    kotlin("jvm") version "1.2.70"
    kotlin("plugin.spring").version("1.2.70")
    id("org.springframework.boot").version("1.5.9.RELEASE")
    id("io.spring.dependency-management") version "1.0.5.RELEASE"
}


buildscript {

    val springBootVersion: String by extra { "2.0.0.M7" }

    dependencies {
        classpath("org.springframework.boot:spring-boot-gradle-plugin:$springBoot Version")
    }

    repositories {
        mavenCentral()
        jcenter()
        maven {
            url = uri("https://repo.spring.io/milestone")
        }
    }
}

extra["kotlin.version"] = plugins.getPlugin(KotlinPluginWrapper::class.java).kotlin PluginVersion

repositories {
    mavenCentral()
    jcenter()
    maven {
        url = uri("https://repo.spring.io/snapshot")
    }
    maven {
        url = uri("https://repo.spring.io/milestone")
    }
}

dependencies {
    compile(kotlin("stdlib-jdk8"))
    compile("org.jetbrains.kotlin:kotlin-reflect")
    compile("org.jetbrains.kotlin:kotlin-stdlib-jdk8")
    compile("com.github.jasync-sql:jasync-mysql:0.8.32")  //添加jasync-mysql依赖
    compile("com.samskivert:jmustache")
    compile("org.springframework.boot:spring-boot-starter-actuator:$springBoot Version")
    compile("org.springframework.boot:spring-boot-starter-webflux:$springBoot Version")
    compile("org.springframework.boot:spring-boot-starter-thymeleaf:$springBoot Version")
}

configure<JavaPluginConvention> {
    sourceCompatibility = JavaVersion.VERSION_1_8
}
tasks.withType<KotlinCompile> {
    kotlinOptions.jvmTarget = "1.8"  //指定编译版本
}

接下来,我们先来定义一下model,这里我们使用data class:

data class StockQuotation(
    val id: Long,
    val stock_id: Long, //股票代码
    val stock_name: String, //股票名称
    val price: Int,  //股票价格
    val time: String //当前时间
)

data class StockQuotationResult(
    val queryTime: String,  //时间
    val stockQuotation: StockQuotation //当前股票信息
)

因为我们这里需要使用Jasync-sql,所以我们需要配置相应的数据库连接池:

@Component
class DB {
    private val configuration = Configuration(  //数据库配置
            "test",
            "localhost",
            3306,
            "123456",
            "test")
    private val poolConfiguration = PoolConfiguration( //连接池配置
            maxObjects = 100,
            maxIdle = TimeUnit.MINUTES.toMillis(15),
            maxQueueSize = 10_000,
            validationInterval = TimeUnit.SECONDS.toMillis(30)
    )
    val connectionPool = ConnectionPool(factory = MySQLConnectionFactory(configuration), configuration = poolConfiguration) //创建数据库连接池
}

接下来,是整个项目中比较重要的一部分,那就是如何以响应式编程的方式获取我们所需的数据:

@Component
class StockService(val db: DB) {

    val repeat = Flux.interval(Duration.ofMillis(1000)) //(1)

    fun getStockQuotation(): Flux<StockQuotationResult> {
        val query = "select * from stock_quotation order by id desc limit 1;"
        fun stockQuotation(time: DateTime) = Mono.fromFuture(db.connectionPool.sendPreparedStatement(query)).map {//(2)
            StockQuotationResult(time.toString("YYYY-MM-dd hh:mm:ss"), transRowData ToStockQuotation(it.rows.orEmpty().first())) //(3)
        } //(5)
        return repeat.flatMap {
            insertStockQuotation()
            stockQuotation(DateTime.now()) }
    }

    private fun transRowDataToStockQuotation(rowData: RowData): StockQuotation { //(4)
        return StockQuotation(
                rowData.get("id").toString().toLong(),
                rowData.get("stock_id").toString().toLong(),
                rowData.get("stock_name").toString(),
                rowData.get("price").toString().toInt(),
                (rowData.get("time") as LocalDateTime).toString("YYYY-MM-dd hh:mm:ss"))
    }

    private fun insertStockQuotation() {  (6)
        val max = 74000
        val min = 72000
        val price = Random().nextInt(max - min) + min
        val query = "insert into stock_quotation (stock_id, stock_name, price, time) values (600519, '贵州茅台', ${price}, '${DateTime.now().toString("YYYY-MM-dd hh:mm:ss")}')"
        db.connectionPool.sendPreparedStatement(query)
    }
}

我们一步一步来看着代码:

1)我们创建了一个定时循环的Flux,用来控制模拟定时循环查询数据库;

2)利用创建的数据库连接池从数据库中查询数据,返回的数据类型是:Completable-Future<QueryResult>;

3)因为我们只需要第1列数据,所以这里使用it.rows.orEmpty().first()来获取第1列数据;

4)将RowData类型数据转化为我们定义的data class对象,这里我们没有使用第三方ORM框架,需要自己手动转换;

5)使用Mono.fromFuture将一个CompletableFuture<T>类型数据转换为Mono<T>类型;

6)模拟定时生成股票价格。

以上是这段代码的一个大致逻辑,其实重点是第2步和第5步。我们知道,CompletableFuture相比Future一个很大的优势就是它获取值的时候不必阻塞等待,这便保证了整个查询过程是异步非阻塞的。同时Reactor提供了将CompletableFuture转化为Mono的方法,这样就可以完全适配Spring Webflux所要求的返回数据的类型格式。当然,前面我们说过,这个例子我们使用的是Server Sent Event,所以我们需要指定返回数据的形式:

ok().bodyToServerSentEvents(stockService.getStockQuotation())

在写完业务逻辑后,还有一块比较重要的就是Router,这里我们将会使用Spring 5最新的函数式Router。当然你也可以使用传统的、基于Spring MVC的注解方式。所以最终我们的Router如下:

@Configuration
class Routes(val userHandler: StockHandler) {
    @Bean
    fun Router() = router {
        accept(MediaType.TEXT_HTML).nest {
            GET("/") { ok().render("index") }
        }
        "/api".nest {  //api开头的请求
            GET("/getStockQuotation").nest {
                accept(TEXT_EVENT_STREAM, userHandler::getStockQuotation)
            }
        }
        resources("/**", ClassPathResource("static/"))  //静态文件访问路径
    }.filter { request, next ->
        next.handle(request).flatMap {
            if (it is RenderingResponse) RenderingResponse.from(it).build() else it.toMono()
        }
    }
}

用这种方式定义router相对传统方式来说,更加语义化也更容易管理,而且还支持对Response进行不同处理。

至此,整个项目的后台开发已经完成,下面我们来看一下前端部分应该怎么做。

虽然目前很多项目多是采用前后端分离的架构,但是这里为了更方便演示示例,以便让大家更容易搭建这个项目,这里前端页面渲染采用了Mustache模板引擎。另外,需要注意的一点是浏览器是否支持Server Sent Event这种传输格式,当前IE及Edge的所有版本都不支持,所以要测试这个例子最好使用其他浏览器。最终我们的前端代码包括以下两个部分。

模板如下:

<!DOCTYPE html>
<html lang="en">
<head>
    <meta charset="UTF-8">
    <title>股票行情</title>
    <script src="index.js"></script>
    <style>
    #stockQuotations{
        margin: 0 auto;
        text-align: center;
    }
    #stockQuotations li {
        list-style-type:none;
        margin-bottom: 5px;
    }
    </style>
</head>
<body>

<div id="stockQuotations">
</div>
</body>
</html>

js文件如下:

var eventSource = new EventSource("/api/getStockQuotation");
eventSource.onmessage = function(e) {
    var li = document.createElement("li");
    var data = JSON.parse(e.data);
    li.innerText = "股票代码: " + data.stockQuotation.stock_id + " 股票名称:" + data.stockQuotation.stock_name + " 当前价格: " + (data.stockQuotation.price / 100.0).toFixed(2);
    document.getElementById("stockQuotations").appendChild(li);
}

这里我们需要使用EventSource,而不是我们常见的Ajax方式请求,同时用eventSource.onmessage来监听返回的数据来进行处理。

最后我们运行程序。通过浏览器打开页面:http://localhost:8282/,我们会看到图13-1所示界面。 

源码已经上传到github上面,有兴趣的读者可以去看看:https://github.com/godpan/reactive-spring-kotlin-app。

本节主要是带大家将之前所讲的知识点进行一个总结串联。自己动手去写一个Demo,能帮助大家更好地理解相关知识点,加深印象。

13.4 本章小结

(1)响应式编程

了解什么是响应式编程的关键,响应式编程相对于传统编程范式的优势,同时如何利用一些第三方类库来帮助我们在程序中进行响应式开发。

(2)Spring 5支持响应式编程

简单了解Spring 5支持响应式编程的背景,同时介绍了它的一些新特性,比如函数式路由以及适配Kotlin等。

(3)异步非阻塞MySQL数据库驱动

介绍了一个基于Netty且用Kotlin实现的全异步非阻塞的MySQL数据库驱动:Jasync-sql,以及如何在Spring Webflux中使用它。

(4)Spring Webflux+Kotlin示例

了解如何Kotlin如何使用Spring Webflux来进行响应式Web应用开发。

更多推荐