简介
Parallel 并行编程是 .NET 中利用多核 CPU 进行并发执行的编程模型,主要通过 System.Threading.Tasks 命名空间中的 Parallel 类实现。它允许将任务分解成多个子任务,在多个线程上同时执行,以加速 CPU 密集型操作(如循环计算、数据处理)。
核心组件:
-
Parallel类:提供静态方法如Parallel.For、Parallel.ForEach、Parallel.Invoke,用于并行执行循环或方法。 -
PLINQ(Parallel LINQ):LINQ的并行版本,通过AsParallel()扩展方法启用并行查询。 -
底层依赖:基于 Task Parallel Library (TPL),利用线程池(
ThreadPool)管理线程,避免手动创建线程的开销。
关键概念:
-
并行度(
Degree of Parallelism):控制并发线程数,默认基于 CPU 核心(e.g., 4 核 CPU 可能 4 线程)。 -
数据并行:将数据分成块(
chunks),每个线程处理一块(e.g.,Parallel.For)。 -
任务并行:同时执行独立方法(e.g.,
Parallel.Invoke)。
Parallel 的核心定位与价值
Parallel 类位于 System.Threading.Tasks 命名空间,是 .NET 提供的 “高层级并行工具”,核心价值在于简化多线程开发的复杂性。它屏蔽了底层线程管理的繁琐细节,如线程创建、销毁及同步原语的使用,使开发者能够专注于业务逻辑本身。通过声明式 API,开发者只需指定任务集合,框架会自动处理负载均衡、线程调度及异常处理。这种抽象不仅提高了代码的可读性和可维护性,还显著降低了因并发错误导致的 Bug 概率。此外,Parallel 与 .NET 生态系统深度集成,能够无缝配合异步编程模型(async/await)及 LINQ 查询,为现代高性能应用开发提供了坚实基础。它特别适用于那些计算密集且任务相互独立的场景,通过最大化硬件利用率,实现性能的大幅提升。
Parallel 的优势
自动线程调度:无需手动创建或管理线程,由
TPL自动分配线程池线程,实现负载均衡;简化并行逻辑:使用类似串行循环的语法实现并行执行,显著降低并行编程的门槛;
适配多核硬件:默认根据
CPU核心数动态调整并行度,最大化利用硬件资源;支持取消与超时:可通过
ParallelOptions精细控制并行过程,例如取消执行或限制并行数量。
注意:Parallel 仅适合 CPU 密集型任务(如数据计算、图像处理、复杂逻辑运算)。对于 I/O 密集型任务(如文件读写、网络请求),应使用 async/await,否则线程阻塞会浪费资源。
核心概念与基础
并行 vs. 并发 vs. 异步
并发:多个任务在重叠的时间段内执行(不一定同时)。在单核处理器上,通过时间片切换来实现并发。
并行:多个任务真正同时执行,这需要多核或多处理器的硬件支持。并行实际上是并发的一种特例。
异步:一种编程模式,允许启动一个操作后不阻塞当前线程,待操作完成后再处理结果。异步操作可以利用并发或并行(通常通过线程池),但其核心目标是非阻塞和高响应性。
线程 vs. 任务
Parallel 核心 API 一览
| API | 用途 |
|---|---|
| Parallel.For | 并行 for 循环 |
| Parallel.ForEach | 并行 foreach |
| Parallel.Invoke | 并行执行多个 Action |
| ParallelOptions | 控制并行度、取消等 |
Parallel.Invoke:并行执行多个独立任务
该方法用于一次性并行执行多个无关联的方法或任务,非常适合“多任务并行执行,等待全部完成”的场景。它简化了并发控制的复杂性,让开发者无需手动管理线程同步即可实现并行处理。
语法:
public static void Invoke(params Action[] actions);
public static void Invoke(ParallelOptions options, params Action[] actions);
示例:并行执行三个独立的 CPU 密集型方法
using System;
using System.Threading.Tasks;
class ParallelInvokeDemo
{
static void Main()
{
// 记录开始时间
var watch = System.Diagnostics.Stopwatch.StartNew();
// 并行执行三个方法
Parallel.Invoke(
() => CalculateSum(1, 100000000), // 任务1:计算1~1亿的和
() => CalculatePrimeCount(1, 100000), // 任务2:统计1~10万的质数数量
() => GenerateRandomData(1000000) // 任务3:生成100万条随机数据
);
watch.Stop();
Console.WriteLine($"并行执行耗时:{watch.ElapsedMilliseconds}ms");
// 对比:串行执行(耗时远高于并行)
watch.Restart();
CalculateSum(1, 100000000);
CalculatePrimeCount(1, 100000);
GenerateRandomData(1000000);
watch.Stop();
Console.WriteLine($"串行执行耗时:{watch.ElapsedMilliseconds}ms");
}
// 模拟CPU密集型任务1:计算累加和
static void CalculateSum(int start, int end)
{
long sum = 0;
for (long i = start; i <= end; i++) sum += i;
Console.WriteLine($"累加和:{sum}");
}
// 模拟CPU密集型任务2:统计质数数量
static void CalculatePrimeCount(int start, int end)
{
int count = 0;
for (int i = start; i <= end; i++)
{
if (IsPrime(i)) count++;
}
Console.WriteLine($"质数数量:{count}");
}
// 模拟CPU密集型任务3:生成随机数据
static void GenerateRandomData(int count)
{
var random = new Random();
double[] data = new double[count];
for (int i = 0; i < count; i++) data[i] = random.NextDouble();
Console.WriteLine($"随机数据生成完成,长度:{data.Length}");
}
// 辅助方法:判断是否为质数
static bool IsPrime(int num)
{
if (num < 2) return false;
for (int i = 2; i <= Math.Sqrt(num); i++)
{
if (num % i == 0) return false;
}
return true;
}
}
关键说明:
-
执行顺序不确定:虽然所有任务会被提交到线程池并行执行,但具体的执行顺序是不确定的。开发者不应依赖任何特定的执行顺序,除非任务之间没有数据依赖关系。
-
阻塞等待:Parallel.Invoke 是一个同步方法,它会阻塞当前线程,直到所有传入的操作委托全部执行完毕。这意味着调用线程会一直等待,直到并行任务完成,这与异步编程模型中的非阻塞特性不同。
-
异常处理:如果在并行执行过程中,任何一个任务抛出未处理的异常,Parallel.Invoke 会将这些异常封装在 AggregateException 中抛出。调用者需要捕获并检查 InnerExceptions 属性,以获取具体的错误信息。这确保了即使部分任务失败,其他任务的结果也不会被静默忽略。
-
线程池调度:该方法内部由 TaskScheduler 管理,默认使用线程池线程。这意味着它会自动利用线程池的优化机制,如线程复用和负载平衡,从而避免频繁创建和销毁线程带来的性能开销。
-
适用场景限制:仅适用于相互独立、无共享状态修改的任务。如果任务之间存在数据竞争或需要共享可变状态,必须使用锁或其他同步原语来保护共享资源,否则可能导致数据不一致或死锁。
Parallel.For:并行执行 for 循环
该方法旨在替代传统的 for 循环,通过将循环迭代分配到多个线程并行执行,从而显著提升性能。它特别适用于“固定次数的循环,且迭代间无依赖”的场景,能够有效利用多核处理器的计算能力。
核心语法如下:
// 基础版:从fromInclusive到toExclusive(不包含)的并行循环
public static ParallelLoopResult For(
int fromInclusive,
int toExclusive,
Action<int> body
);
// 带配置版:支持取消、限制并行度
public static ParallelLoopResult For(
int fromInclusive,
int toExclusive,
ParallelOptions options,
Action<int> body
);
以下是一个并行累加数组元素的示例,展示了如何确保线程安全:
using System;
using System.Threading;
using System.Threading.Tasks;
class ParallelForDemo
{
static void Main()
{
// 初始化1000万个元素的数组
int[] numbers = new int[10_000_000];
Random random = new Random();
for (int i = 0; i < numbers.Length; i++) numbers[i] = random.Next(1, 100);
long total = 0; // 共享累加变量(需保证线程安全)
var watch = System.Diagnostics.Stopwatch.StartNew();
// 并行循环累加
ParallelLoopResult result = Parallel.For(
0, // 起始索引(包含)
numbers.Length, // 结束索引(不包含)
// 循环体:i为当前迭代索引
(i) => Interlocked.Add(ref total, numbers[i]) // 用Interlocked保证原子累加
);
watch.Stop();
Console.WriteLine($"并行累加结果:{total}");
Console.WriteLine($"耗时:{watch.ElapsedMilliseconds}ms");
Console.WriteLine($"循环是否完成:{result.IsCompleted}");
}
}
在使用 Parallel.For 时,需注意以下关键说明:
-
Parallel.For的迭代索引i是线程局部的,这意味着每个线程拥有独立的索引副本,因此无需担心索引冲突。然而,对于共享变量(如total),必须保证线程安全,通常需使用Interlocked、lock或ConcurrentBag等同步机制来保护共享数据,防止数据竞争。 -
方法的返回值
ParallelLoopResult包含了循环执行的详细状态信息。其中,IsCompleted属性指示循环是否全部完成,而LowestBreakIteration属性则表明是否因调用 Stop 或 Break 而提前中断。开发者可根据这些状态进行后续处理或错误排查。 -
迭代间绝对不能存在依赖关系。例如,第
i次迭代的结果不应依赖于第i-1次迭代的输出。如果存在此类依赖,并行执行将导致结果不可预测或错误,此时应考虑使用串行循环或调整算法逻辑。
Parallel.ForEach:并行执行 foreach 循环
Parallel.ForEach 用于替代传统的 foreach 循环,它能够遍历 IEnumerable 集合,并将集合中的元素分配到多个线程进行并行处理。这种模式非常适合处理大量独立数据项的场景,能够大幅缩短处理时间。
核心语法如下:
// 基础版:遍历IEnumerable集合
public static ParallelLoopResult ForEach<TSource>(
IEnumerable source,
Action body
);
// 带索引版:获取元素的索引
public static ParallelLoopResult ForEach<TSource>(
IEnumerable source,
Actionlong > body
);
以下是一个并行处理文件列表的示例,特别适用于 CPU 密集型的文件内容解析任务:
using System;
using System.Collections.Generic;
using System.IO;
using System.Threading.Tasks;
class ParallelForEachDemo
{
static void Main()
{
// 获取指定目录下的所有文本文件
string[] files = Directory.GetFiles(@"D:test", "*.txt");
// 存储解析结果(线程安全集合)
var parseResults = new System.Collections.Concurrent.ConcurrentDictionary<string, int>();
var watch = System.Diagnostics.Stopwatch.StartNew();
// 并行遍历文件列表
Parallel.ForEach(
files, // 要遍历的集合
(file, state, index) => // 循环体:file=当前文件,state=循环状态,index=当前索引
{
try
{
// 解析文件:统计文件中的数字数量(CPU密集型)
int numberCount = CountNumbersInFile(file);
// 将结果存入线程安全字典
parseResults.TryAdd(file, numberCount);
Console.WriteLine($"已处理第{index+1}个文件:{file},数字数量:{numberCount}");
}
catch (Exception ex)
{
Console.WriteLine($"处理文件{file}失败:{ex.Message}");
// 可选:终止所有迭代
// state.Stop();
}
}
);
watch.Stop();
Console.WriteLine($"n全部处理完成,共{parseResults.Count}个文件,耗时:{watch.ElapsedMilliseconds}ms");
}
// 模拟CPU密集型任务:统计文件中的数字数量
static int CountNumbersInFile(string filePath)
{
string content = File.ReadAllText(filePath);
int count = 0;
foreach (char c in content)
{
if (char.IsDigit(c)) count++;
}
// 模拟复杂计算(放大CPU消耗)
for (int i = 0; i < 100000; i++) { Math.Sqrt(i); }
return count;
}
}
关键说明:
-
ParallelLoopState:用于控制循环(Stop()终止所有迭代、Break()终止后续迭代、IsStopped判断是否终止); -
推荐使用线程安全集合(如
ConcurrentDictionary、ConcurrentBag)存储并行处理的结果,避免共享集合的线程安全问题; -
若集合元素数量少、循环体执行时间极短,并行开销可能超过收益,此时应使用串行循环。
ParallelOptions:配置并行行为
ParallelOptions 用于自定义并行执行的规则,核心属性:
MaxDegreeOfParallelism | 限制最大并行度(线程数),默认值为-1(自动适配 CPU 核心数),可设置为具体数值(如4表示最多 4 个线程并行) |
|---|---|
CancellationToken | 取消令牌,用于取消并行执行 |
TaskScheduler | 指定任务调度器(默认使用线程池调度器) |
示例:限制并行度 + 取消并行执行
using System;
using System.Threading;
using System.Threading.Tasks;
class ParallelOptionsDemo
{
static void Main()
{
CancellationTokenSource cts = new CancellationTokenSource();
// 5秒后取消执行
cts.CancelAfter(5000);
ParallelOptions options = new ParallelOptions
{
MaxDegreeOfParallelism = 4, // 最多4个线程并行
CancellationToken = cts.Token // 绑定取消令牌
};
try
{
Parallel.For(
0,
1000000,
options,
(i) =>
{
// 模拟耗时操作
Thread.Sleep(10);
if (i % 100000 == 0) Console.WriteLine($"已处理{i}次");
// 检查取消令牌(可选,TPL会自动检查,但手动检查更及时)
options.CancellationToken.ThrowIfCancellationRequested();
}
);
}
catch (OperationCanceledException)
{
Console.WriteLine("并行执行被取消");
}
catch (AggregateException ex)
{
foreach (var innerEx in ex.InnerExceptions)
{
Console.WriteLine($"异常:{innerEx.Message}");
}
}
finally
{
cts.Dispose();
}
}
}
并发度说明:
-
默认:≈ CPU 核心数
-
并不是越大越好
-
过大 → 线程切换成本上升
CPU 密集型 ≈ 核心数
轻计算 ≈ 核心数 × 1.5
PLINQ(Parallel LINQ)
并行查询:
var query = from num in data.AsParallel() // 启用并行
where num % 2 == 0
select num * 2;
var results = query.ToArray(); // 强制执行
-
选项:
WithDegreeOfParallelism(4)控制并行度;WithExecutionMode(ParallelExecutionMode.ForceParallelism)强制并行。 -
有序 vs. 无序:默认无序;用
AsOrdered()保持顺序(性能稍低)。
底层原理:Parallel 的执行机制
Parallel 的高效性源于 TPL 的核心设计:
Parallel 与线程安全
错误示例:共享变量
int sum = 0;
Parallel.For(0, 1000, i =>
{
sum += i; // 线程不安全
});
结果:不确定
正确方式一:Interlocked
int sum = 0;
Parallel.For(0, 1000, i =>
{
Interlocked.Add(ref sum, i);
});
正确方式二:局部变量 + 聚合(推荐)
int sum = 0;
Parallel.For(0, 1000,
() => 0,
(i, state, local) => local + i,
local => Interlocked.Add(ref sum, local)
);
高性能、低竞争
替代方案:
-
小任务:用
Task.Run或Parallel LINQ。 -
数据并行:
System.Numerics (SIMD)或GPU (CUDA.NET)。 -
高并发:
Actor模型 (Akka.NET) 或Channels。
Parallel.ForEachAsync
在受控并发度下,并行执行异步操作
它本质是:
-
async / await友好 -
支持并发限制
-
内置
CancellationToken -
自动调度,不用手写
SemaphoreSlim
基本用法
最简单示例
-
必须
await -
lambda参数里有CancellationToken -
返回
ValueTask
限制并发度
这相当于:
但 更简洁、更安全
Parallel.ForEachAsync vs Task.WhenAll
Task.WhenAll(无并发控制)
特点:
-
一次性创建所有
Task -
不限制并发
-
数量大 → 容易压垮资源
Parallel.ForEachAsync(有并发控制)
特点:
-
滑动窗口式并发
-
同时最多
N个任务 -
更适合:
-
HTTP
-
DB
-
文件 IO
-
调用第三方接口
-
什么时候选哪个?
典型实战场景
批量调用第三方 API
批量文件处理(IO)
await Parallel.ForEachAsync(items, async (item, ct) =>
{
await ProcessAsync(item, ct);
});
var options = new ParallelOptions
{
MaxDegreeOfParallelism = 5
};
await Parallel.ForEachAsync(items, options, async (item, ct) =>
{
await CallApiAsync(item, ct);
});
SemaphoreSlim(5) + Task.WhenAll
await Task.WhenAll(items.Select(item =>
ProcessAsync(item)
));
await Parallel.ForEachAsync(items,
new ParallelOptions { MaxDegreeOfParallelism = 5 },
async (item, ct) =>
{
await ProcessAsync(item, ct);
});
| 场景 | 推荐 |
|---|---|
| 少量任务(<20) | Task.WhenAll |
| 大量任务 | Parallel.ForEachAsync |
| 需要限流 | Parallel.ForEachAsync |
| CPU 密集 | Parallel.For / PLINQ |
await Parallel.ForEachAsync(userIds,
new ParallelOptions { MaxDegreeOfParallelism = 10 },
async (id, ct) =>
{
await apiClient.SyncUserAsync(id, ct);
});
await Parallel.ForEachAsync(files,
async (file, ct) =>
{
var content = await File.ReadAllTextAsync(file, ct);
await SaveAsync(content, ct);
});
CPU
Thread
TPL
Task
Task
TaskScheduler
ContinueWith, WhenAll, WhenAny
AggregateException
CancellationTokenSource/CancellationToken
Created, Running, RanToCompletion, Canceled, Faulted
.NET
TPL
ThreadPool.QueueUserWorkItem
Task.Run/Task.Factory.StartNew
Parallel.Invoke
Action
TPL
AggregateException
Partitioner
Work-Stealing
PLINQ
ParallelQuery
Barrier
AggregateException
InnerExceptions









