【发布时间】:2013-12-19 19:00:56
【问题描述】:
观察以下函数:
public Task RunInOrderAsync<TTaskSeed>(IEnumerable<TTaskSeed> taskSeedGenerator,
CreateTaskDelegate<TTaskSeed> createTask,
OnTaskErrorDelegate<TTaskSeed> onError = null,
OnTaskSuccessDelegate<TTaskSeed> onSuccess = null) where TTaskSeed : class
{
Action<Exception, TTaskSeed> onFailed = (exc, taskSeed) =>
{
if (onError != null)
{
onError(exc, taskSeed);
}
};
Action<Task> onDone = t =>
{
var taskSeed = (TTaskSeed)t.AsyncState;
if (t.Exception != null)
{
onFailed(t.Exception, taskSeed);
}
else if (onSuccess != null)
{
onSuccess(t, taskSeed);
}
};
var enumerator = taskSeedGenerator.GetEnumerator();
Task task = null;
while (enumerator.MoveNext())
{
if (task == null)
{
try
{
task = createTask(enumerator.Current);
Debug.Assert(ReferenceEquals(task.AsyncState, enumerator.Current));
}
catch (Exception exc)
{
onFailed(exc, enumerator.Current);
}
}
else
{
task = task.ContinueWith((t, taskSeed) =>
{
onDone(t);
var res = createTask((TTaskSeed)taskSeed);
Debug.Assert(ReferenceEquals(res.AsyncState, taskSeed));
return res;
}, enumerator.Current).TaskUnwrap();
}
}
if (task != null)
{
task = task.ContinueWith(onDone);
}
return task;
}
其中TaskUnwrap 是标准Task.Unwrap 的状态保留版本:
public static class Extensions
{
public static Task TaskUnwrap(this Task<Task> task, object state = null)
{
return task.Unwrap().ContinueWith((t, _) =>
{
if (t.Exception != null)
{
throw t.Exception;
}
}, state ?? task.AsyncState);
}
}
RunInOrderAsync 方法允许异步运行 N 个任务,但顺序是 - 一个接一个。实际上,它运行从给定种子创建的任务,并发限制为 1。
让我们假设createTask委托从种子创建的任务不对应于多个并发任务。
现在,我想输入 maxConcurrencyLevel 参数,因此函数签名如下所示:
Task RunInOrderAsync<TTaskSeed>(int maxConcurrencyLevel,
IEnumerable<TTaskSeed> taskSeedGenerator,
CreateTaskDelegate<TTaskSeed> createTask,
OnTaskErrorDelegate<TTaskSeed> onError = null,
OnTaskSuccessDelegate<TTaskSeed> onSuccess = null) where TTaskSeed : class
在这里我有点卡住了。
SO 有如下问题:
- System.Threading.Tasks - Limit the number of concurrent Tasks
- Task based processing with a limit for concurrent task number with .NET 4.5 and c#
- .Net TPL: Limited Concurrency Level Task scheduler with task priority?
基本上提出了两种解决问题的方法:
- 将
Parallel.ForEach与ParallelOptions一起使用,指定MaxDegreeOfParallelism属性值等于所需的最大并发级别。 - 使用具有所需
MaximumConcurrencyLevel值的自定义TaskScheduler。
第二种方法并没有减少它,因为涉及的所有任务都必须使用相同的任务调度程序实例。为此,所有用于返回 Task 的方法都必须具有接受自定义 TaskScheduler 实例的重载。不幸的是,微软在这方面并不是很一致。例如,SqlConnection.OpenAsync 不接受这样的论点(但 TaskFactory.FromAsync 接受)。
第一种方法意味着我必须将任务转换为操作,如下所示:
() => t.Wait()
我不确定这是一个好主意,但我很乐意就此获得更多意见。
另一种方法是使用TaskFactory.ContinueWhenAny,但这很麻烦。
有什么想法吗?
编辑 1
我想澄清想要限制的原因。我们的任务最终对同一个 SQL 服务器执行 SQL 语句。我们想要的是一种限制并发传出 SQL 语句数量的方法。完全有可能在其他代码段中同时执行其他 SQL 语句,但这是一个批处理器,可能会淹没服务器。
现在,请注意,虽然我们讨论的是同一个 SQL 服务器,但同一台服务器上有许多数据库。所以,这并不是要限制打开同一个数据库的 SQL 连接的数量,因为数据库可能根本不一样。
这就是为什么像 ThreadPool.SetMaxThreads() 这样的末日解决方案是无关紧要的。
现在,关于SqlConnection.OpenAsync。它之所以异步是有原因的——它可能会往返于服务器,因此可能会受到网络延迟和分布式环境的其他可爱副作用的影响。因此,它与接受TaskScheduler 参数的其他异步方法没有什么不同。我倾向于认为不接受只是一个错误。
编辑 2
我想保留原始函数的异步精神。因此,我希望避免任何明确的阻塞解决方案。
编辑 3
感谢@fsimonazzi's answer 我现在有了所需功能的有效实现。代码如下:
var sem = new SemaphoreSlim(maxConcurrencyLevel);
var tasks = new List<Task>();
var enumerator = taskSeedGenerator.GetEnumerator();
while (enumerator.MoveNext())
{
tasks.Add(sem.WaitAsync().ContinueWith((_, taskSeed) =>
{
Task task = null;
try
{
task = createTask((TTaskSeed)taskSeed);
if (task != null)
{
Debug.Assert(ReferenceEquals(task.AsyncState, taskSeed));
task = task.ContinueWith(t =>
{
sem.Release();
onDone(t);
});
}
}
catch (Exception exc)
{
sem.Release();
onFailed(exc, (TTaskSeed)taskSeed);
}
return task;
}, enumerator.Current).TaskUnwrap());
}
return Task.Factory.ContinueWhenAll(tasks.ToArray(), _ => sem.Dispose());
【问题讨论】:
-
BlockingCollection 允许你设置一个 BoundedCapacity
-
好吧,当然 SqlConnection.OpenAsync 没有那个选项。它不会烧掉 N 个线程。事实上,它不消耗任何的可能性很大。如果您的任务非常不守规矩,以至于您必须提供这种保证,那么 ThreadPool.SetMaxThreads() 是核选项。
-
请参阅EDIT 1。
-
@Blam - 如果您提供代码,您的回复将是一个很好的答案。
-
这就是为什么它是评论而不是答案。请参阅 MSDN 上的文档。有一个使用 BoundedCapacity 的示例。 SO 不是代码生成工具。
标签: .net asynchronous