【发布时间】: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