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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 06:36:00