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

如何将IAsyncEnumerable作为Akka Streams的Source使用

IAsyncEnumerable转Akka Streams Source实现方案

方案1:使用原生内置API(推荐)

Akka.NET Streams 从1.4.0版本开始已经原生支持IAsyncEnumerable<T>作为输入源,直接调用Source.FromAsyncEnumerable方法即可,无需额外封装,还能天然兼容Akka Streams的背压机制。
你示例中的代码可以直接修改为如下可运行版本:

using System.Collections.Generic;
using System.Threading.Tasks;
using Akka.Streams.Dsl;

namespace ConsoleApp1
{
    class Program
    {
        static async Task Main(string[] args)
        { 
            // 直接传入IAsyncEnumerable实例,不需要await
            Source.FromAsyncEnumerable(AsyncEnumerable())
                .Via(/*你的自定义处理逻辑*/)
                .RunForeach(Console.WriteLine, materializer);
            // 后续业务逻辑
        }

        private static async IAsyncEnumerable<int> AsyncEnumerable()
        {
            for (int i = 0; i < 100; i++)
            {
                await Task.Delay(100); // 模拟异步IO操作
                yield return i;
            }
        }
    }
}

注意:你原有代码中的Source.From(await AsyncEnumerable())存在编译错误,IAsyncEnumerable不能直接被await,就算你通过ToListAsync把所有元素加载到内存再转Source,也会完全丢失异步流的增量加载、低内存占用的优势。

方案2:低版本Akka.NET自定义实现

如果你使用的Akka.NET版本低于1.4.0,没有内置的FromAsyncEnumerable方法,可以通过Source.UnfoldAsync自行封装扩展方法:

using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Akka;
using Akka.Streams;
using Akka.Streams.Dsl;

public static class AkkaStreamsExtensions
{
    public static Source<T, NotUsed> FromCustomAsyncEnumerable<T>(Func<IAsyncEnumerable<T>> enumerableFactory)
    {
        return Source.UnfoldAsync<IAsyncEnumerator<T>, T>(
            initialState: null,
            async (enumerator) =>
            {
                // 第一次执行时初始化枚举器
                if (enumerator == null)
                {
                    enumerator = enumerableFactory().GetAsyncEnumerator();
                }

                // 异步迭代下一个元素
                if (await enumerator.MoveNextAsync())
                {
                    return (enumerator, enumerator.Current);
                }

                // 迭代结束,释放枚举器资源
                await enumerator.DisposeAsync();
                return null;
            });
    }
}

使用时直接调用扩展方法即可:

Source.FromCustomAsyncEnumerable(AsyncEnumerable())
    .Via(/*你的处理逻辑*/)
    // 后续操作

内容的提问来源于stack exchange,提问作者Dmitry

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 17:06:03