如何实现并行任务独立循环执行及动态任务管理与控制?
问题描述
现有代码与痛点
我编写了如下代码:
var Options = new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount * 10, CancellationToken = CTS.Token }; while (!CTS.IsCancellationRequested) { var TasksZ = new[] { "Task A", "Task B", "Task C" }; await Parallel.ForEachAsync(TasksZ, Options, async (Comando, Token) => { await MyFunction(Comando); await Task.Delay(1000, Token); });
当前代码存在两个核心问题:
- 所有任务(Task A、B、C)会同时启动,且外层循环必须等待所有任务完成才会再次执行。比如Task A、B仅需10秒完成,但Task C需要2分钟,A和B也得等2分钟才能再次启动。
- 无法动态调整任务列表(添加/移除任务),也没法单独停止/暂停每个任务。
额外需求与场景补充
我基于微软插件支持示例开发应用,通过共用Interface加载不同DLL,每个DLL执行独立任务:读取文件 -> 处理在线POST请求 -> 保存文件 -> 通过自定义类向主应用返回JSON -> 重复执行。
- 不同任务执行频率差异大:部分需高频执行,部分仅需每2分钟执行一次。
- 需要避免双核服务器CPU占用过高。
- 为单独停止/暂停每个任务,需配置独立的TaskCompletionSource,但
MyFunction是主应用与各DLL共用的Interface,疑惑这些TCS是该在各DLL中单独声明,还是在共用Interface中声明一个即可?
解决方案
一、让任务独立循环执行,支持动态调整
放弃外层统一循环+Parallel.ForEachAsync的模式,改为给每个任务单独维护独立的循环逻辑,同时用线程安全集合管理任务列表,实现动态添加/移除。
核心思路
- 用
ConcurrentDictionary<string, CancellationTokenSource>管理每个任务的取消令牌,键为任务标识(如"Task A"),值为对应任务的取消源,方便单独停止任务。 - 每个任务启动后进入自身独立循环,完成一次执行后按自身频率延迟,无需等待其他任务。
- 提供线程安全的添加/移除任务方法。
示例代码
// 线程安全的任务管理集合:键为任务标识,值为任务的取消源 private readonly ConcurrentDictionary<string, CancellationTokenSource> _taskCtsDict = new(); // 全局退出的主取消源 private readonly CancellationTokenSource _globalCts = new(); // 添加任务的方法 public void AddTask(string taskId, Func<string, CancellationToken, Task> taskFunc, int delayMs) { if (_taskCtsDict.TryAdd(taskId, out var cts)) { // 启动独立的任务循环 _ = Task.Run(async () => { var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(_globalCts.Token, cts.Token); try { while (!linkedCts.Token.IsCancellationRequested) { await taskFunc(taskId, linkedCts.Token); // 按任务自身频率延迟 await Task.Delay(delayMs, linkedCts.Token); } } catch (OperationCanceledException) { // 任务被取消,正常退出 } finally { // 清理资源 _taskCtsDict.TryRemove(taskId, out _); linkedCts.Dispose(); cts.Dispose(); } }, _globalCts.Token); } } // 移除/停止单个任务的方法 public void RemoveTask(string taskId) { if (_taskCtsDict.TryRemove(taskId, out var cts)) { cts.Cancel(); } } // 全局停止所有任务 public void StopAllTasks() { _globalCts.Cancel(); } // 使用示例:添加不同频率的任务 AddTask("Task A", MyFunction, 1000); // 高频执行,间隔1秒 AddTask("Task B", MyFunction, 1000); AddTask("Task C", MyFunction, 120000); // 每2分钟执行一次,间隔120秒
二、TaskCompletionSource的声明位置
结论:每个DLL单独声明TCS
原因如下:
- 任务独立性:每个DLL对应一个独立任务,暂停/停止状态是任务自身属性,共用Interface中声明无法区分不同任务的状态。
- 封装性:每个DLL应负责自身任务的控制逻辑,TCS放在DLL内部能更好封装状态,避免主应用或其他DLL干扰。
- 接口职责清晰:共用Interface只需定义任务执行的核心逻辑(如
MyFunction),无需承担任务控制职责,保持接口简洁。
实现方式建议
在每个DLL的任务实现类中声明独立的TaskCompletionSource,同时在Interface中增加控制方法定义(如Pause()、Resume()),让主应用通过接口调用,TCS具体实现由DLL自行处理:
// 共用Interface public interface ITaskPlugin { Task ExecuteAsync(CancellationToken token); void Pause(); void Resume(); } // DLL中的任务实现类 public class TaskA : ITaskPlugin { private TaskCompletionSource<bool>? _pauseTcs; public async Task ExecuteAsync(CancellationToken token) { while (!token.IsCancellationRequested) { // 检查是否需要暂停 if (_pauseTcs != null) { await _pauseTcs.Task; } // 执行核心逻辑:读取文件->POST->保存->返回JSON await ReadFileAsync(); await PostRequestAsync(); await SaveFileAsync(); await ReturnJsonToMainAppAsync(); // 按自身频率延迟 await Task.Delay(1000, token); } } public void Pause() { if (_pauseTcs == null || _pauseTcs.Task.IsCompleted) { _pauseTcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); } } public void Resume() { _pauseTcs?.TrySetResult(true); } // 其他私有方法:ReadFileAsync等 }
三、CPU占用控制建议
- 合理设置并行度:双核服务器上,CPU密集型任务并行数不超过2;IO密集型任务可适当增加(如4-6),避免
Environment.ProcessorCount *10这类过高配置。 - 避免无意义循环:任务执行间隔要合理,通过
Task.Delay给CPU留出空闲时间,不要让高频任务完成后立即循环。 - 全异步IO操作:文件读写、POST请求全部使用异步方法(
async/await),避免阻塞线程,提高线程利用率同时降低CPU占用。
内容的提问来源于stack exchange,提问作者Tyler
相关产品推荐
相关产品推荐

