【发布时间】:2022-01-20 22:35:27
【问题描述】:
我尝试使用 Spring Reactor Webflux 学习响应式编程。特别是在以下示例中,我尝试弄清楚如何正确组合处理管道,以使用 ReactorNettyWebSocketClient 消耗 Websocket 连接的入站流。
以下 sn-p 工作正常:
WebSocketClient client = new ReactorNettyWebSocketClient();
client.execute(
URI.create("wss://stream.binance.com:9443/ws/btcusdt@trade"),
session -> {
Flux<String> tradesFlux = session.receive()
.map(WebSocketMessage::getPayloadAsText)
.doOnNext(event -> log.info(event));
return tradesFlux.then();
}
).subscribe();
在上面的 sn-p 中,我在入站流 Flux 上使用 doOnNext() - tradesFlux 来消耗每个传入的事件。它适用于上面的 sn-p,但据我所知,@ 987654324@ 是一个副作用操作,所以我尝试执行以下操作:
WebSocketClient client = new ReactorNettyWebSocketClient();
client.execute(
URI.create("wss://stream.binance.com:9443/ws/btcusdt@trade"),
session -> {
Flux<String> tradesFlux = session.receive()
.map(WebSocketMessage::getPayloadAsText);
tradesFlux.subscribe(new Subscriber<String>() {
@Override
public void onSubscribe(Subscription s) {
}
@Override
public void onNext(String s) {
log.info("replicate binance trade {}", s);
}
@Override
public void onError(Throwable t) {
}
@Override
public void onComplete() {
}
});
return tradesFlux.then();
}
).subscribe();
在第二个 sn-p 中,我尝试通过在入站流 - tradesFlux 上调用 subscribe 来使用带有 Subscriber 的入站流。但是有了这个 sn-p,我得到了以下异常:
reactor.core.Exceptions$ErrorCallbackNotImplemented: java.lang.IllegalStateException: Only one connection receive subscriber allowed.
Caused by: java.lang.IllegalStateException: Only one connection receive subscriber allowed.
不知何故,无法通过他自己的Subscriber 消耗入站流 Flux,所以我想问的问题是:第一个 sn-p 是否已经是使用 sideEffect 消耗ReactorNettyWebSocketClient 入站流的正确方法-op doOnNext() 或者我在这里遗漏了什么?
非常感谢! 董
【问题讨论】:
-
订阅消费,返回的是一次性的,不能再订阅。记录(你正在做的)是一个副作用,所以使用 doOnNext 是正确的选择。
标签: java spring reactive-programming spring-webflux reactor-netty