You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现并行任务独立循环执行及动态任务管理与控制?

问题描述

现有代码与痛点

我编写了如下代码:

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的模式,改为给每个任务单独维护独立的循环逻辑,同时用线程安全集合管理任务列表,实现动态添加/移除。

核心思路

  1. 用ConcurrentDictionary<string, CancellationTokenSource>管理每个任务的取消令牌,键为任务标识(如"Task A"),值为对应任务的取消源,方便单独停止任务。
  2. 每个任务启动后进入自身独立循环,完成一次执行后按自身频率延迟,无需等待其他任务。
  3. 提供线程安全的添加/移除任务方法。

示例代码

// 线程安全的任务管理集合:键为任务标识,值为任务的取消源
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

原因如下:

  1. 任务独立性:每个DLL对应一个独立任务,暂停/停止状态是任务自身属性,共用Interface中声明无法区分不同任务的状态。
  2. 封装性:每个DLL应负责自身任务的控制逻辑,TCS放在DLL内部能更好封装状态,避免主应用或其他DLL干扰。
  3. 接口职责清晰:共用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占用控制建议

  1. 合理设置并行度:双核服务器上,CPU密集型任务并行数不超过2;IO密集型任务可适当增加(如4-6),避免Environment.ProcessorCount *10这类过高配置。
  2. 避免无意义循环:任务执行间隔要合理,通过Task.Delay给CPU留出空闲时间,不要让高频任务完成后立即循环。
  3. 全异步IO操作:文件读写、POST请求全部使用异步方法(async/await),避免阻塞线程,提高线程利用率同时降低CPU占用。

内容的提问来源于stack exchange,提问作者Tyler

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 00:46:03