use*_*129 40 .net c# task-parallel-library async-await
我们使用IEnumerables从数据库中返回大量数据集:
public IEnumerable<Data> Read(...)
{
using(var connection = new SqlConnection(...))
{
// ...
while(reader.Read())
{
// ...
yield return item;
}
}
}
Run Code Online (Sandbox Code Playgroud)
现在我们想使用异步方法来做同样的事情.但是,async没有IEnumerables,因此我们必须将数据收集到列表中,直到加载整个数据集:
public async Task<List<Data>> ReadAsync(...)
{
var result = new List<Data>();
using(var connection = new SqlConnection(...))
{
// ...
while(await reader.ReadAsync().ConfigureAwait(false))
{
// ...
result.Add(item);
}
}
return result;
}
Run Code Online (Sandbox Code Playgroud)
这将消耗服务器上的大量资源,因为所有数据必须在返回之前在列表中.IEnumerables处理大数据流的最佳且易于使用的异步替代方法是什么?我想避免在处理时将所有数据存储在内存中.
i3a*_*non 26
最简单的选择是使用TPL Dataflow.您需要做的就是配置一个ActionBlock处理处理(如果您愿意,并行处理)并将项目异步地"逐条"发送到其中.
我还建议设置一个BoundedCapacity,当处理无法处理速度时,它会限制读取器从数据库中读取数据.
var block = new ActionBlock<Data>(
data => ProcessDataAsync(data),
new ExecutionDataflowBlockOptions
{
BoundedCapacity = 1000,
MaxDegreeOfParallelism = Environment.ProcessorCount
});
using(var connection = new SqlConnection(...))
{
// ...
while(await reader.ReadAsync().ConfigureAwait(false))
{
// ...
await block.SendAsync(item);
}
}
Run Code Online (Sandbox Code Playgroud)
您也可以使用Reactive Extensions,但这是一个比您可能需要的更复杂和更健壮的框架.
Sim*_*ier 11
在处理异步/等待方法的大部分时间里,我发现更容易解决问题,并使用函数(Func<...>)或actions(Action<...>)而不是ad-hoc代码,尤其是使用IEnumerable和yield.
换句话说,当我认为"异步"时,我试图忘记功能"返回值"的旧概念,否则它是如此明显并且我们如此熟悉.
例如,如果您将初始同步代码更改为此代码(processor最终将执行您对一个数据项执行的操作的代码):
public void Read(..., Action<Data> processor)
{
using(var connection = new SqlConnection(...))
{
// ...
while(reader.Read())
{
// ...
processor(item);
}
}
}
Run Code Online (Sandbox Code Playgroud)
然后,异步版本编写起来非常简单:
public async Task ReadAsync(..., Action<Data> processor)
{
using(var connection = new SqlConnection(...))
{
// note you can use connection.OpenAsync()
// and command.ExecuteReaderAsync() here
while(await reader.ReadAsync())
{
// ...
processor(item);
}
}
}
Run Code Online (Sandbox Code Playgroud)
如果您可以通过这种方式更改代码,则不需要任何扩展或额外的库或IAsyncEnumerable内容.
这将消耗服务器上的大量资源,因为所有数据必须在返回之前在列表中.IEnumerables处理大数据流的最佳且易于使用的异步替代方法是什么?我想避免在处理时将所有数据存储在内存中.
如果您不想立即将所有数据发送到客户端,您可以考虑使用Reactive Extensions (Rx)(在客户端上)和SignalR(在客户端和服务器上)来处理此问题.
SignalR将允许异步发送数据到客户端.Rx允许在数据项到达客户端时将LINQ应用于异步数据项序列.但是,这会改变客户端 - 服务器应用程序的整个代码模型.
示例(Samuel Jack的博客文章):
相关问题(如果不是重复):
正如其他一些海报所提到的,这可以用Rx实现.使用Rx,该函数将返回IObservable<Data>可以订阅的函数,并在订阅者可用时将数据推送到订阅者.IObservable还支持LINQ并添加了一些自己的扩展方法.
更新
我添加了一些通用帮助方法,以使读取器的使用可重用,并支持取消.
public static class ObservableEx
{
public static IObservable<T> CreateFromSqlCommand<T>(string connectionString, string command, Func<SqlDataReader, Task<T>> readDataFunc)
{
return CreateFromSqlCommand(connectionString, command, readDataFunc, CancellationToken.None);
}
public static IObservable<T> CreateFromSqlCommand<T>(string connectionString, string command, Func<SqlDataReader, Task<T>> readDataFunc, CancellationToken cancellationToken)
{
return Observable.Create<T>(
async o =>
{
SqlDataReader reader = null;
try
{
using (var conn = new SqlConnection(connectionString))
using (var cmd = new SqlCommand(command, conn))
{
await conn.OpenAsync(cancellationToken);
reader = await cmd.ExecuteReaderAsync(CommandBehavior.CloseConnection, cancellationToken);
while (await reader.ReadAsync(cancellationToken))
{
var data = await readDataFunc(reader);
o.OnNext(data);
}
o.OnCompleted();
}
}
catch (Exception ex)
{
o.OnError(ex);
}
return reader;
});
}
}
Run Code Online (Sandbox Code Playgroud)
的实施ReadData,现在被大大简化.
private static IObservable<Data> ReadData()
{
return ObservableEx.CreateFromSqlCommand(connectionString, "select * from Data", async r =>
{
return await Task.FromResult(new Data()); // sample code to read from reader.
});
}
Run Code Online (Sandbox Code Playgroud)
用法
你可以通过给它一个订阅Observable,IObserver但也有一些带有lambdas的重载.随着数据变得可用,将OnNext调用回调.如果存在异常,则OnError调用回调.最后,如果没有更多数据,则OnCompleted调用回调.
如果要取消observable,只需处理订阅即可.
void Main()
{
// This is an asyncrhonous call, it returns straight away
var subscription = ReadData()
.Skip(5) // Skip first 5 entries, supports LINQ
.Delay(TimeSpan.FromSeconds(1)) // Rx operator to delay sequence 1 second
.Subscribe(x =>
{
// Callback when a new Data is read
// do something with x of type Data
},
e =>
{
// Optional callback for when an error occurs
},
() =>
{
//Optional callback for when the sequenc is complete
}
);
// Dispose subscription when finished
subscription.Dispose();
Console.ReadKey();
}
Run Code Online (Sandbox Code Playgroud)