如何将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
相关产品推荐
相关产品推荐

