【问题标题】:How to wait for variables to change before sending the next content , websocket by reactor-netty如何在发送下一个内容之前等待变量更改,websocket by reactor-netty
【发布时间】:2019-08-10 02:49:05
【问题描述】:

怎么办? ReadFlushMessage() 会等待这个。消息以获取新内容。

没有立即退出 websocket 连接

    private Set<String> message = new HashSet<>();

    private void writeMessage(String message) {
        this.message.add(message);
    }

    private String[] readFlushMessage() {
        String[] _message = (String[])this.message.toArray();
        this.message = new HashSet<>();
        return _message;
    }

    private Publisher<Void> websocketPublisherA(HttpServerRequest request, HttpServerResponse response, WebSocketServerHandle handleObject) {
        return response
            .header("content-type", "text/plain")
            .sendWebsocket((in, out) ->
                out.options(NettyPipeline.SendOptions::flushOnEach)
                    .sendString(
                        Flux.just(readFlushMessage())
                    )
            );
    }

【问题讨论】:

    标签: websocket reactor-netty


    【解决方案1】:

    尝试使用Reactor Core 类型。

    这是一个非常简单的示例,说明您想要实现的目标:

        @Test
        public void test() {
            FluxProcessor<String, String> serverMsg =
                    ReplayProcessor.<String>create();
    
            Flux.range(1, 20)
                    .map(Object::toString)
                    .subscribe(serverMsg::onNext);
    
            DisposableServer server =
                    HttpServer.create()
                            .port(0)
                            .handle((req, resp) ->
                                    resp.header("content-type", "text/plain")
                                        .sendWebsocket((in, out) ->
                                                out.options(NettyPipeline.SendOptions::flushOnEach)
                                                        .sendString(serverMsg)
                                        ))
                            .bindNow();
    
            HttpClient.create()
                    .port(server.port())
                    .websocket()
                    .receive()
                    .asString()
                    .doOnNext(System.out::println)
                    .blockLast();
    
            server.disposeNow();
        }
    

    【讨论】:

    • 像魅力一样工作!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-11-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-11-12
    • 2020-06-15
    相关资源
    最近更新 更多