在 .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. 确保高效、稳定、健壮的实践

  1. 高效性:
    • 并行异步:使用 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();
  2. 稳定性:
    • 异常处理:捕获 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.
  3. 健壮性:
    • 取消支持: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(); }

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 应用

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐