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

如何在Akka Graph中使用基于Actor的Source并向其发送数据

问题原因

你拿不到ActorRef的核心原因是:Source.actorRef()生成的是流的蓝图,ActorRef是流运行物化之后才会生成的对象,你之前构建RunnableGraph时仅将sink作为物化值保留,直接丢弃了source对应的物化值,自然拿不到发消息用的Actor引用。

修复方法

构建图的时候同时保留source和sink的物化值,流启动后从返回的物化结果里取出ActorRef即可,具体修改如下:

  1. 调用GraphDSL.create时同时传入source、sink,以及合并两个物化值的函数,不要在builder内部重复add传入的source
  2. 调用run()方法启动流后,从返回结果中取出ActorRef,直接调用tell发消息

修正后的核心代码

// 引入Akka自带的二元组类
import akka.japi.Pair;

// 图构建部分:同时保留source和sink的物化值
RunnableGraph<Pair<ActorRef, CompletionStage<Done>>> graph = RunnableGraph.fromGraph(
        GraphDSL.create(
                integerSource,
                sink,
                (actorRef, doneStage) -> Pair.create(actorRef, doneStage), // 合并两个物化值返回
                (builder, sourceShape, out) -> {
                    FlowShape<Integer, Integer> flow1Shape = builder.add(flow1);
                    // 注意:原代码这里两个分支都加了flow1,会导致两个分支执行相同逻辑,需要区分的话请单独定义flow2
                    FlowShape<Integer, Integer> flow2Shape = builder.add(flow1);
                    UniformFanOutShape<Integer, Integer> broadcast = builder.add(Broadcast.create(2));
                    UniformFanInShape<Integer, Integer> merge = builder.add(Merge.create(2));

                    builder.from(sourceShape)
                            .viaFanOut(broadcast)
                            .via(flow1Shape);
                    builder.from(broadcast).via(flow2Shape);
                    builder.from(flow1Shape)
                            .viaFanIn(merge)
                            .to(out);
                    builder.from(flow2Shape).viaFanIn(merge);

                    return ClosedShape.getInstance();
                }
        )
);

// 启动流,拿到物化值
Pair<ActorRef, CompletionStage<Done>> runResult = graph.run(actorSystem);
ActorRef sourceActor = runResult.first();

// 发送数据
sourceActor.tell(1, ActorRef.noSender());
sourceActor.tell(2, ActorRef.noSender());
sourceActor.tell(3, ActorRef.noSender());

// 所有数据发完后发送Done触发流正常结束
sourceActor.tell(Done.done(), ActorRef.noSender());
注意事项
  • 不要在builder内部重复调用builder.add(integerSource),否则会生成新的source副本,拿到的ActorRef和实际运行的source不对应,消息无法进入流处理链路
  • tell方法需要传入发送方的ActorRef,没有指定发送方时传ActorRef.noSender()即可
  • Source本身是流的设计蓝图,不是运行时的实体,必须等流run物化之后才能拿到对应的ActorRef,不能直接对Source实例调用tell方法

内容的提问来源于stack exchange,提问作者blue-sky

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 22:33:16