如何在Akka Graph中使用基于Actor的Source并向其发送数据
问题原因
你拿不到ActorRef的核心原因是:Source.actorRef()生成的是流的蓝图,ActorRef是流运行物化之后才会生成的对象,你之前构建RunnableGraph时仅将sink作为物化值保留,直接丢弃了source对应的物化值,自然拿不到发消息用的Actor引用。
修复方法
构建图的时候同时保留source和sink的物化值,流启动后从返回的物化结果里取出ActorRef即可,具体修改如下:
- 调用
GraphDSL.create时同时传入source、sink,以及合并两个物化值的函数,不要在builder内部重复add传入的source - 调用
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
相关产品推荐
相关产品推荐

