如何基于akka-stream Source创建仅在订阅时启动的Flow.Publisher
你的理解完全正确。Akka Stream 中的 Source 本质是惰性的流蓝图,只有触发物化操作(调用 run/runWith 等方法)时,才会实际创建运行态的流实例,开始消费上游元素。你现有实现的问题恰恰是提前触发了物化:runWith 会立刻执行物化逻辑,和后续返回的 Publisher 有没有被订阅完全无关,所以 LivePacketIterator 会被提前初始化、拉取元素。
解决方案
你只需要把流的物化逻辑延迟到订阅者发起订阅请求时即可,不需要修改原有的流定义,只需要加一层极薄的 Publisher 代理:
// 首先把原来的流定义为未触发物化的 RunnableGraph private final RunnableGraph<Flow.Publisher<Packet>> liveStreamGraph = Source .fromIterator(() -> new LivePacketIterator(maximumPacketSize)) .filter(OrderSnapshotPacket::isNotEmpty) .async(dbDispatcher) .toMat(JavaFlowSupport.Sink.asPublisher(AsPublisher.WITHOUT_FANOUT), Keep.right()); private final ActorSystem actorSystem; // 注入你的ActorSystem实例 @Override public Flow.Publisher<Packet> getLivePublisher() { return new Flow.Publisher<>() { // 因为用了 WITHOUT_FANOUT,仅支持单订阅,加标记防止重复订阅 private final AtomicBoolean subscribed = new AtomicBoolean(false); @Override public void subscribe(Flow.Subscriber<? super Packet> subscriber) { if (subscribed.compareAndSet(false, true)) { // 仅在首次订阅时才触发流物化,启动迭代器拉取元素 Flow.Publisher<Packet> actualPublisher = liveStreamGraph.run(actorSystem); actualPublisher.subscribe(subscriber); } else { subscriber.onError(new IllegalStateException("当前Publisher仅支持单个订阅者")); } } }; }
方案说明
- 完全符合需求:无订阅时流不会物化,
LivePacketIterator不会被初始化;只有订阅者发起订阅请求后,才会触发流的启动。 - 天然兼容 Flow API 的背压规范:订阅者调用
Subscription.next(n)发出请求后,流才会按照请求量推送元素,不会提前生产数据。 - 如果需要支持多订阅者,只需要把
AsPublisher.WITHOUT_FANOUT改为AsPublisher.WITH_FANOUT,同时去掉单订阅标记,每次订阅都执行一次liveStreamGraph.run(actorSystem)即可。
内容的提问来源于stack exchange,提问作者LordOfThePigs
相关产品推荐
相关产品推荐

