【问题标题】:Cannot use 'subscribe' or 'subscribeWith' with 'ReactorNettyWebSocketClient' in Kotlin不能在 Kotlin 中将 'subscribe' 或 'subscribeWith' 与 'ReactorNettyWebSocketClient' 一起使用
【发布时间】: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 应该是什么是。我看不出后者是如何编译的。

【问题讨论】:

标签: websocket kotlin spring-webflux project-reactor


【解决方案1】:

subscribe 和 subscribeWith 应始终在运算符链的末尾使用,而不是作为中间运算符。

【讨论】:

  • 这就是我的意思。我想指定将对 WebSocketMessage 的有效负载执行的最终操作。到目前为止,在我对 Rx 的所有使用中,'subscribe' 或 'subscribeWith' 将是合适的选择,但这里似乎必须是 'doOnNext'。因此我在问为什么我想做的最后一件事不符合运营商链的末端?
  • 啊,明白了。在 reactor-netty API 中,您的最终处理实际上是包含更多服务器生命周期的更大链的一部分。这就是为什么你必须返回一个Flux(或Publisher),因此不能使用subscribe。在框架订阅的 Spring WebFlux 中也是如此。
  • 我认为它必须是这样的。但这似乎明显违反了“最少意外原则”——考虑到它是由 API 传递给我的一个流,并且我有一个我想要执行的终端操作。至少我希望能找到更清晰的文档。
【解决方案2】:

Simon 已经提供了答案,但我会添加一些额外的上下文。

当使用 Reactor(和 ReactiveX 模式)组合异步逻辑时,您构建了一个端到端的处理步骤链,其中不仅包括 WebSocketHandler 本身的逻辑,还包括负责负责的底层 WebSocket 框架代码的逻辑向套接字发送和从套接字接收消息。将整个链连接在一起非常重要,这样在运行时“信号”将从头到尾流过它(onNext、onError 或 onComplete)并传达最终结果,即.block() 在结束。

在 WebSocket 的情况下,这看起来有点令人生畏,因为您实际上是将两个或多个流合并为一个。您不能只订阅其中一个流(例如,用于入站消息),因为这会阻止组成统一的处理流,并且信号不会流到预期最终结果的末尾。

另一方面是subscribe() 触发流上的消费,而您真正想要的是在延迟模式下继续编写异步逻辑,即声明数据实现时将发生的所有事情。这是组成单个统一链很重要的另一个原因。所以完全声明后可以触发。

简而言之,与 Servlet 世界的命令式 WebSocketHandler 的主要区别在于,它不是单个消息的处理程序,而是组合完整流的处理程序。在这里,单个消息的处理只是整个处理链的一个步骤。所以订阅的唯一位置是在最后,.block() 所在的位置,以便开始处理。

顺便说一句,自从几个月前首次发布此问题以来,文档已改进为 provide more guidance,关于如何实现 WebSocketHandler。

【讨论】:

    猜你喜欢
    • 2018-02-01
    • 1970-01-01
    • 1970-01-01
    • 2014-01-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-09
    • 1970-01-01
    相关资源
    最近更新 更多