Spring WebFlux整合WebSocket:优雅融合认证数据与事件处理方案
问题描述
我开发了一款基于WebSocket的消息交互应用,采用Spring WebFlux与Reactive WebSockets技术栈。安全组件通过以下方式提供包含当前Principal的Mono对象:
Mono<Principal> principalMono = session.getHandshakeInfo().getPrincipal();
目前我只能通过zip操作将该Mono与入站消息Flux、出站消息Flux结合,但这种方式繁琐重复且不够优雅,现有实现代码如下:
public class ChatSocketHandler implements WebSocketHandler { private Flux<DownstreamMessage> outputEvents; public ChatSocketHandler(Sinks.Many<DownstreamMessage> downstreamPublisher) { this.outputEvents = downstreamPublisher.asFlux(); } @Override public Mono<Void> handle(WebSocketSession session) { var input = session.receive() .map(WebSocketMessage::getPayloadAsText) .map(this::toEvent) .flatMap(event -> Mono.just(event).zipWith(session.getHandshakeInfo().getPrincipal())) .doOnNext(objects -> /* do something with objects.getT1() and objects.getT2()*/) .then(); var output = session.send( outputEvents .flatMap(event -> Mono.just(event).zipWith(session.getHandshakeInfo().getPrincipal())) .filter(objects -> /* Filter events based on event objects.getT1() and Principal objects.getT2() */) .map(Tuple2::getT1) .map(this::toJSON) .map(session::textMessage) ) .then(); return Mono.zip(input, output).then(); } }
请问是否存在更优雅的解决方案?
优雅解决方案
1. 提前获取Principal并复用
session.getHandshakeInfo().getPrincipal()是冷Mono,仅在订阅时获取。可以先获取Principal实例,再将其与输入、输出流绑定,避免重复的zip操作:
@Override public Mono<Void> handle(WebSocketSession session) { Mono<Principal> principalMono = session.getHandshakeInfo().getPrincipal(); // 处理入站消息:直接复用已获取的Principal Mono<Void> inputHandler = principalMono.flatMap(principal -> session.receive() .map(WebSocketMessage::getPayloadAsText) .map(this::toEvent) .doOnNext(event -> /* 直接使用event和principal处理业务逻辑 */) .then() ); // 处理出站消息:用Principal过滤后发送 Mono<Void> outputHandler = principalMono.flatMap(principal -> session.send( outputEvents .filter(event -> /* 直接使用event和principal做过滤判断 */) .map(this::toJSON) .map(session::textMessage) ) ); return Mono.zip(inputHandler, outputHandler).then(); }
这种方式只需要获取一次Principal,代码结构更简洁,避免了重复的flatMap+zip模板代码。
2. 封装通用工具方法(可选)
如果需要在多个WebSocketHandler中复用Principal与流的绑定逻辑,可以封装一个通用方法:
private <T> Flux<Tuple2<T, Principal>> withPrincipal(Flux<T> flux, Mono<Principal> principalMono) { return principalMono.flatMapMany(principal -> flux.map(event -> Tuples.of(event, principal))); }
之后在处理流时直接调用该方法,减少重复代码:
@Override public Mono<Void> handle(WebSocketSession session) { Mono<Principal> principalMono = session.getHandshakeInfo().getPrincipal(); var input = withPrincipal( session.receive() .map(WebSocketMessage::getPayloadAsText) .map(this::toEvent), principalMono ) .doOnNext(tuple -> /* 使用tuple.getT1()(事件)和tuple.getT2()(Principal)处理 */) .then(); var output = session.send( withPrincipal(outputEvents, principalMono) .filter(tuple -> /* 用tuple.getT1()和tuple.getT2()做过滤 */) .map(Tuple2::getT1) .map(this::toJSON) .map(session::textMessage) ) .then(); return Mono.zip(input, output).then(); }
3. 利用Reactor上下文传递Principal(进阶)
如果需要在嵌套操作中共享Principal,可以通过contextWrite将其存入Reactor上下文,后续通过Mono.deferContextual获取:
@Override public Mono<Void> handle(WebSocketSession session) { return session.getHandshakeInfo().getPrincipal() .flatMap(principal -> { // 将Principal存入上下文 Context context = Context.of(Principal.class, principal); Mono<Void> input = session.receive() .map(WebSocketMessage::getPayloadAsText) .map(this::toEvent) .contextWrite(context) .flatMap(event -> Mono.deferContextual(ctx -> { Principal currentPrincipal = ctx.get(Principal.class); /* 处理event和currentPrincipal */ return Mono.empty(); })) .then(); Mono<Void> output = session.send( outputEvents .contextWrite(context) .filterWhen(event -> Mono.deferContextual(ctx -> { Principal currentPrincipal = ctx.get(Principal.class); /* 基于event和currentPrincipal做过滤,返回Mono<Boolean> */ return Mono.just(true); })) .map(this::toJSON) .map(session::textMessage) ); return Mono.zip(input, output).then(); }); }
这种方式适合需要在多层嵌套操作中访问Principal的场景,无需显式传递参数。
内容的提问来源于stack exchange,提问作者VladRia
相关产品推荐
相关产品推荐

