【问题标题】:Convert a callback into a reactive publisher (Flux)将回调转换为反应式发布者 (Flux)
【发布时间】:2018-05-16 06:14:29
【问题描述】:

我正在使用第三方库来注册一个MessageListener,当某些事件发生时它们会调用注册的监听器onMessage方法

public interface MessageListener {  
   // third party code, it auto-scans for all MessageListeners and registers them
    void onMessage(Message message);
}


public class SimpleMessageListener implements MessageListener {
   public void onMessage(Message message) {
      //do something non blocking
      //is it possible to 'transmit' to messagePublisher
}
   public Flux<Message> messagePublisher() {
       // a method to which to subscribeOn    
   }
}

所以我的问题是将它变成 Flux 的最佳方法是什么

最后我希望能够做这样的事情

messagePublisher().subscribe(System.out::println);

******************编辑 我的第一次尝试是这样的

private List<FluxSink<Message>> handlers = new ArrayList<>();
public void onMessage(Message message) {
   handlers.forEach(han -> han.next(message));
}
public Flux<Message> messagePublisher() {
        return Flux.create(sink -> {
            handlers.add(sink);
            sink.onDispose(() -> handlers.remove(sink));
        });
    }

这可行 - 但我觉得这不是一个很好的解决方案,让类实现 FluxSink 并手动处理会更好 - 目前我不希望有很多订阅者。 但是很多 MessageListeners(每种类型一个)

【问题讨论】:

  • 您可以将您的编辑添加为答案并接受它,因为这是 IMO 的最佳解决方案。实现FluxSink 将无济于事,因为它不是Flux。此外,每个订阅者需要一个接收器,以跟踪订阅者请求。

标签: java spring-webflux project-reactor


【解决方案1】:

您可以创建单个Flux 实例来桥接MessageListener 观察到的消息,例如

public class SimpleMessageListener implements MessageListener {
   private FluxSink<Message> handler;
   private Flux<Message> flux;

   public SimpleMessageListener() {
      flux = Flux.create(emitter -> {
          handler = emitter;
      }, OverflowStrategy.DROP); // or some other overflow strategy
   }

   public void onMessage(Message message) {
       if (handler != null) {
           /* 
            * null check is required to avoid NPE if a message is received 
            * before any subscription occurs since handler is instantiated
            * lazily when the first subscription is requested
            */
           handler.next(message);
       }
   }

   public Flux<Message> messagePublisher() {
       return flux;
   }
}

现在所有监听器都可以使用 Flux 的 publish() 方法及其返回的 ConnectableFlux 订阅同一个 messsagePublisher() Flux 实例:

// fetch message publisher
Flux<Message> messagePublisher = messageListener.messagePublisher();

// prepare ConenctableFlux
ConnectableFlux<Message> connectableFlux = messagePublisher().publish();

// register subscribers
connectableFlux.subscribe(/* aConsumer */);
connectableFlux.subscribe(/* aCoreSubscriber */);
connectableFlux.subscribe(/* aSubscriber */);

// connect the ConnectableFlux to messagePublisher
connectableFlux.connect();

【讨论】:

    猜你喜欢
    • 2021-12-09
    • 1970-01-01
    • 1970-01-01
    • 2021-12-02
    • 1970-01-01
    • 2017-12-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多