You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何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的创建时机和作用域不同:

  1. Bean形式的WebSocketHandler

    • Flux<ProfileCreatedEvent> flux = Flux.create(eventListener).share();在Spring容器初始化时执行(@Bean方法触发),创建一个全局唯一的共享事件流。
    • share()操作符将这个Flux转换为热流,所有后续的WebSocket连接都会订阅同一个Flux实例。新连接会从订阅时刻开始接收后续所有事件,自然实现多连接共享消息的效果。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 06:50:00