【问题标题】:Spring Boot Reactive Websoket - Block out flux until received all information from clientSpring Boot Reactive Websocket - 中断通量,直到收到来自客户端的所有信息
【发布时间】:2018-11-18 08:16:54
【问题描述】:

我正在使用响应式 websocket(在 spring boot 2.1.0 上),但我在尝试阻止等待客户端信息的输出通量时遇到问题。

我知道阻塞不是处理反应流的正确方法,但我需要在继续之前从客户端接收一些信息(即:身份验证密钥、ID),有一种可接受的方法来管理它反应式?

例如:我希望客户端发送授权密钥和订阅 ID(仅订阅特定事件),并且仅当我拥有这两个信息时才发送输出通量,或者如果信息不存在则关闭会话有效

我已经尝试在 handle 方法中管理检查

     webSocketSession.receive().subscribe(inMsg -> {

            if(!inMsg.getPayloadAsText().equals("test-key")) {
                log.info("AUTHORIZATION ERROR");
                webSocketSession.close();
            }
        });

但是这种方式不起作用并且不正确,因为它以“异步”方式管理会话终止,并且无论如何在接收到带有错误密钥的消息时会话仍然保持活动状态

另一种方法是使用内存存储来保持“会话”,以跟踪接收到的信息并在业务逻辑级别进行处理

我一直在寻找一种以被动方式管理它的正确方法

我的出发点是这个例子:http://www.baeldung.com/spring-5-reactive-websockets

提前致谢

其他信息:

使用示例https://github.com/eugenp/tutorials/tree/master/spring-5-reactive/src/main/java/com/baeldung/reactive/websocket

我创建了一个基本的 Spring Boot 应用程序,并为 websocket 添加了一个配置类:

@Configuration 
class WebConfig {

    @Bean
    public HandlerMapping handlerMapping() {
        Map<String, WebSocketHandler> map = new HashMap<>();
        map.put("/event-emitter-test", new MyWebSocketHandler());

        SimpleUrlHandlerMapping mapping = new SimpleUrlHandlerMapping();
        mapping.setUrlMap(map);
        mapping.setOrder(-1); // before annotated controllers
        return mapping;
    }

    @Bean
    public WebSocketHandlerAdapter handlerAdapter() {
        return new WebSocketHandlerAdapter();
    }
}

以下是包含主 websocket 方法的主类(用我的实际代码修改):

@Component
public class MyWebSocketHandler  implements WebSocketHandler {

    @Autowired
    private WebSocketHandler webSocketHandler;

    @Bean
    public HandlerMapping webSocketHandlerMapping() {
        Map<String, WebSocketHandler> map = new HashMap<>();
        map.put("/event-emitter-test", webSocketHandler);

        SimpleUrlHandlerMapping handlerMapping = new SimpleUrlHandlerMapping();
        handlerMapping.setOrder(1);
        handlerMapping.setUrlMap(map);
        return handlerMapping;
    }

    @Override
    public Mono<Void> handle(WebSocketSession webSocketSession) {

       webSocketSession.receive().subscribe(inMsg -> {

        if(!inMsg.getPayloadAsText().equals("test-key")) {
            // log.info("AUTHORIZATION ERROR");
            webSocketSession.close();
        }
    });

    List<String> data = new ArrayList<String>(Arrays.asList("{A}", "{B}", "{C}"));
    Flux<String> intervalFlux = Flux
                                .interval(Duration.ofMillis(500))
                                .map(tick -> {
                                    if (tick < data.size())
                                        return "item " + tick + ": " + data.get(tick.intValue());
                                    return "Done (tick == data.size())";
                                });

        return webSocketSession.send(intervalFlux
          .map(webSocketSession::textMessage));
    }


}

【问题讨论】:

标签: spring-boot websocket reactive-programming spring-webflux


【解决方案1】:

您不应该在反应式管道中使用 subscribeblock - 您可以发现您在这样的管道中,因为处理程序方法的返回类型是 Mono&lt;Void&gt;,这意味着处理传入消息已完成。

在您的情况下,您可能希望阅读第一条消息,检查它是否包含您期望的订阅信息,然后发送消息。

public class TestWebSocketHandler implements WebSocketHandler {

    public Mono<Void> handle(WebSocketSession session) {

        return session.receive()
                .map(WebSocketMessage::getPayloadAsText)
                .flatMap(msg -> {
                    SubscriptionInfo info = extract(msg);
                    if (info == null) {
                        return session.close();
                    }
                    else {
                        Mono<WebSocketMessage> message = Mono.just(session.textMessage("message"));
                        return session.send(message);
                    }
                })
                .then();
    }

    SubscriptionInfo extract(String message) {
        //
    }

    class SubscriptionInfo {
        //
    }
}

【讨论】:

  • 谢谢!这正是我试图做的!可能问题是我正在尝试“走出”地图并尝试连接发送而不是在地图内调用它。只有一个理论问题,下一个是必需的,因为从通量中我需要一个单声道(只有一个项目)才能应用该功能?
  • 不,这不是必需的。您可以将它应用于该会话中发送的每条消息。这可能是一个错误,我已经在我的回答中改变了这一点
  • In 在 next() 中确实是正确的,因为 flatMax 正在使用一个函数 我已经在查看代码时回答了我的疑问!谢谢!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-07-17
  • 1970-01-01
  • 2020-05-13
  • 2020-10-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多