【发布时间】:2018-09-10 10:24:25
【问题描述】:
下面的 Kotlin 代码成功连接到 Spring WebFlux 服务器,发送消息并打印通过返回的流发送的每条消息。
fun main(args: Array<String>) {
val uri = URI("ws://localhost:8080/myservice")
val client = ReactorNettyWebSocketClient()
val input = Flux.just(readMsg())
client.execute(uri) { session ->
session.send(input.map(session::textMessage))
.thenMany(
session.receive()
.map(WebSocketMessage::getPayloadAsText)
.doOnNext(::println) // want to replace this call
.then()
).then()
}.block()
}
在之前的响应式编程经验中,我总是使用 subscribe 或 subscribeWith 来调用 doOnNext。但是,在这种情况下它不起作用。我知道这是因为两者都没有返回正在使用的反应流 - subscribe 返回一个 Disposable 并且 subscribeWith 返回 Subscriber 它作为参数接收。
我的问题是调用 doOnNext 是否真的是添加处理程序来处理传入消息的正确方法? 大多数 Spring 5 教程显示调用 this 或 log 的代码,但有些使用 subscribeWith(output).then() 而不指定 output 应该是什么是。我看不出后者是如何编译的。
【问题讨论】:
-
你有这样的 subscribeWith 教程的例子吗?
标签: websocket kotlin spring-webflux project-reactor