SmallRye Mutiny异步订阅异常:无法触发onNext方法求助
问题原因与解决方法
你的问题出在自定义Subscriber的onSubscribe方法中——没有调用Subscription的request()方法来向上游请求数据。
按照Reactive Streams规范,当订阅建立后(进入onSubscribe),下游必须主动调用subscription.request(n)来告诉上游“我准备好接收n条数据了”,否则上游不会发送任何事件到onNext。你现在的代码只打印了日志,没有触发数据请求,自然不会进入onNext方法。
修改后的代码示例
只需要在onSubscribe中添加请求逻辑即可,比如请求所有可用数据:
private final ManagedExecutor managedExecutor; public void writeTo(StreamingInfo streamingInfo) { streamingInfo .getEvents() .runSubscriptionOn(managedExecutor) .subscribe() .withSubscriber( new Subscriber < String > () { @Override public void onSubscribe(Subscription s) { System.out.println("OnSubscription Method"); // 关键:请求上游发送数据,Long.MAX_VALUE表示请求所有可用数据 s.request(Long.MAX_VALUE); System.out.println("ON SUBS END"); } @Override public void onNext(String event) { System.out.println("On Next Method: " + event); } @Override public void onError(Throwable t) { System.out.println("OnError Method: " + t.getMessage()); } @Override public void onComplete() { System.out.println("On Complete Method"); } }); }
补充说明
- 如果不需要一次性请求所有数据,可以根据实际场景请求固定数量(比如
s.request(10)),处理完当前批次后再在onNext中继续请求; - 直接收集到List正常是因为Mutiny的内置收集器已经帮你处理了请求逻辑,无需手动调用
request()。
内容的提问来源于stack exchange,提问作者BATMAN_2008
相关产品推荐
相关产品推荐

