使用Akka.Net Graphs时如何访问Source.Queue的OfferAsync方法?
在Akka.NET Streams中访问Source.Queue的OfferAsync方法
你当前的问题是sourceQueue定义在GraphDsl的内部作用域里,外部拿不到引用,没法调用OfferAsync。得调整图的创建方式,把队列实例暴露到外部。
修改图的创建代码
用GraphDsl.Create的重载,把Source.Queue作为参数传入,让它能被外部获取:
// 创建图,同时返回队列引用和可运行图 var (sourceQueue, runnableGraph) = RunnableGraph.FromGraph(GraphDsl.Create( // 传入要暴露的Source.Queue实例 Source.Queue<int>(100, OverflowStrategy.Fail), (builder, queue) => { // 这里用queue变量来配置流,比如连接到你的处理逻辑 // 示例:连接到一个打印元素的Sink var processingSink = Sink.ForEach<int>(element => Console.WriteLine($"处理事件: {element}")); builder.From(queue).To(processingSink); // 返回Closed.Instance表示图的所有端口都已连接,是封闭的 return Closed.Instance; } )); // 运行图,此时sourceQueue已经可以在外部使用 runnableGraph.Run(materializer);
在外部循环中使用队列
拿到sourceQueue引用后,直接在你的事件获取循环里调用OfferAsync就行,记得处理返回结果:
while (true) { var events = GetEventsFromExternalSource(); foreach(var singleEvent in events) { var offerResult = await sourceQueue.OfferAsync(singleEvent); // 根据返回结果处理不同情况 switch (offerResult) { case QueueOfferResult.Success: // 元素成功加入队列 break; case QueueOfferResult.Failure: // 队列已满或其他错误,按需处理 Console.WriteLine("事件添加失败:队列已满或流出现错误"); break; case QueueOfferResult.QueueClosed: // 流已经终止,队列关闭,退出循环 Console.WriteLine("队列已关闭,停止添加事件"); return; } } }
注意点
OfferAsync是异步方法,最好用await等待结果,避免丢失状态反馈。- 流终止时(比如Sink完成或出现错误),队列会自动关闭,此时调用
OfferAsync会返回QueueClosed。 - 确保
materializer是正确初始化的,比如通过ActorMaterializer.Create(actorSystem)创建(actorSystem是你的Akka.NET ActorSystem实例)。
内容的提问来源于stack exchange,提问作者Bad Banana
相关产品推荐
相关产品推荐

