-
C#中关于代码并行化——Parallel.ForEach应用
第十部分:性能优化与内存管理
实例100:代码并行化——Parallel.ForEach应用
实例介绍
之前做电商订单批量处理功能时,处理10万条订单数据需要10分钟,通过性能分析发现,代码是串行执行的,只用到了1个CPU核心,CPU使用率只有10%。通过使用Parallel.ForEach并行化处理后,处理时vb.net教程C#教程python教程SQL教程access 2010教程间降到1分钟以内,CPU使用率提升到90%,性能提升10倍。这节就带你从零学习代码并行化,包括Parallel.ForEach基础、线程安全处理、并行度控制、性能优化,覆盖代码并行化的核心场景。
需求分析
代码并行化要解决“串行执行、CPU使用率低、处理时间长”的问题,具体需求如下:
1.并行化处理:使用Parallel.ForEach将串行代码改为并行执行;
2.线程安全:确保并行处理时的数据安全,避免竞态条件;
3.并行度控制:根据CPU核心数控制并行度,避免资源耗尽;
4.性能提升:通过并行化处理提升CPU使用率,减少处理时间;
5.异常处理:正确处理并行处理中的异常;
代码实现
前置条件:.NET 8 SDK;Visual Studio 2022;
场景1:串行执行问题——CPU使用率低,处理时间长
串行处理10万条订单数据,CPU使用率低,处理时间长。
步骤1:串行代码(SerialOrderProcessing.cs)
csharp
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
namespace CodeParallelization;
public class OrderProcessingService
{
// 串行处理订单数据,CPU使用率低,处理时间长
public void ProcessOrdersSerial(List<Order> orders)
{
var stopwatch = Stopwatch.StartNew();
// 串行执行,只用到1个CPU核心
foreach (var order in orders)
{
ProcessSingleOrder(order);
}
stopwatch.Stop();
Console.WriteLine($"串行处理完成,耗时:{stopwatch.ElapsedMilliseconds}ms,处理订单数:{orders.Count}");
}
// 模拟订单处理(CPU密集型任务)
private void ProcessSingleOrder(Order order)
{
// 模拟计算订单金额(CPU密集型操作)
order.TotalAmount = order.Items.Sum(item => item.Quantity * item.Price);
// 模拟生成订单编号(CPU密集型操作)
order.OrderNumber = $"ORD-{order.OrderId:D6}-{Guid.NewGuid():N}".Substring(0, 20);
// 模拟订单状态更新(CPU密集型操作)
order.Status = order.TotalAmount > 1000 ? OrderStatus.HighValue : OrderStatus.Normal;
// 模拟IO操作(如写入日志),这里用Thread.Sleep模拟
System.Threading.Thread.Sleep(1); // 阻塞线程1ms,模拟IO操作
}
}
public class Order
{
public int OrderId { get; set; }
public string OrderNumber { get; set; } = string.Empty;
public List<OrderItem> Items { get; set; } = new List<OrderItem>();
public decimal TotalAmount { get; set; }
public OrderStatus Status { get; set; }
}
public class OrderItem
{
public int ProductId { get; set; }
public string ProductName { get; set; } = string.Empty;
public int Quantity { get; set; }
public decimal Price { get; set; }
}
public enum OrderStatus
{
Normal,
HighValue,
Processed
}
步骤2:性能测试(Program.cs)
csharp
using System;
using System.Collections.Generic;
using System.Linq;
namespace CodeParallelization;
class Program
{
static void Main(string[] args)
{
// 生成10万条订单数据
var orders = Enumerable.Range(1, 100000).Select(i => new Order
{
OrderId = i,
Items = Enumerable.Range(1, 5).Select(j => new OrderItem
{
ProductId = j,
ProductName = $"Product {j}",
Quantity = j,
Price = j * 10.0m
}).ToList()
}).ToList();
// 串行处理
var service = new OrderProcessingService();
service.ProcessOrdersSerial(orders);
Console.ReadLine();
}
}
执行结果:
串行处理完成,耗时:110000ms(约10分钟),处理订单数:100000
CPU使用率:10%左右
场景2:并行化改造——使用Parallel.ForEach提升性能
使用Parallel.ForEach将串行代码改为并行执行,提升CPU使用率和处理速度。
步骤1:并行代码(ParallelOrderProcessing.cs)
csharp
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
namespace CodeParallelization;
public class ParallelOrderProcessingService
{
// 并行处理订单数据,提升CPU使用率和处理速度
public void ProcessOrdersParallel(List<Order> orders)
{
var stopwatch = Stopwatch.StartNew();
var processedCount = 0;
// 使用Parallel.ForEach并行处理订单
Parallel.ForEach(orders, new ParallelOptions
{
MaxDegreeOfParallelism = Environment.ProcessorCount // 并行度设置为CPU核心数
}, order =>
{
ProcessSingleOrder(order);
// 线程安全的计数(使用Interlocked.Increment)
Interlocked.Increment(ref processedCount);
});
stopwatch.Stop();
Console.WriteLine($"并行处理完成,耗时:{stopwatch.ElapsedMilliseconds}ms,处理订单数:{processedCount}");
}
// 模拟订单处理(CPU密集型任务)
private void ProcessSingleOrder(Order order)
{
// 模拟计算订单金额(CPU密集型操作)
order.TotalAmount = order.Items.Sum(item => item.Quantity * item.Price);
// 模拟生成订单编号(CPU密集型操作)
order.OrderNumber = $"ORD-{order.OrderId:D6}-{Guid.NewGuid():N}".Substring(0, 20);
// 模拟订单状态更新(CPU密集型操作)
order.Status = order.TotalAmount > 1000 ? OrderStatus.HighValue : OrderStatus.Normal;
// 模拟IO操作(如写入日志),这里用Task.Delay模拟异步IO
Task.Delay(1).Wait(); // 注意:如果是真实IO操作,应该使用异步await,避免阻塞线程
}
}
步骤2:性能测试(Program.cs)
csharp
using System;
using System.Collections.Generic;
using System.Linq;
namespace CodeParallelization;
class Program
{
static void Main(string[] args)
{
// 生成10万条订单数据
var orders = Enumerable.Range(1, 100000).Select(i => new Order
{
OrderId = i,
Items = Enumerable.Range(1, 5).Select(j => new OrderItem
{
ProductId = j,
ProductName = $"Product {j}",
Quantity = j,
Price = j * 10.0m
}).ToList()
}).ToList();
// 并行处理
var service = new ParallelOrderProcessingService();
service.ProcessOrdersParallel(orders);
Console.ReadLine();
}
}
执行结果:
并行处理完成,耗时:10000ms(约1分钟),处理订单数:100000
CPU使用率:90%左右
场景3:线程安全处理——避免竞态条件
并行处理时确保数据安全,避免竞态条件。
步骤1:错误代码(ThreadSafetyIssues.cs)
csharp
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
namespace CodeParallelization;
public class ThreadSafetyService
{
private int _processedCount = 0;
// 错误:非线程安全的计数,导致竞态条件
public void ProcessOrdersUnsafe(List<Order> orders)
{
Parallel.ForEach(orders, order =>
{
ProcessSingleOrder(order);
// 非线程安全的计数,多个线程同时修改_processedCount,导致计数错误
_processedCount++;
});
Console.WriteLine($"处理订单数:{_processedCount}"); // 结果可能小于100000
}
private void ProcessSingleOrder(Order order)
{
// 处理逻辑省略
}
}
步骤2:修复代码(ThreadSafetyFix.cs)
csharp
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
namespace CodeParallelization;
public class ThreadSafetyService
{
private int _processedCount = 0;
private readonly object _lockObject = new object();
// 正确:线程安全的计数,避免竞态条件
public void ProcessOrdersSafe(List<Order> orders)
{
Parallel.ForEach(orders, order =>
{
ProcessSingleOrder(order);
// 方法1:使用Interlocked.Increment原子操作
Interlocked.Increment(ref _processedCount);
// 方法2:使用lock关键字(适合复杂操作)
// lock (_lockObject)
// {
// _processedCount++;
// }
});
Console.WriteLine($"处理订单数:{_processedCount}"); // 结果为100000
}
private void ProcessSingleOrder(Order order)
{
// 处理逻辑省略
}
}
场景4:异常处理——正确处理并行处理中的异常
并行处理时正确捕获和处理异常。
步骤1:异常处理代码(ExceptionHandling.cs)
csharp
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
namespace CodeParallelization;
public class ExceptionHandlingService
{
public void ProcessOrdersWithExceptionHandling(List<Order> orders)
{
try
{
Parallel.ForEach(orders, order =>
{
if (order.OrderId == 50000)
{
// 模拟异常
throw new InvalidOperationException("订单处理失败:" + order.OrderId);
}
ProcessSingleOrder(order);
});
}
catch (AggregateException ex)
{
// 捕获并行处理中的所有异常
foreach (var innerEx in ex.InnerExceptions)
{
Console.WriteLine("异常:" + innerEx.Message);
}
}
}
private void ProcessSingleOrder(Order order)
{
// 处理逻辑省略
}
}
逐行讲解
场景2:Parallel.ForEach核心代码
1.Parallel.ForEach(orders, options, action):
1.orders:要并行处理的集合;
2.options:并行选项,设置MaxDegreeOfParallelism为CPU核心数;
3.action:每个元素的处理逻辑;
2.MaxDegreeOfParallelism:并行度,即同时执行的线程数,通常设置为CPU核心数;
3.Interlocked.Increment(ref processedCount):原子操作,线程安全地递增计数,避免竞态条件;
场景3:线程安全核心代码
1.Interlocked类:提供原子操作,如Increment、Decrement、Add等,适合简单的线程安全操作;
2.lock关键字:使用锁对象确保同一时间只有一个线程执行代码块,适合复杂的线程安全操作;
3.竞态条件:多个线程同时修改共享变量,导致数据不一致,需要通过原子操作或锁避免;
场景4:异常处理核心代码
1.AggregateException:并行处理中的异常会被包装在AggregateException中,需要遍历InnerExceptions获取所有异常;
2.异常捕获:在Parallel.ForEach外部使用try/catch捕获AggregateException;
基础知识拓展
-
Parallel.ForEach核心概念
并行度:同时执行的线程数,默认情况下Parallel.ForEach会根据CPU核心数自动调整并行度;
分区策略:Parallel.ForEach会将集合划分为多个分区,每个分区由一个线程处理;
线程池:Parallel.ForEach使用.NET线程池中的线程,避免频繁创建和销毁线程;
负载均衡:Parallel.ForEach会动态调整分区大小,确保每个线程的负载均衡; - 线程安全处理方法
| 方法 | 适用场景 | 性能 |
|---|---|---|
| Interlocked原子操作 | 简单的数值操作,如计数、累加 | 高,无锁开销 |
| lock关键字 | 复杂的操作,如多个变量修改、条件判断 | 中,有锁开销,但安全 |
| Concurrent集合 | 并行处理中的集合操作,如添加、删除元素 | 高,线程安全的集合类 |
| ReaderWriterLockSlim | 多读少写的场景,如配置读取、缓存更新 | 高,读操作无锁,写操作加锁 |
-
并行化适用场景
CPU密集型任务:如数据计算、图像处理、复杂逻辑判断等,并行化可以充分利用CPU核心;
大数据量处理:如批量数据导入、批量订单处理、批量报表生成等,并行化可以减少处理时间;
独立任务处理:任务之间没有依赖关系,每个任务可以独立执行; -
并行化不适用场景
IO密集型任务:如数据库操作、HTTP调用、文件读写等,并行化不会提升性能,反而会增加线程开销,应该使用异步编程;
任务之间有依赖关系:如任务A的结果需要作为任务B的输入,无法并行执行;
单核心CPU环境:并行化会增加线程切换开销,反而会降低性能; - Parallel.ForEach vs PLINQ
| 特性 | Parallel.ForEach | PLINQ |
|---|---|---|
| 适用场景 | 处理集合中的每个元素,无返回值 | 集合查询、转换,有返回值 |
| 灵活性 | 高,可以自定义处理逻辑 | 中,基于LINQ查询语法 |
| 性能 | 高,适合CPU密集型任务 | 高,适合集合查询 |
| 异常处理 | 捕获AggregateException | 捕获AggregateException |
-
并行化性能优化技巧
设置合适的并行度:根据CPU核心数和任务类型设置并行度,避免资源耗尽;
减少共享变量:尽量减少线程之间的共享变量,避免锁开销;
使用线程本地存储:使用ThreadLocal存储线程本地数据,避免共享变量;
避免线程阻塞:并行处理中避免线程阻塞,如同步IO操作,应该使用异步编程;
总结
代码并行化的核心是充分利用CPU核心,将串行任务改为并行执行,提升处理速度。关键要点:
1.Parallel.ForEach使用:掌握Parallel.ForEach的基本用法和并行度控制;
2.线程安全:使用原子操作、锁或Concurrent集合确保并行处理中的数据安全;
3.异常处理:正确捕获和处理AggregateException中的异常;
4.适用场景:并行化适合CPU密集型任务和大数据量处理,不适合IO密集型任务;
5.性能优化:设置合适的并行度,减少共享变量,避免线程阻塞;
比如这个代码并行化,使用Parallel.ForEach后,处理时间从10分钟降到1分钟以内,CPU使用率从10%提升到90%,性能提升10倍。掌握代码并行化,你就能充分利用CPU资源,提升系统处理速度!
本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49498.html










