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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 08:01:24