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.ForEach 包 async 方法:
// ❌ 错误用法:Parallel.ForEach 不会 await 异步方法
Parallel.ForEach(urls, async url =>
{
await httpClient.GetAsync(url); // 这个 await 没人等,提前返回
});
// 上面这段代码会在请求还没发出时就"完成"
Parallel.ForEach 把 async 委托当作返回 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,而是维护固定数量的工作单元,每个工作单元从共享分区里拉取数据。这种设计避免了任务创建开销,特别适合处理海量小任务。
关键点在于:并发度不等于线程数。当 body 在 await 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,已经在执行的任务会通过 body 的 ct 参数收到取消信号。
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);
});
最佳实践总结
- 按任务类型选 API:CPU 密集用
Parallel.ForEach,I/O 密集用Parallel.ForEachAsync。 - 显式设置 DOP:I/O 场景根据下游能力设置,别用默认的
ProcessorCount。 - body 内传递 CancellationToken:所有
await都要带上ct,确保能及时取消。 - 共享状态用并发集合:
ConcurrentBag、ConcurrentDictionary,计数用Interlocked。 - 异常处理要兜底:单个任务失败不应中断整体,用
try-catch包住 body 内部逻辑。 - 配合限流和重试:调外部服务时加上 Polly 等重试策略,避免瞬时故障导致批量失败。
小结
Parallel.ForEachAsync 是 .NET 6 对异步并行编程的重要补充,它填补了 Parallel.ForEach 无法正确处理 async/await 的空白。掌握它的核心机制——分区并行、信号量控并发、协作式取消——能在批量 I/O 场景下写出既高效又可控的代码。关键是根据任务类型选对工具,并合理设置 MaxDegreeOfParallelism,避免”为了并行而并行”反而拖垮系统。