C#中async/await为何无法直接实现并行处理需手动创建线程?
我编写了一款Windows服务,希望在保留现有业务逻辑完全不变的前提下,实现所有任务的并行处理。由于实际生产代码抽象程度较高且涉及涉密内容,无法直接贴出完整源码,以下是核心逻辑概要:
该应用是常驻运行的进程调度器,基于Entity Framework 6扫描数据库中的任务记录,记录核心字段包括:
- 待运行进程的路径
- 进程计划执行时间
- 任务调度周期
基础功能逻辑
- 循环查询数据库中所有激活状态的记录,返回全部计划任务详情
- 设定缓冲阈值,比对任务计划时间与当前系统时间判断是否到期
- 对到期需要执行的任务,原有逻辑通过
new Process().Start(...)传入记录中的路径启动进程,校验目标文件存在且可执行后启动进程,等待进程退出或达到配置的超时阈值 - 进程执行的退出码(或进程挂起无退出码的状态)是唯一判断依据:决定任务记录是否保持激活状态、继续循环动态调度到后续执行时间,或是标记为失活、将错误日志写入数据库对应记录
- 调度流程会永久持续运行,除非收到显式停止指令
当前并行实现情况(速度提升1000%,但疑似漏处理记录)
我最初怀疑是访问数据库前未加锁导致问题,后续排查发现问题根源是我使用using (var process) {...}包裹进程实例,导致对象被提前释放抛出异常,盯着代码排查数天后才发现这个为了规范资源释放写出的低级错误 ;p
当前可运行的并行版本通过手动创建Thread实现,代码如下:
var threads = new List<Thread>(); schedules.ForEach(schedule => { // 我也尝试过ThreadPool.QueueUserWorkerItem(...),但查阅资料后了解到它本质上是Task.Run()的冗长写法,实际测试中该实现和手动new Thread的运行效果并不一致 var thread = new Thread(() => await ProcessSchedule(schedule)); // 生产环境中我实际使用了SemaphoreSlim控制并发,上述代码做了简化 thread.Start(); threads.Add(thread); }); // 退出前等待所有线程执行完成 while (!threads.All(instance => !instance.IsAlive)) { await Task.Delay(debounceValue); continue; }
串行运行版本(无逻辑错误但阻塞严重、速度极慢)
串行版本代码如下:
var tasks = new List<Task>(); schedules.ForEach(schedule => { // 我也尝试过直接在这里await,但显然会阻塞同线程上的其他操作,所以我把任务加到列表里,等循环结束、进程都独立运行后再统一等待完成 tasks.Add(ProcessSchedule(schedule)); }); // 退出前等待 // 我原本预期这段代码能实现并行,但实际运行仍然是逐记录串行执行 :*( // 也试过使用Task.Run(() => await ProcessSchedule(...)),依然没有效果 await Task.WhenAll(tasks);
我原本预期将ProcessSchedule返回的Task加入列表后调用Task.WhenAll等待所有任务完成即可实现并行,但实际运行仍然是逐记录串行执行,即便包裹Task.Run(() => await task(...))也没有效果。
注:实际生产代码中我会将任务/线程列表向上层传递,以便在所有任务运行期间统一等待处理,上述代码是简化后的近似伪代码,仅用于尽可能清晰地展示我遇到的核心问题。
ProcessSchedule方法内部逻辑
该方法为async方法,核心逻辑是启动新进程并等待其退出,进程退出后通过Entity Framework 6将执行成功/失败的状态写入对应调度记录的数据库字段中,核心代码如下:
new Process(startInfo).Start(); // 监控进程运行状态,持久化退出状态 dbContext.SaveChangesAsync(); process.StandardError += handleProcExitListener; process.StandardOutput += handleProcExitListener; process.Exited += (...) => handleProcExitListener(...);
我可以确认的是: 代码中不存在未await的async方法,所有await写法均为await Task.Run(MethodAsync)、Main方法中的await Task.WhenAll(task);这类符合规范的写法。
我疑惑是否是因为DbContext默认非线程安全的特性,导致async/await出现阻塞?如果是这个原因,希望能得到正确的实现方案指引。
我尝试了多种技术方案,但只有使用多线程(直接new Thread或是使用ThreadPool)的情况下,才能实现启动所有进程后统一等待、根据进程结束状态做后续处理的并行效果,否则始终串行执行。
自从C#引入async/await特性后我已经很久没有手动操作线程了,因此在没有彻底搞懂原理的情况下直接使用线程让我很疑惑,希望能得到解答理清我遗漏的知识点。
在我看来async本质上是对状态机特性的封装语法糖,但查阅资料时看到有说法称TPL async/await出现后ThreadPool.QueueUserWorkerItem(...)已经过时。如果async/await本身不会创建新线程,那么不使用手动线程是否可以实现进程并行执行?另外我场景中的单个进程运行时长在10分钟到45分钟不等,因此并行执行对性能提升至关重要。
受限于项目运行环境只能使用.NET Framework 4.8,我无法使用.NET 5及以上版本新增的WaitForExitAsync()异步方法。
已尝试的解决方案
我参考技术社区的《Async process start and wait for it to finish》问题方案,封装了WaitForExitAsync扩展方法,代码如下:
public static Task<bool> WaitForExitAsync(this Process process, TimeSpan timeout) { ManualResetEvent processWaitObject = new ManualResetEvent(false); processWaitObject.SafeWaitHandle = new SafeWaitHandle(process.Handle, false); TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); RegisteredWaitHandle registeredProcessWaitHandle = null; registeredProcessWaitHandle = ThreadPool.RegisterWaitForSingleObject( processWaitObject, delegate(object state, bool timedOut) { if (!timedOut) { registeredProcessWaitHandle.Unregister(null); } processWaitObject.Dispose(); tcs.SetResult(!timedOut); }, null /* state */, timeout, true /* executeOnlyOnce */); return tcs.Task; }
内容的提问来源于stack exchange,提问作者AlphaG33k

