在 .NET 的异步编程中,Task.WhenAll 是一个核心方法,用于并行执行多个异步任务并等待所有任务完成,特别适合处理 I/O 密集型任务(如网络请求、文件操作)和硬件通讯等场景
在 .NET 的异步编程中,Task.WhenAll 是一个核心方法,用于并行执行多个异步任务并等待所有任务完成,特别适合处理 I/O 密集型任务(如网络请求、文件操作)和硬件通讯等场景。
在 WinForms 应用程序中,Task.WhenAll 可以与并发工具(如 SemaphoreSlim)、并发集合(如 ConcurrentBag<T>, ConcurrentQueue<T>, ConcurrentDictionary<TKey, TValue>)和原子操作(如 Interlocked)结合,优化性能、保持 UI 响应性,并确保线程安全。
本节将深入分析 Task.WhenAll 的定义、机制、优缺点,与其他并发机制(如 Task.WhenAny、锁、Parallel.For)进行比较,提供 WinForms 环境下的详细代码示例,明确应用场景,并确保高效、稳定、健壮。
内容涵盖线程安全、性能优化、异常处理、资源管理,解决 SemaphoreFullException 和异步死锁等问题,适合需要深入理解 Task.WhenAll 的开发者。
1. Task.WhenAll 深入分析
1.1 定义与机制
- 定义:Task.WhenAll 是 .NET 异步编程模型(TAP)中的静态方法,接受一组 Task 或 Task<T>,返回一个 Task 或 Task<T[]>,在所有输入任务完成后完成。它用于并行执行多个异步操作并等待其结果。
- 机制:
- 输入:接受 IEnumerable<Task> 或 Task[],可以是无返回值的 Task 或有返回值的 Task<T>。
- 输出:
- 如果输入是 Task[],返回 Task,表示所有任务完成。
- 如果输入是 Task<T>[],返回 Task<T[]>,包含所有任务的结果。
- 执行:Task.WhenAll 不启动任务,仅等待已启动的任务完成。任务在调用时通常已通过 Task.Run 或异步方法(如 HttpClient.GetStringAsync)启动。
- 并发性:任务并行执行,依赖线程池或底层异步 I/O(如 IOCP)。
- 异常处理:如果任一任务失败,Task.WhenAll 将抛出 AggregateException,包含所有异常。
- 关键 API:
- Task.WhenAll(IEnumerable<Task>):等待无返回值任务。
- Task.WhenAll(IEnumerable<Task<T>>):等待有返回值任务,返回结果数组。
- 结合 CancellationToken 支持取消。
- 关键特性:
- 并行执行:所有任务并发运行,适合 I/O 密集型任务。
- 非阻塞:await Task.WhenAll 释放调用线程(如 UI 线程)。
- 结果收集:自动收集 Task<T> 的结果,简化批量处理。
- 异常聚合:统一处理多个任务的异常。
1.2 优点
- 高性能:并行执行多个任务,减少总等待时间,适合 I/O 密集型任务。
- UI 响应性:结合 async/await,避免阻塞 WinForms 主线程。
- 简单性:简化多任务等待逻辑,代码清晰。
- 结果收集:自动收集 Task<T> 结果,适合批量数据处理。
- 取消支持:通过 CancellationToken 支持统一取消。
- 与并发集合集成:结合 ConcurrentDictionary 或 ConcurrentBag 存储结果,线程安全。
1.3 缺点
- 等待所有任务:必须等待所有任务完成,不适合需要部分结果的场景(对比 Task.WhenAny)。
- 异常复杂性:抛出 AggregateException,需解析多个异常。
- 资源消耗:高并发可能导致资源竞争,需结合 SemaphoreSlim 限制。
- 不适合 CPU 密集型:Task.WhenAll 更适合 I/O 密集型任务,CPU 密集型需结合 Task.Run。
- 调试难度:并发任务的调试可能复杂,需日志支持。
1.4 与其他并发机制的比较以下是 Task.WhenAll 与 Task.WhenAny、锁(lock 语句)、信号量(SemaphoreSlim)、原子操作和 Parallel.For 的比较,结合 WinForms 和并发场景:
|
特性 |
Task.WhenAll |
Task.WhenAny |
锁(lock 语句) |
信号量(SemaphoreSlim) |
原子操作 |
Parallel.For |
|---|---|---|---|---|---|---|
|
实现方式 |
等待所有异步任务完成 |
等待任一异步任务完成 |
互斥锁(Monitor) |
计数器限制并发 |
硬件指令或细粒度锁 |
并行循环 |
|
线程安全 |
是(结合并发集合/原子操作) |
是(结合并发集合/原子操作) |
是(互斥访问) |
是(限制并发任务数) |
是(硬件或集合级保证) |
是(内部管理) |
|
性能 |
高(并行异步,线程池优化) |
高(仅等待一个任务) |
中等(锁竞争严重时性能下降) |
中等(计数器管理有开销) |
高(无锁或细粒度锁) |
高(并行 CPU 密集型) |
|
复杂性 |
中等(异常处理、状态机) |
中等(需手动迭代) |
高(需手动管理锁) |
中等(需管理计数器) |
低(简单 API) |
中等(需配置选项) |
|
适用场景 |
批量 I/O 任务(如多 API 请求) |
快速响应任一任务 |
复杂多步骤操作 |
限制并发任务数 |
简单操作(如计数、键值更新) |
CPU 密集型并行计算 |
|
死锁风险 |
有(误用 Task.Result) |
有(误用 Task.Result) |
有(需避免死锁) |
低(需正确释放) |
无 |
低(内部管理) |
|
UI 响应性 |
优秀(异步非阻塞) |
优秀(异步非阻塞) |
差(锁可能阻塞) |
中等(并发控制) |
优秀(无锁) |
差(可能阻塞) |
|
与并发集合 |
优秀(结果收集) |
适中(单任务结果) |
有限(锁降低并行性) |
优秀(控制并发) |
优秀(高效存储) |
适中(需手动收集) |
|
典型 WinForms 场景 |
批量 API 请求、文件读写 |
快速响应(如竞态 API) |
复杂状态管理 |
限制 API 请求并发 |
计数、键值更新 |
数据处理、图像计算 |
分析:
- Task.WhenAll:适合批量异步任务(如多 API 请求、文件读写),并行执行,收集所有结果。
- Task.WhenAny:适合快速响应场景(如竞态请求),仅等待一个任务完成。
- 锁(lock 语句):适合复杂多步骤操作,但锁竞争和死锁风险高。
- 信号量(SemaphoreSlim):适合限制并发任务数,结合 Task.WhenAll 控制资源。
- 原子操作:适合简单计数或键值更新,性能高,无死锁风险。
- Parallel.For:适合 CPU 密集型并行计算,需结合 Task.Run 用于异步。
选择建议:
- 批量 I/O 任务:选择 Task.WhenAll(如多 API 请求、文件操作),结合 SemaphoreSlim 限制并发。
- 快速响应:选择 Task.WhenAny(如选择最快 API 响应)。
- 复杂多步骤操作:使用 lock 语句。
- 简单计数/键值更新:使用原子操作(如 Interlocked、ConcurrentDictionary)。
- CPU 密集型并行:使用 Parallel.For。
2. Task.WhenAll 使用场景与代码示例以下针对前文提到的场景(实时数据刷新、批量数据同步、I/O 密集型任务、硬件通讯),提供结合 Task.WhenAll、异步编程(async/await)、SemaphoreSlim、并发集合(ConcurrentQueue<T>、ConcurrentDictionary<TKey, TValue>)和原子操作(Interlocked)的 WinForms 代码示例,融入最佳实践(线程安全、并发限制、动态调整、取消支持、UI 更新、资源管理、避免 SemaphoreFullException 和死锁)。
2.1 实时数据刷新(多股票价格刷新)场景:每秒从多个 API 获取股票价格,限制并发,动态调整并发,缓存结果,批量更新 UI。
- 适用方法:Task.WhenAll + async/await + SemaphoreSlim + ConcurrentDictionary + Interlocked.
- 最佳实践:
- 使用 Task.WhenAll 并行执行多个异步 API 请求。
- ConcurrentDictionary 缓存价格,AddOrUpdate 原子更新。
- Interlocked 跟踪请求计数。
- 检查 SemaphoreSlim.CurrentCount 避免 SemaphoreFullException.
代码示例:
using System;
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Net.Http;
using System.Threading;
using System.Threading.Tasks;
using System.Windows.Forms;
public partial class MainForm : Form
{
private readonly System.Windows.Forms.Timer _timer = new System.Windows.Forms.Timer();
private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(3, 6);
private readonly ConcurrentDictionary<string, string> _priceCache = new ConcurrentDictionary<string, string>();
private readonly CancellationTokenSource _cts = new CancellationTokenSource();
private readonly HttpClient _httpClient = new HttpClient();
private int _currentConcurrency = 3;
private readonly int _maximumCount = 6;
private long _requestCount = 0; // 原子计数器
public MainForm()
{
InitializeComponent();
_timer.Interval = 1000; // 1秒
_timer.Tick += async (s, e) => await UpdateStockPricesAsync(_cts.Token);
_timer.Start();
}
private async Task UpdateStockPricesAsync(CancellationToken token)
{
var stockCodes = new[] { "stock1", "stock2", "stock3" };
var tasks = new ConcurrentBag<Task>();
foreach (var stockCode in stockCodes)
{
bool acquired = await _semaphore.WaitAsync(0, token);
if (!acquired)
{
Debug.WriteLine($"跳过任务 {stockCode}:信号量已满");
continue;
}
tasks.Add(Task.Run(async () =>
{
try
{
Interlocked.Increment(ref _requestCount); // 原子递增
var data = await _httpClient.GetStringAsync($"https://api.example.com/stock/{stockCode}", token); // 异步请求
_priceCache.AddOrUpdate(stockCode, data, (k, v) => data); // 原子更新
}
catch (HttpRequestException ex)
{
Debug.WriteLine($"任务 {stockCode} 失败: {ex.Message}, 当前计数: {_semaphore.CurrentCount}");
}
finally
{
if (acquired)
{
_semaphore.Release();
Debug.WriteLine($"释放信号量,当前计数: {_semaphore.CurrentCount}");
}
}
}, token));
}
try
{
var stopwatch = Stopwatch.StartNew();
await Task.WhenAll(tasks); // 等待所有任务完成
BeginInvoke((Action)(() => stockPriceLabel.Text = $"{string.Join("\n", _priceCache.Values)} (请求: {Interlocked.Read(ref _requestCount)})"));
// 动态调整并发
var successRate = (double)_priceCache.Count / stockCodes.Length;
if (successRate < 0.5 && _currentConcurrency > 1)
{
Interlocked.Decrement(ref _currentConcurrency);
_timer.Interval = (int)(_timer.Interval * 1.2);
Debug.WriteLine($"成功率低 ({successRate:P0}),减少并发到 {_currentConcurrency}");
}
else if (successRate > 0.9 && _currentConcurrency < _maximumCount)
{
if (_semaphore.CurrentCount < _maximumCount)
{
_semaphore.Release();
Interlocked.Increment(ref _currentConcurrency);
_timer.Interval = Math.Max(500, (int)(_timer.Interval / 1.2));
Debug.WriteLine($"成功率高 ({successRate:P0}),增加并发到 {_currentConcurrency}, 当前计数: {_semaphore.CurrentCount}");
}
}
}
catch (OperationCanceledException)
{
BeginInvoke((Action)(() => stockPriceLabel.Text = "刷新取消"));
}
catch (AggregateException ex)
{
Debug.WriteLine($"批量错误: {ex.InnerException?.Message}, 当前计数: {_semaphore.CurrentCount}");
BeginInvoke((Action)(() => stockPriceLabel.Text = $"错误: {ex.InnerException?.Message}"));
}
}
private void CancelButton_Click(object sender, EventArgs e)
{
_cts.Cancel();
_cts = new CancellationTokenSource();
}
protected override void OnFormClosing(FormClosingEventArgs e)
{
base.OnFormClosing(e);
_timer.Stop();
_cts.Cancel();
_cts.Dispose();
_semaphore.Dispose();
_httpClient.Dispose();
}
}
说明:
- 适用性:实时股票价格、传感器数据批量刷新。
- Task.WhenAll:并行执行多个 API 请求,等待所有完成。
- 异步编程:HttpClient.GetStringAsync 异步请求,await Task.WhenAll 释放 UI 线程。
- 原子操作:Interlocked.Increment 计数,ConcurrentDictionary.AddOrUpdate 缓存。
- 最佳实践:
- Task.WhenAll 高效并行处理,批量收集结果。
- ConcurrentDictionary 高效缓存,Interlocked 线程安全计数。
- 检查 acquired 和 CurrentCount 避免 SemaphoreFullException.
- 复用 HttpClient,批量 UI 更新。
- 稳定性:捕获 HttpRequestException 和 AggregateException,finally 块释放。
- 效率:动态调整并发,减少 UI 更新频率。
- Task.WhenAll 优势:并行执行多任务,简化批量等待逻辑。
2.2 批量数据同步(多 API 同步)场景:每 5 秒同步多个 API,限制并发,动态调整并发,缓存结果,批量更新 UI。
- 适用方法:Task.WhenAll + async/await + SemaphoreSlim + ConcurrentDictionary + Interlocked.
- 最佳实践:
- 使用 Task.WhenAll 并行执行多个异步 API 请求。
- ConcurrentDictionary 存储结果,AddOrUpdate 原子更新。
- Interlocked 跟踪成功计数。
代码示例:csharp
using System;
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Net.Http;
using System.Threading;
using System.Threading.Tasks;
using System.Windows.Forms;
public partial class MainForm : Form
{
private readonly System.Windows.Forms.Timer _timer = new System.Windows.Forms.Timer();
private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(3, 6);
private readonly CancellationTokenSource _cts = new CancellationTokenSource();
private readonly HttpClient _httpClient = new HttpClient();
private readonly ConcurrentDictionary<string, string> _results = new ConcurrentDictionary<string, string>();
private int _currentConcurrency = 3;
private readonly int _maximumCount = 6;
private long _successCount = 0; // 原子计数器
public MainForm()
{
InitializeComponent();
_timer.Interval = 5000; // 5秒
_timer.Tick += async (s, e) => await SyncDataAsync(_cts.Token);
_timer.Start();
}
private async Task SyncDataAsync(CancellationToken token)
{
var endpoints = new[] { "data1", "data2", "data3", "data4" };
var tasks = new ConcurrentBag<Task>();
foreach (var endpoint in endpoints)
{
bool acquired = await _semaphore.WaitAsync(1000, token);
if (!acquired)
{
Debug.WriteLine($"跳过任务 {endpoint}:超时或取消");
continue;
}
tasks.Add(Task.Run(async () =>
{
try
{
var data = await _httpClient.GetStringAsync($"https://api.example.com/{endpoint}", token); // 异步请求
_results.AddOrUpdate(endpoint, data, (k, v) => data); // 原子更新
Interlocked.Increment(ref _successCount); // 原子递增
}
catch (HttpRequestException ex)
{
Debug.WriteLine($"任务 {endpoint} 失败: {ex.Message}, 当前计数: {_semaphore.CurrentCount}");
}
finally
{
if (acquired)
{
_semaphore.Release();
Debug.WriteLine($"释放信号量,当前计数: {_semaphore.CurrentCount}");
}
}
}, token));
}
try
{
var stopwatch = Stopwatch.StartNew();
await Task.WhenAll(tasks); // 等待所有任务完成
BeginInvoke((Action)(() => dataLabel.Text = $"{string.Join("\n", _results.Values)} (成功: {Interlocked.Read(ref _successCount)})"));
// 动态调整并发
var successRate = (double)Interlocked.Read(ref _successCount) / endpoints.Length;
if (successRate < 0.5 && _currentConcurrency > 1)
{
Interlocked.Decrement(ref _currentConcurrency);
_timer.Interval = (int)(_timer.Interval * 1.2);
Debug.WriteLine($"成功率低 ({successRate:P0}),减少并发到 {_currentConcurrency}");
}
else if (successRate > 0.9 && _currentConcurrency < _maximumCount)
{
if (_semaphore.CurrentCount < _maximumCount)
{
_semaphore.Release();
Interlocked.Increment(ref _currentConcurrency);
_timer.Interval = Math.Max(2000, (int)(_timer.Interval / 1.2));
Debug.WriteLine($"成功率高 ({successRate:P0}),增加并发到 {_currentConcurrency}, 当前计数: {_semaphore.CurrentCount}");
}
}
}
catch (OperationCanceledException)
{
BeginInvoke((Action)(() => dataLabel.Text = "同步取消"));
}
catch (AggregateException ex)
{
Debug.WriteLine($"批量错误: {ex.InnerException?.Message}, 当前计数: {_semaphore.CurrentCount}");
BeginInvoke((Action)(() => dataLabel.Text = $"错误: {ex.InnerException?.Message}"));
}
}
private void CancelButton_Click(object sender, EventArgs e)
{
_cts.Cancel();
_cts = new CancellationTokenSource();
}
protected override void OnFormClosing(FormClosingEventArgs e)
{
base.OnFormClosing(e);
_timer.Stop();
_cts.Cancel();
_cts.Dispose();
_semaphore.Dispose();
_httpClient.Dispose();
}
}
说明:
- 适用性:批量 API 请求、数据库同步。
- Task.WhenAll:并行执行多个 API 请求,等待所有完成。
- 异步编程:HttpClient.GetStringAsync 异步请求,await Task.WhenAll 释放 UI 线程。
- 原子操作:Interlocked.Increment 计数,ConcurrentDictionary.AddOrUpdate 缓存。
- 最佳实践:
- Task.WhenAll 高效并行处理,批量收集结果。
- ConcurrentDictionary 高效存储,Interlocked 线程安全计数。
- 检查 acquired 和 CurrentCount 避免 SemaphoreFullException.
- 复用 HttpClient,批量 UI 更新。
- 稳定性:捕获 HttpRequestException 和 AggregateException,finally 块释放。
- 效率:动态调整并发,异步并行处理。
- Task.WhenAll 优势:简化多任务并行等待,适合批量同步。
2.3 I/O 密集型任务(批量文件读写)场景:每 5 秒读取多个文件,限制并发,动态调整并发,缓存结果,批量更新 UI。
- 适用方法:Task.WhenAll + async/await + SemaphoreSlim + ConcurrentQueue + ConcurrentDictionary + Interlocked.
- 最佳实践:
- 使用 Task.WhenAll 并行读取多个文件。
- ConcurrentQueue 管理任务队列,ConcurrentDictionary 缓存结果。
- Interlocked 跟踪成功计数。
代码示例:csharp
using System;
using System.Collections.Concurrent;
using System.Diagnostics;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using System.Windows.Forms;
public partial class MainForm : Form
{
private readonly System.Windows.Forms.Timer _timer = new System.Windows.Forms.Timer();
private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(3, 6);
private readonly ConcurrentQueue<string> _fileQueue = new ConcurrentQueue<string>();
private readonly ConcurrentDictionary<string, string> _fileCache = new ConcurrentDictionary<string, string>();
private readonly CancellationTokenSource _cts = new CancellationTokenSource();
private int _currentConcurrency = 3;
private readonly int _maximumCount = 6;
private long _successCount = 0; // 原子计数器
public MainForm()
{
InitializeComponent();
_timer.Interval = 5000; // 5秒
_timer.Tick += async (s, e) => await ProcessFilesAsync(_cts.Token);
_timer.Start();
_fileQueue.Enqueue("file1.txt");
_fileQueue.Enqueue("file2.txt");
_fileQueue.Enqueue("file3.txt");
}
private async Task ProcessFilesAsync(CancellationToken token)
{
var tasks = new ConcurrentBag<Task>();
while (_fileQueue.TryDequeue(out var file))
{
bool acquired = await _semaphore.WaitAsync(1000, token);
if (!acquired)
{
_fileQueue.Enqueue(file); // 重新入队
Debug.WriteLine($"跳过任务 {file}:超时或取消");
continue;
}
tasks.Add(Task.Run(async () =>
{
try
{
var content = await File.ReadAllTextAsync($"C:\\Files\\{file}", token); // 异步读取
_fileCache.AddOrUpdate(file, content, (k, v) => content); // 原子更新
Interlocked.Increment(ref _successCount); // 原子递增
}
catch (IOException ex)
{
Debug.WriteLine($"文件 {file} 错误: {ex.Message}, 当前计数: {_semaphore.CurrentCount}");
_fileQueue.Enqueue(file); // 重新入队
}
finally
{
if (acquired)
{
_semaphore.Release();
Debug.WriteLine($"释放信号量,当前计数: {_semaphore.CurrentCount}");
}
}
}, token));
}
try
{
var stopwatch = Stopwatch.StartNew();
await Task.WhenAll(tasks); // 等待所有任务完成
BeginInvoke((Action)(() => fileContentLabel.Text = $"{string.Join("\n", _fileCache.Values)} (成功: {Interlocked.Read(ref _successCount)})"));
// 动态调整并发
var successRate = (double)Interlocked.Read(ref _successCount) / _fileQueue.Count;
if (successRate < 0.5 && _currentConcurrency > 1)
{
Interlocked.Decrement(ref _currentConcurrency);
_timer.Interval = (int)(_timer.Interval * 1.2);
Debug.WriteLine($"成功率低 ({successRate:P0}),减少并发到 {_currentConcurrency}");
}
else if (successRate > 0.9 && _currentConcurrency < _maximumCount)
{
if (_semaphore.CurrentCount < _maximumCount)
{
_semaphore.Release();
Interlocked.Increment(ref _currentConcurrency);
_timer.Interval = Math.Max(2000, (int)(_timer.Interval / 1.2));
Debug.WriteLine($"成功率高 ({successRate:P0}),增加并发到 {_currentConcurrency}, 当前计数: {_semaphore.CurrentCount}");
}
}
}
catch (OperationCanceledException)
{
BeginInvoke((Action)(() => fileContentLabel.Text = "处理取消"));
}
catch (AggregateException ex)
{
Debug.WriteLine($"批量错误: {ex.InnerException?.Message}, 当前计数: {_semaphore.CurrentCount}");
BeginInvoke((Action)(() => fileContentLabel.Text = $"错误: {ex.InnerException?.Message}"));
}
}
private void CancelButton_Click(object sender, EventArgs e)
{
_cts.Cancel();
_cts = new CancellationTokenSource();
}
protected override void OnFormClosing(FormClosingEventArgs e)
{
base.OnFormClosing(e);
_timer.Stop();
_cts.Cancel();
_cts.Dispose();
_semaphore.Dispose();
}
}
说明:
- 适用性:批量文件读写、数据库查询。
- Task.WhenAll:并行读取多个文件,等待所有完成。
- 异步编程:File.ReadAllTextAsync 异步读取,await Task.WhenAll 释放 UI 线程。
- 原子操作:Interlocked.Increment 计数,ConcurrentDictionary.AddOrUpdate 缓存。
- 最佳实践:
- Task.WhenAll 高效并行处理,批量收集结果。
- ConcurrentQueue 管理任务队列,ConcurrentDictionary 高效缓存。
- 检查 acquired 和 CurrentCount 避免 SemaphoreFullException.
- 失败重新入队,批量 UI 更新。
- 稳定性:捕获 IOException 和 AggregateException,finally 块释放。
- 效率:动态调整并发,异步并行处理。
- Task.WhenAll 优势:简化多文件并行读取,适合 I/O 密集型任务。
2.4 硬件通讯(TCP 数据采集)场景:通过 TCP 连接从多个硬件设备读取数据,限制并发,动态调整频率,缓存结果,批量更新 UI。
- 适用方法:Task.WhenAll + async/await + SemaphoreSlim + ConcurrentQueue + ConcurrentDictionary + Interlocked + TcpClient.
- 最佳实践:
- 使用 Task.WhenAll 并行处理多个 TCP 连接。
- ConcurrentQueue 缓冲数据,ConcurrentDictionary 缓存结果。
- Interlocked 跟踪读取计数。
代码示例:csharp
using System;
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using System.Windows.Forms;
public partial class MainForm : Form
{
private readonly System.Windows.Forms.Timer _timer = new System.Windows.Forms.Timer();
private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(2, 4);
private readonly ConcurrentQueue<(string, string)> _dataQueue = new ConcurrentQueue<(string, string)>();
private readonly ConcurrentDictionary<string, string> _dataCache = new ConcurrentDictionary<string, string>();
private readonly CancellationTokenSource _cts = new CancellationTokenSource();
private readonly TcpClient[] _tcpClients = { new TcpClient(), new TcpClient() };
private readonly NetworkStream[] _streams = new NetworkStream[2];
private int _currentConcurrency = 2;
private readonly int _maximumCount = 4;
private long _readCount = 0; // 原子计数器
public MainForm()
{
InitializeComponent();
_timer.Interval = 1000; // 1秒
_timer.Tick += async (s, e) => await ProcessTcpDataAsync(_cts.Token);
Task.Run(() => StartTcpListenersAsync(_cts.Token));
try
{
_tcpClients[0].Connect("192.168.1.100", 5000); // 设备 1
_tcpClients[1].Connect("192.168.1.101", 5000); // 设备 2
_streams[0] = _tcpClients[0].GetStream();
_streams[1] = _tcpClients[1].GetStream();
_timer.Start();
}
catch (SocketException ex)
{
Debug.WriteLine($"TCP 连接失败: {ex.Message}");
BeginInvoke((Action)(() => tcpLabel.Text = $"错误: {ex.Message}"));
}
}
private async Task StartTcpListenersAsync(CancellationToken token)
{
var tasks = new Task[_tcpClients.Length];
for (int i = 0; i < _tcpClients.Length; i++)
{
int index = i;
tasks[i] = Task.Run(async () =>
{
var buffer = new byte[1024];
while (!token.IsCancellationRequested)
{
try
{
int bytesRead = await _streams[index].ReadAsync(buffer, 0, buffer.Length, token); // 异步读取
if (bytesRead > 0)
{
var data = Encoding.UTF8.GetString(buffer, 0, bytesRead);
_dataQueue.Enqueue(($"device{index + 1}", data)); // 数据入队
}
}
catch (Exception ex)
{
Debug.WriteLine($"TCP 设备 {index + 1} 读取失败: {ex.Message}");
}
}
}, token);
}
try
{
await Task.WhenAll(tasks); // 等待所有监听任务
}
catch (OperationCanceledException)
{
Debug.WriteLine("TCP 监听取消");
}
}
private async Task ProcessTcpDataAsync(CancellationToken token)
{
bool acquired = await _semaphore.WaitAsync(0, token);
if (!acquired)
{
Debug.WriteLine("跳过数据处理:信号量已满");
return;
}
try
{
var stopwatch = Stopwatch.StartNew();
var tasks = new ConcurrentBag<Task>();
while (_dataQueue.TryDequeue(out var data))
{
tasks.Add(Task.Run(() =>
{
var (deviceId, content) = data;
_dataCache.AddOrUpdate(deviceId, content, (k, v) => content); // 原子更新
Interlocked.Increment(ref _readCount); // 原子递增
}, token));
}
await Task.WhenAll(tasks); // 等待所有数据处理任务
BeginInvoke((Action)(() => tcpLabel.Text = $"{string.Join("\n", _dataCache.Values)} (读取: {Interlocked.Read(ref _readCount)})"));
// 动态调整频率
if (stopwatch.ElapsedMilliseconds > 1000 && _timer.Interval < 5000)
{
_timer.Interval = (int)(_timer.Interval * 1.2);
Debug.WriteLine($"读取耗时高,降低频率到 {_timer.Interval}ms");
}
else if (stopwatch.ElapsedMilliseconds < 500 && _timer.Interval > 500)
{
if (_semaphore.CurrentCount < _maximumCount)
{
_semaphore.Release();
_timer.Interval = Math.Max(500, (int)(_timer.Interval / 1.2));
Debug.WriteLine($"读取耗时低,增加频率到 {_timer.Interval}ms, 当前计数: {_semaphore.CurrentCount}");
}
}
}
catch (OperationCanceledException)
{
BeginInvoke((Action)(() => tcpLabel.Text = "数据处理取消"));
}
catch (AggregateException ex)
{
Debug.WriteLine($"批量错误: {ex.InnerException?.Message}, 当前计数: {_semaphore.CurrentCount}");
BeginInvoke((Action)(() => tcpLabel.Text = $"错误: {ex.InnerException?.Message}"));
}
finally
{
if (acquired)
{
_semaphore.Release();
Debug.WriteLine($"释放信号量,当前计数: {_semaphore.CurrentCount}");
}
}
}
private void CancelButton_Click(object sender, EventArgs e)
{
_cts.Cancel();
_cts = new CancellationTokenSource();
}
protected override void OnFormClosing(FormClosingEventArgs e)
{
base.OnFormClosing(e);
_timer.Stop();
_cts.Cancel();
_cts.Dispose();
_semaphore.Dispose();
foreach (var stream in _streams) stream?.Dispose();
foreach (var client in _tcpClients) client?.Dispose();
}
}
说明:
- 适用性:TCP 硬件通讯、传感器数据采集。
- Task.WhenAll:并行处理多个 TCP 连接和数据处理任务。
- 异步编程:NetworkStream.ReadAsync 异步读取,await Task.WhenAll 释放 UI 线程。
- 原子操作:Interlocked.Increment 计数,ConcurrentDictionary.AddOrUpdate 缓存。
- 最佳实践:
- Task.WhenAll 高效并行处理多设备数据。
- ConcurrentQueue 缓冲数据,ConcurrentDictionary 高效缓存。
- 检查 acquired 和 CurrentCount 避免 SemaphoreFullException.
- 批量 UI 更新,严格管理 TCP 资源。
- 稳定性:捕获 SocketException 和 AggregateException,finally 块释放。
- 效率:动态调整频率,异步并行处理。
- Task.WhenAll 优势:简化多设备数据并行处理,适合硬件通讯。
3. 具体应用场景
|
场景 |
适用方法 |
具体应用 |
Task.WhenAll 适用性 |
|---|---|---|---|
|
实时数据刷新 |
Task.WhenAll + SemaphoreSlim + Timer |
多股票价格、传感器数据 |
优秀(并行刷新多数据源) |
|
批量数据同步 |
Task.WhenAll + SemaphoreSlim |
多 API 请求、数据库同步 |
优秀(并行处理多任务) |
|
I/O 密集型任务 |
Task.WhenAll + SemaphoreSlim + ConcurrentQueue |
文件读写、数据库查询 |
优秀(并行 I/O 操作) |
|
硬件通讯 |
Task.WhenAll + SemaphoreSlim + ConcurrentQueue |
串口/TCP 通信 |
优秀(并行处理多设备) |
4. 确保高效、稳定、健壮的实践
- 高效性:
- 并行异步:使用 Task.WhenAll 并行执行异步任务:csharp
await Task.WhenAll(tasks); - 并发控制:SemaphoreSlim 限制并发,动态调整:csharp
if (_semaphore.CurrentCount < _maximumCount) { _semaphore.Release(); Interlocked.Increment(ref _currentConcurrency); } - 批量更新:收集结果后一次性更新 UI:csharp
BeginInvoke((Action)(() => dataLabel.Text = string.Join("\n", _results.Values))); - 资源复用:复用 HttpClient 或 TcpClient:csharp
private readonly HttpClient _httpClient = new HttpClient();
- 并行异步:使用 Task.WhenAll 并行执行异步任务:csharp
- 稳定性:
- 异常处理:捕获 AggregateException 和具体异常(如 HttpRequestException、IOException),记录日志:csharp
catch (AggregateException ex) { Debug.WriteLine($"批量错误: {ex.InnerException?.Message}, 当前计数: {_semaphore.CurrentCount}"); BeginInvoke((Action)(() => dataLabel.Text = $"错误: {ex.InnerException?.Message}")); } - 线程安全:使用 ConcurrentQueue<T>、ConcurrentDictionary<TKey, TValue> 和 Interlocked:csharp
_results.AddOrUpdate(endpoint, data, (k, v) => data); Interlocked.Increment(ref _successCount); - 资源释放:finally 块释放 SemaphoreSlim,Dispose 清理资源:csharp
if (acquired) { _semaphore.Release(); Debug.WriteLine($"释放信号量,当前计数: {_semaphore.CurrentCount}"); } - 避免 SemaphoreFullException:csharp
bool acquired = await _semaphore.WaitAsync(0, token); if (acquired) { try { /* 任务逻辑 */ } finally { _semaphore.Release(); } } - 避免死锁:避免在 UI 线程使用 Task.Result 或 Wait,始终使用 await.
- 异常处理:捕获 AggregateException 和具体异常(如 HttpRequestException、IOException),记录日志:csharp
- 健壮性:
- 取消支持:CancellationToken 统一取消任务:csharp
private void CancelButton_Click(object sender, EventArgs e) { _cts.Cancel(); _cts = new CancellationTokenSource(); } - 重试机制:失败任务重新入队(文件操作):csharp
_fileQueue.Enqueue(file); - 状态监控:跟踪成功率、计数:csharp
var successRate = (double)Interlocked.Read(ref _successCount) / endpoints.Length; - 资源管理:严格管理 HttpClient、TcpClient 等资源:csharp
foreach (var stream in _streams) stream?.Dispose(); foreach (var client in _tcpClients) client?.Dispose(); - 测试性:封装接口以便单元测试:csharp
public interface ITimer { void Start(); void Stop(); int Interval { get; set; } event EventHandler Tick; } public interface IConcurrencyControl { Task<bool> WaitAsync(TimeSpan timeout, CancellationToken token); void Release(); }
- 取消支持:CancellationToken 统一取消任务:csharp
5. 总结Task.WhenAll 深入分析:
- 功能:并行执行多个异步任务,等待所有完成,收集结果。
- 优点:高性能,保持 UI 响应性,简化批量等待,适合 I/O 密集型任务。
- 缺点:必须等待所有任务完成,异常处理复杂,需限制并发。
- 与 Task.WhenAny 比较:Task.WhenAll 适合批量完成,Task.WhenAny 适合快速响应。
- 与锁比较:Task.WhenAll 异步非阻塞,锁适合复杂同步逻辑。
- 与信号量比较:Task.WhenAll 结合 SemaphoreSlim 控制并发。
- 与原子操作比较:Task.WhenAll 处理复杂任务,原子操作适合简单计数/键值更新。
- 与 Parallel.For 比较:Task.WhenAll 适合 I/O 密集型,Parallel.For 适合 CPU 密集型。
应用场景:
- 实时数据刷新:Task.WhenAll 并行刷新多数据源,ConcurrentDictionary 缓存。
- 批量数据同步:Task.WhenAll 并行 API 请求,ConcurrentDictionary 存储。
- I/O 密集型任务:Task.WhenAll 并行文件读写,ConcurrentQueue 调度。
- 硬件通讯:Task.WhenAll 并行处理多设备数据,ConcurrentQueue 缓冲。
高效、稳定、健壮:
- 高效性:并行异步、并发控制、批量更新、资源复用。
- 稳定性:异常处理、线程安全、资源释放、避免 SemaphoreFullException 和死锁。
- 健壮性:取消支持、重试机制、状态监控、资源管理、接口封装。
通过深入分析和最佳实践,Task.WhenAll 在 WinForms 中是处理批量异步任务的理想工具,结合 SemaphoreSlim、并发集合和原子操作,可实现高效、稳定、健壮的并发控制,满足实时数据处理、批量同步和硬件通讯的复杂需求。
深入 Task.WhenAny
Parallel.For 应用
更多推荐
所有评论(0)