【发布时间】: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
提前致谢
其他信息:
我创建了一个基本的 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));
}
}
【问题讨论】:
-
你能展示一个完整的代码sn-p(类和方法)吗?该代码应该在哪个方法中调用?
-
嗨@BrianClozel,我从这段代码github.com/eugenp/tutorials/tree/master/spring-5-reactive/src/… 开始,这应该在公共Mono
句柄(WebSocketSession webSocketSession)中的github.com/eugenp/tutorials/blob/master/spring-5-reactive/src/… 中提交我已经添加了sn-p上面这段代码(我也使用声明性方式进行了其他试验,但没有结果) -
请编辑您的问题。很难从多个链接中收集数据并弄清楚您尝试了什么。
标签: spring-boot websocket reactive-programming spring-webflux