为何WebFlux中Lambda实现的WebSocketHandler共享响应,组件实现则不?
WebFlux WebSocket多连接事件共享问题解析
问题场景
以下返回Lambda形式WebSocket#handle的配置可将Entity共享给所有浏览器连接:
@Bean public WebSocketHandler webSocketHandler (ObjectMapper objectMapper, ProfileCreatedEventListener eventListener) { Flux<ProfileCreatedEvent> flux = Flux.create(eventListener).share(); return session -> { Flux<WebSocketMessage> messageFlux = flux.map(evt -> { try { return objectMapper.writeValueAsString(evt.getSource()); } catch(IOException e) { throw new RuntimeException(e); } }) .map(str -> { return session.textMessage(str); //returns WebSocketMessage }); return session.send(messageFlux); } ; //end WebSocket#handle } //end @Bean webSocketHandler
但将完全相同的WebSocketHandler#handle方法编写在实现WebSocketHandler的@Component中时,仅第一个连接的浏览器能收到JSON实体。请问这两种实现的差异是什么?
@Configuration @Bean @Autowired public HandlerMapping handlerMapping(MyWebSocketHandler wsh) { Map<String, WebSocketHandler> map = new HashMap<>(); map.put("/ws/profiles", wsh); int order = 10; // after ProfileEndpointConfiguration return new SimpleUrlHandlerMapping(map, order); }
@Component public class MyWebSocketHandler implements WebSocketHandler { @Override public Mono<Void> handle(WebSocketSession session) { Flux<ProfileCreatedEvent> fluxEvent = Flux.create(listener).share(); Flux<WebSocketMessage> messageFlux = fluxEvent.map ( evt -> { ObjectMapper objectMapper = new ObjectMapper(); String json = null; try { json = objectMapper.writeValueAsString(evt.getSource()); } catch(IOException e) { throw new RuntimeException(e); } return json; }) .map(str -> session.textMessage(str)); return session.send(messageFlux); } //end handle }
整体架构
JSON实体上传至RouterFunction,路由绑定的处理器通过Service组件调用ReactiveCrudRepository#save,Service的create方法会发布插入的实体。
//Model public class Profile { @Id private String id; private String email; } public interface ProfileRepository extends ReactiveMongoRepository<Profile, String> {} //Service Component public ProfileService(ApplicationEventPublisher publisher, ProfileRepository profileRepository) {...} public Mono<Profile> create(String email){ return profileRepository.save( new Profile(null, email) ) .doOnSuccess( profile -> publisher.publishEvent( new ProfileCreatedEvent( profile) ) ); } public class ProfileCreatedEvent extends ApplicationEvent { public ProfileCreatedEvent(Profile source) { super( source ); } } //Handler public Mono<ServerResponse> create(ServerRequest request){ final Flux<Profile> profileFlux = request.bodyToFlux( Profile.class ) .flatMap( p -> profileService.create( p.getEmail() ) ); return defaultWriteResponse( profileFlux ); } //Note: The body is not returned in the ServerResponse, so that Profile entity can be published. private static Mono<ServerResponse> defaultWriteResponse(Publisher<Profile> profiles) { return Mono .from(profiles) .flatMap(p -> ServerResponse .created( URI.create( "/profiles/" + p.getId())) .contentType(MediaType.APPLICATION_JSON) .build() ); }
实现ApplicationListener与Consumer<FluxSink>的组件中编写了onApplicationEvent方法,该组件用作Flux#create的Consumer:
@Component public class ProfileCreatedEventListener implements ApplicationListener <ProfileCreatedEvent>, Consumer <FluxSink<ProfileCreatedEvent>> { @Override public void onApplicationEvent(ProfileCreatedEvent event) { this.queue.offer(event); } @Override public void accept(FluxSink<ProfileCreatedEvent> sink) { this.executor.execute(() -> { while(true) { try { ProfileCreatedEvent event = que.take(); sink.next(event) ; } catch (InterruptedException e) {} }//end while } //end lambda ) ; //end execute } //end accept }//end class
测试方式为通过curl上传实体至路由,前端HTML页面通过JavaScript的WebSocket#onMessage回调展示JSON文本。
核心差异解析
问题的根源在于事件流Flux的创建时机和作用域不同:
Bean形式的WebSocketHandler
Flux<ProfileCreatedEvent> flux = Flux.create(eventListener).share();在Spring容器初始化时执行(@Bean方法触发),创建一个全局唯一的共享事件流。share()操作符将这个Flux转换为热流,所有后续的WebSocket连接都会订阅同一个Flux实例。新连接会从订阅时刻开始接收后续所有事件,自然实现多连接共享消息的效果。
Component形式的MyWebSocketHandler
Flux<ProfileCreatedEvent> fluxEvent = Flux.create(listener).share();在handle()方法内部执行,每次新WebSocket连接建立时都会创建一个全新的Flux实例。- 虽然调用了
share(),但每个Flux实例都是独立的。而ProfileCreatedEventListener中的队列是单例的,第一个连接创建的Flux会通过que.take()阻塞消费队列中的所有事件,后续连接的Flux永远无法获取到新事件,因此无法收到消息。
修复方案
将事件流的创建提升至类级别,让所有连接复用同一个共享Flux实例,同时避免每次创建新的ObjectMapper:
@Component public class MyWebSocketHandler implements WebSocketHandler { private final Flux<ProfileCreatedEvent> sharedEventFlux; private final ObjectMapper objectMapper; // 构造注入Listener和ObjectMapper public MyWebSocketHandler(ProfileCreatedEventListener listener, ObjectMapper objectMapper) { this.sharedEventFlux = Flux.create(listener).share(); this.objectMapper = objectMapper; } @Override public Mono<Void> handle(WebSocketSession session) { Flux<WebSocketMessage> messageFlux = sharedEventFlux.map(evt -> { try { return objectMapper.writeValueAsString(evt.getSource()); } catch(IOException e) { throw new RuntimeException(e); } }) .map(str -> session.textMessage(str)); return session.send(messageFlux); } }
内容的提问来源于stack exchange,提问作者dinah_foster
相关产品推荐
相关产品推荐

