替代ParallelForEach的方案:程序退出时立即终止并行进程
我开发了一个简单的控制台应用程序,从数据库加载文件至HashSet,随后通过Parallel.ForEach循环并行处理这些文件。为避免不同线程日志互相覆盖的问题,程序会为每个待处理文件启动一个新的Process对象,即打开新控制台窗口运行外部程序。
当前问题是:关闭应用程序时,Parallel.ForEach仍会继续尝试处理更多文件,无法立即终止所有任务。我希望在终止程序时,能立即停止所有正在执行的处理任务。
退出捕获逻辑参考自Stack Overflow的《Capture console exit C#》,程序在收到CTRL+C、点击窗口关闭按钮等取消指令时会执行清理操作。
相关代码片段:
class Program { private static bool _isFileLoadingDone; static ConcurrentDictionary<int, Tuple<Tdx2KlarfParserProcInfo, string>> _currentProcessesConcurrentDict = new ConcurrentDictionary<int, Tuple<Tdx2KlarfParserProcInfo, string>>(); static void Main(string[] args) { try { if (args.Length == 0) { // 响应窗口关闭、CTRL-C、进程终止等事件的模板代码 LaunchFolderMode(); } } } }
Main调用LaunchFolderMode(),后者调用循环执行ParseFiles()的ParseFilesUntilEmpty(),最终在ParseFiles()中执行并行逻辑:
private static void ParseFiles() { filesToProcess = new HashSet<string>(){@"file1", "file2", "file3", "file4"}; // 实际从数据库获取文件,此处为示例 int parallelCount = 2; Parallel.ForEach(filesToProcess, new ParallelOptions { MaxDegreeOfParallelism = parallelCount }, tdxFile =>{ ConfigureAndStartProcess(tdxFile); }); }
ConfigureAndStartProcess()会启动外部进程TDXXMLParser.exe并等待其退出,同时将进程信息存入ConcurrentDictionary。
我曾尝试使用Parallel State Stop方法终止,但程序仍会继续处理部分文件后才退出。想请教:为何Parallel.ForEach在程序关闭后仍会运行?有没有更合适的并行处理方法,能在程序终止时立即停止当前处理的文件且不启动新进程?
为何Parallel.ForEach无法立即终止?
- State.Stop的局限性:
ParallelLoopState.Stop()仅能阻止新任务被调度,但已经进入执行流程的委托(即已调用ConfigureAndStartProcess的任务)会继续运行,直到完成。而你的ConfigureAndStartProcess会等待外部进程退出,这就导致即使触发停止,正在运行的外部进程仍会继续,Parallel循环也会等待这些委托执行完毕才结束。 - 缺乏取消信号传递:如果退出捕获逻辑没有向Parallel循环传递明确的取消信号,循环无法感知程序即将退出,仍会继续调度剩余任务。
改进方案
1. 用CancellationToken实现精准取消
Parallel.ForEach支持通过ParallelOptions传入CancellationToken,结合退出捕获逻辑触发取消,既能阻止新任务启动,又能终止正在运行的外部进程。
示例代码:
private static CancellationTokenSource _cts = new CancellationTokenSource(); // 在退出捕获逻辑中添加取消触发 // 比如在窗口关闭/CTRL-C的处理方法中调用: _cts.Cancel(); // 修改ParseFiles方法 private static void ParseFiles() { filesToProcess = new HashSet<string>(){@"file1", "file2", "file3", "file4"}; int parallelCount = 2; var parallelOptions = new ParallelOptions { MaxDegreeOfParallelism = parallelCount, CancellationToken = _cts.Token // 传入取消令牌 }; try { Parallel.ForEach(filesToProcess, parallelOptions, tdxFile => { // 先检查取消状态,避免启动新进程 parallelOptions.CancellationToken.ThrowIfCancellationRequested(); var process = ConfigureAndStartProcess(tdxFile); try { // 等待进程退出时轮询取消信号 while (!process.WaitForExit(100)) { parallelOptions.CancellationToken.ThrowIfCancellationRequested(); } } catch (OperationCanceledException) { // 取消时强制终止外部进程 process.Kill(); process.WaitForExit(); throw; // 抛出异常让Parallel循环感知取消 } }); } catch (OperationCanceledException) { // 执行必要的资源清理 } } // 修改ConfigureAndStartProcess,返回Process对象以便后续操作 private static Process ConfigureAndStartProcess(string tdxFile) { var process = new Process(); process.StartInfo.FileName = "TDXXMLParser.exe"; process.StartInfo.Arguments = tdxFile; // 其他进程配置... process.Start(); _currentProcessesConcurrentDict.TryAdd(process.Id, Tuple.Create(new Tdx2KlarfParserProcInfo(), tdxFile)); return process; }
2. 退出时主动清理所有外部进程
在程序退出的清理逻辑中,直接遍历_currentProcessesConcurrentDict,终止所有未退出的外部进程,确保没有残留任务:
// 退出清理逻辑 foreach (var kvp in _currentProcessesConcurrentDict) { try { var process = Process.GetProcessById(kvp.Key); if (!process.HasExited) { process.Kill(); process.WaitForExit(); } } catch (Exception) { // 忽略进程已退出等异常 } } _currentProcessesConcurrentDict.Clear();
3. 替换Parallel.ForEach为Task.WhenAll(可选)
如果需要更灵活的并发控制和取消能力,可以改用Task结合信号量管理并发,通过CancellationToken统一取消所有任务:
private static async Task ParseFilesAsync() { filesToProcess = new HashSet<string>(){@"file1", "file2", "file3", "file4"}; int parallelCount = 2; var semaphore = new SemaphoreSlim(parallelCount); // 控制并发数 var tasks = new List<Task>(); foreach (var tdxFile in filesToProcess) { _cts.Token.ThrowIfCancellationRequested(); await semaphore.WaitAsync(_cts.Token); tasks.Add(Task.Run(async () => { try { var process = ConfigureAndStartProcess(tdxFile); while (!process.WaitForExit(100)) { _cts.Token.ThrowIfCancellationRequested(); } } catch (OperationCanceledException) { // 终止对应外部进程 // 可从ConcurrentDictionary中查找并处理 } finally { semaphore.Release(); } }, _cts.Token)); } try { await Task.WhenAll(tasks); } catch (OperationCanceledException) { // 处理取消后的清理 } }
内容的提问来源于stack exchange,提问作者edo101

