Parallel.ForEach vs Parallel.ForEachAsync

在 .NET 的并行编程世界里,Parallel.ForEach 是老牌选手,从 .NET Framework 时代就存在;而 Parallel.ForEachAsync 是 .NET 6 才加入的新秀。它们名字相似,但适用场景完全不同,选错一个可能让程序性能腰斩甚至死锁。

同步的 Parallel.ForEach

Parallel.ForEach 是为 CPU 密集型 任务设计的,它基于同步委托(Action<T>),内部用线程池并行执行。它的执行模型是”分区+阻塞”——把数据源切分成多个分区,每个分区由一个线程处理,主线程会阻塞等待所有分区完成。

// Parallel.ForEach:同步阻塞,适合 CPU 密集型
var numbers = Enumerable.Range(1, 1000);
Parallel.ForEach(numbers, num =>
{
    // 这是 CPU 计算,同步执行
    var result = HeavyCompute(num);
    Console.WriteLine($"处理 {num},结果 {result}");
});
// 这行代码会等所有任务完成后才执行
Console.WriteLine("全部完成");

异步的 Parallel.ForEachAsync

Parallel.ForEachAsync 是为 I/O 密集型 任务设计的,它接受异步委托(Func<T, CancellationToken, ValueTask>),能正确处理 async/await。它的执行模型是”信号量控并发”——用 SemaphoreSlim 控制同时进行的任务数量,不占用线程等待 I/O 完成。

// Parallel.ForEachAsync:异步非阻塞,适合 I/O 密集型
var urls = new[] { "https://api.a.com", "https://api.b.com", "https://api.c.com" };
await Parallel.ForEachAsync(urls, async (url, ct) =>
{
    // 这是 I/O 等待,异步执行,不占线程
    var response = await httpClient.GetAsync(url, ct);
    var content = await response.Content.ReadAsStringAsync(ct);
    Console.WriteLine($"{url} 返回 {content.Length} 字符");
});
// await 保证全部完成后才执行
Console.WriteLine("全部完成");

核心区别一览

特性 Parallel.ForEach Parallel.ForEachAsync
委托类型 Action<T>(同步) Func<T, CancellationToken, ValueTask>(异步)
适用场景 CPU 密集型 I/O 密集型
阻塞调用方 是,阻塞主线程 否,可 await
线程占用 执行期间占用线程 I/O 等待时释放线程
并发控制 Partitioner 分区 SemaphoreSlim 信号量
取消支持 ParallelOptions.CancellationToken 同左
引入版本 .NET Framework 4.0 .NET 6

一个常见误区是用 Parallel.ForEachasync 方法:

// ❌ 错误用法:Parallel.ForEach 不会 await 异步方法
Parallel.ForEach(urls, async url =>
{
    await httpClient.GetAsync(url); // 这个 await 没人等,提前返回
});
// 上面这段代码会在请求还没发出时就"完成"

Parallel.ForEachasync 委托当作返回 Task 的普通方法调用,根本不等待任务完成。这就是 Parallel.ForEachAsync 存在的根本原因。

ForEachAsync 的核心机制

理解 ForEachAsync 的内部机制,能帮助我们写出更可控的并发代码。它的源码非常简洁,核心就两部分:分区并行DOP(Degree of Parallelism)控制

分区并行

ForEachAsync 接受 IAsyncEnumerable<T>IEnumerable<T> 作为数据源。对于 IEnumerable<T>,它会通过 Partitioner 把数据切分成多个块,每个块由一个工作单元顺序消费。这样既能并行处理,又避免了为每个元素都创建任务的开销。

// 模拟 ForEachAsync 内部简化逻辑
public static async Task ForEachAsyncSimple<T>(
    IEnumerable<T> source,
    int maxDop,
    Func<T, CancellationToken, ValueTask> body,
    CancellationToken ct = default)
{
    // 用信号量控制最大并发度
    using var semaphore = new SemaphoreSlim(maxDop);
    var tasks = new List<Task>();

    foreach (var item in source)
    {
        ct.ThrowIfCancellationRequested();
        await semaphore.WaitAsync(ct); // 等待可用槽位

        tasks.Add(Task.Run(async () =>
        {
            try
            {
                await body(item, ct); // 执行异步工作
            }
            finally
            {
                semaphore.Release(); // 释放槽位
            }
        }, ct));
    }

    await Task.WhenAll(tasks); // 等待所有任务完成
}

DOP 控制:信号量控并发

真正的 ForEachAsync 比上面更优化,不会为每个元素都 Task.Run,而是维护固定数量的工作单元,每个工作单元从共享分区里拉取数据。这种设计避免了任务创建开销,特别适合处理海量小任务。

关键点在于:并发度不等于线程数。当 bodyawait I/O 时,工作单元不会占用线程,信号量释放后可以立即启动下一个任务。所以 DOP=10 在 I/O 密集场景下,可能只占用 1-2 个线程,却能维持 10 个并发请求。

MaxDegreeOfParallelism 参数详解

ParallelOptions.MaxDegreeOfParallelism 是控制并发的核心参数,默认值是 Environment.ProcessorCount(CPU 逻辑核心数)。

var options = new ParallelOptions
{
    MaxDegreeOfParallelism = 4 // 最多 4 个并发
};
await Parallel.ForEachAsync(data, options, async (item, ct) =>
{
    await ProcessAsync(item, ct);
});

如何选择合适的值

CPU 密集型任务:DOP 等于或略大于 CPU 核心数。设太大反而因线程切换降低性能。

// CPU 密集型:用核心数即可
var cpuOptions = new ParallelOptions
{
    MaxDegreeOfParallelism = Environment.ProcessorCount
};

I/O 密集型任务:DOP 取决于下游服务的承受能力,而不是 CPU。比如调外部 API,可以设 10-50;写数据库,要看连接池大小(通常 100-200)。

// I/O 密集型:根据下游限流设置
var ioOptions = new ParallelOptions
{
    MaxDegreeOfParallelism = 20 // 20 个并发 HTTP 请求
};

不设置 DOP 的风险:默认值是 ProcessorCount,对于 I/O 任务来说并发度太低。比如 8 核机器默认只允许 8 个并发,批量调 1000 个 API 会非常慢。

// ❌ 不设 DOP,I/O 场景下会很慢
await Parallel.ForEachAsync(urls, async (url, ct) =>
{
    await httpClient.GetAsync(url, ct); // 默认只能 8 并发
});

动态调整 DOP

某些场景需要根据运行时情况动态调整,可以自己用 SemaphoreSlim 实现:

// 动态并发控制:根据响应时间自适应
int currentDop = 10;
using var semaphore = new SemaphoreSlim(currentDop);

async Task ProcessWithAdaptiveDop(string url)
{
    await semaphore.WaitAsync();
    try
    {
        var sw = Stopwatch.StartNew();
        await httpClient.GetAsync(url);
        sw.Stop();

        // 响应快则提高并发,慢则降低
        if (sw.ElapsedMilliseconds < 100 && currentDop < 50)
            Interlocked.Increment(ref currentDop);
        else if (sw.ElapsedMilliseconds > 1000 && currentDop > 5)
            Interlocked.Decrement(ref currentDop);
    }
    finally
    {
        semaphore.Release();
    }
}

CancellationToken 支持

Parallel.ForEachAsync 原生支持取消,通过 ParallelOptions.CancellationToken 传入。取消时会抛 OperationCanceledException,已经在执行的任务会通过 bodyct 参数收到取消信号。

using var cts = new CancellationTokenSource();

// 5 秒后自动取消
cts.CancelAfter(TimeSpan.FromSeconds(5));

var options = new ParallelOptions
{
    MaxDegreeOfParallelism = 10,
    CancellationToken = cts.Token // 关键:传入取消令牌
};

try
{
    await Parallel.ForEachAsync(urls, options, async (url, ct) =>
    {
        // body 内部也要检查 ct,确保及时退出
        var response = await httpClient.GetAsync(url, ct);
        var stream = await response.Content.ReadAsStreamAsync(ct);
        using var reader = new StreamReader(stream);
        // 长耗时操作期间主动检查取消
        while (!reader.EndOfStream)
        {
            ct.ThrowIfCancellationRequested();
            var line = await reader.ReadLineAsync(ct);
            ProcessLine(line);
        }
    });
}
catch (OperationCanceledException)
{
    Console.WriteLine("任务被取消");
}

注意:取消是协作式的ForEachAsync 本身会在分配新任务前检查 ct,但如果 body 内部不检查 ct,已经启动的任务仍会跑完。所以 body 里所有 await 调用都要传入 ct

实际应用场景

场景一:批量 HTTP 请求

最常见的场景——并发拉取多个接口或爬取页面:

var httpClient = new HttpClient();
var urls = LoadUrlsFromFile("urls.txt"); // 假设有 1000 个 URL

var results = new ConcurrentBag<(string Url, string Content)>(); // 线程安全集合

var options = new ParallelOptions
{
    MaxDegreeOfParallelism = 20 // 20 并发,避免压垮目标服务器
};

await Parallel.ForEachAsync(urls, options, async (url, ct) =>
{
    try
    {
        var response = await httpClient.GetAsync(url, ct);
        response.EnsureSuccessStatusCode();
        var content = await response.Content.ReadAsStringAsync(ct);
        results.Add((url, content));
    }
    catch (HttpRequestException ex)
    {
        // 单个失败不影响整体,记录日志继续
        Console.WriteLine($"[失败] {url}: {ex.Message}");
    }
});

Console.WriteLine($"成功 {results.Count}/{urls.Count}");

场景二:批量数据处理

并发处理数据库批量更新,控制好连接池压力:

var records = await LoadRecordsAsync(); // 从数据库加载待处理记录
var connectionStr = "Server=...;Max Pool Size=100";

var options = new ParallelOptions
{
    MaxDegreeOfParallelism = 20 // 配合连接池大小,留余量
};

await Parallel.ForEachAsync(records, options, async (record, ct) =>
{
    // 每个任务用独立连接,避免共享连接的锁竞争
    await using var conn = new SqlConnection(connectionStr);
    await conn.OpenAsync(ct);

    var sql = "UPDATE Orders SET Status=@s WHERE Id=@id";
    await using var cmd = new SqlCommand(sql, conn);
    cmd.Parameters.AddWithValue("@s", "Processed");
    cmd.Parameters.AddWithValue("@id", record.Id);
    await cmd.ExecuteNonQueryAsync(ct);
});

场景三:带进度报告的批量处理

var items = Enumerable.Range(1, 10000).ToList();
var processed = 0;
var sw = Stopwatch.StartNew();

await Parallel.ForEachAsync(items, new ParallelOptions
{
    MaxDegreeOfParallelism = 8
}, async (item, ct) =>
{
    await ProcessItemAsync(item, ct);
    var done = Interlocked.Increment(ref processed); // 原子计数
    if (done % 100 == 0)
    {
        Console.WriteLine($"进度 {done}/{items.Count},已用 {sw.ElapsedMilliseconds}ms");
    }
});
Console.WriteLine($"全部完成,耗时 {sw.ElapsedMilliseconds}ms");

常见陷阱与最佳实践

陷阱一:共享状态未加锁

Parallel.ForEachAsync 是并发的,多个任务可能同时访问同一资源。普通 List<T>、计数器、字典都不是线程安全的:

// ❌ 错误:List 不是线程安全的
var list = new List<int>();
await Parallel.ForEachAsync(data, async (x, ct) =>
{
    await Task.Delay(10, ct);
    list.Add(x); // 并发 Add 会丢数据或抛异常
});

// ✅ 正确:用 ConcurrentBag
var bag = new ConcurrentBag<int>();
await Parallel.ForEachAsync(data, async (x, ct) =>
{
    await Task.Delay(10, ct);
    bag.Add(x); // 线程安全
});

// ✅ 正确:计数用 Interlocked
var counter = 0;
await Parallel.ForEachAsync(data, async (x, ct) =>
{
    await Task.Delay(10, ct);
    Interlocked.Increment(ref counter); // 原子操作
});

陷阱二:在 body 里用 async void

async void 无法被 await,异常也无法捕获,会导致进程崩溃:

// ❌ 错误:async void 会让异常丢失
await Parallel.ForEachAsync(data, async (x, ct) =>
{
    // 编译器允许但这是 async void,异常会抛到 SynchronizationContext
    SomeFireAndForgetMethod(x);
});

// ✅ 正确:完整 await
await Parallel.ForEachAsync(data, async (x, ct) =>
{
    await SomeAsyncMethod(x, ct);
});

陷阱三:DOP 设置过大压垮下游

// ❌ 错误:DOP=1000 调外部 API,会被限流封 IP
await Parallel.ForEachAsync(urls, new ParallelOptions
{
    MaxDegreeOfParallelism = 1000
}, async (url, ct) => await httpClient.GetAsync(url, ct));

// ✅ 正确:根据目标服务能力设置,并加重试
await Parallel.ForEachAsync(urls, new ParallelOptions
{
    MaxDegreeOfParallelism = 10
}, async (url, ct) =>
{
    await RetryAsync(async () =>
    {
        var resp = await httpClient.GetAsync(url, ct);
        resp.EnsureSuccessStatusCode();
        return resp;
    }, maxRetries: 3);
});

陷阱四:用 Task.Run 包裹 I/O 任务

有些开发者习惯 Task.Run(() => ...) 来”并行化”,但对 I/O 任务这是反模式:

// ❌ 错误:Task.Run 占线程等 I/O,浪费资源
var tasks = urls.Select(url => Task.Run(async () =>
{
    return await httpClient.GetAsync(url); // 占着线程等响应
}));
await Task.WhenAll(tasks);

// ✅ 正确:ForEachAsync 直接异步,I/O 等待时释放线程
await Parallel.ForEachAsync(urls, async (url, ct) =>
{
    await httpClient.GetAsync(url, ct);
});

最佳实践总结

  1. 按任务类型选 API:CPU 密集用 Parallel.ForEach,I/O 密集用 Parallel.ForEachAsync
  2. 显式设置 DOP:I/O 场景根据下游能力设置,别用默认的 ProcessorCount
  3. body 内传递 CancellationToken:所有 await 都要带上 ct,确保能及时取消。
  4. 共享状态用并发集合ConcurrentBagConcurrentDictionary,计数用 Interlocked
  5. 异常处理要兜底:单个任务失败不应中断整体,用 try-catch 包住 body 内部逻辑。
  6. 配合限流和重试:调外部服务时加上 Polly 等重试策略,避免瞬时故障导致批量失败。

小结

Parallel.ForEachAsync 是 .NET 6 对异步并行编程的重要补充,它填补了 Parallel.ForEach 无法正确处理 async/await 的空白。掌握它的核心机制——分区并行、信号量控并发、协作式取消——能在批量 I/O 场景下写出既高效又可控的代码。关键是根据任务类型选对工具,并合理设置 MaxDegreeOfParallelism,避免”为了并行而并行”反而拖垮系统。