【问题标题】:What is proper way to consume inbound stream of ReactorNettyWebSocketClient使用 ReactorNettyWebSocketClient 的入站流的正确方法是什么
【发布时间】: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


【解决方案1】:

我发现这个 Websocket 连接不能被订阅两次,就像异常说的那样。对我来说,解决方案是使用map 操作——实际上是在入站流 Flux 上的链,以一对一的方式处理事件。之后上面的问题听起来很愚蠢,但我发现所有关于使用ReactorNettyWebSocketClient 的示例总是使用doOnNext() 来处理来自入站流的事件。

运行代码如下所示:

@Bean
public WebSocketClient binanceTradesWS(Function<String, BinanceTrade> binanceTradeMapper,
                                       Function<BinanceTrade, Trade> binanceTradeReplicator) {
    WebSocketClient client = new ReactorNettyWebSocketClient();
    client.execute(
            URI.create("wss://stream.binance.com:9443/ws/btcusdt@trade"),
            session -> {
                return session.receive()
                        .map(WebSocketMessage::getPayloadAsText)
                        .map(binanceTradeMapper)
                        .map(binanceTradeReplicator)
                        .then();
            }
    ).subscribe();

    return client;
}

@Bean
public ObjectMapper objectMapper() {
    ObjectMapper mapper = new ObjectMapper().registerModule(new KotlinModule());
    mapper.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
    return mapper;
}

@Bean
public Function<String, BinanceTrade> binanceTradeMapper(ObjectMapper objectMapper) {
    return (String tradeJSON) -> {
        log.info("map binance trade {}", tradeJSON);
        try {
            return objectMapper.readValue(tradeJSON, BinanceTrade.class);
        } catch (JsonProcessingException e) {
            e.printStackTrace();
            return null;
        }
    };
}

@Bean
public Function<BinanceTrade, Trade> binanceTradeReplicator() {
    return (BinanceTrade binanceTrade) -> {
        Trade trade =new Trade(
                binanceTrade.getT(),
                binanceTrade.getS(),
                new BigDecimal(binanceTrade.getQ()),
                new BigDecimal(binanceTrade.getP()),
                LocalDateTime.ofInstant(Instant.ofEpochMilli(binanceTrade.getTime()), TimeZone.getTimeZone("GMT").toZoneId())
        );
        log.info("replicated binance trade {}", trade);
        return trade;
    };
}

更新: 在深入研究材料之后,我觉得处理ReactorNettyWebSocketClient#execute 的 Websocket 连接的方式与它的反应性生态系统不兼容。如果在带有 Kafka-Binder 的 Spring Cloud Stream 上,您无法将其入站流 Flux 转发到模块 Reactor Kafka 的反应性 KafkaSender 或实现此 Flux 的 Supplier 函数。因为这两种机制都会在 Flux 上调用 subcribe 并且会抛出上面提到的异常。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-17
    • 2018-02-04
    • 2016-10-27
    • 2019-09-04
    • 2016-01-11
    • 2019-09-14
    相关资源
    最近更新 更多