使用手动重置事件释放多个线程

知识点

什么是ManualResetEvent

ManualResetEvent是一种同步原语,与AutoResetEvent不同,它需要手动重置状态。当事件处于信号状态时,所有等待的线程都会被释放,事件保持信号状态直到手动调用Reset()。

ManualResetEvent vs AutoResetEvent

特性 ManualResetEvent AutoResetEvent
重置方式 手动重置 自动重置
唤醒线程数 全部等待线程 单个线程
状态保持 直到手动重置 自动重置
适用场景 广播通知 一对一通知

核心方法

  • WaitOne():等待事件信号
  • WaitOne(timeout):带超时的等待
  • Set():设置事件为信号状态,释放所有等待线程
  • Reset():手动将事件重置为非信号状态

ManualResetEventSlim

.NET 4.0引入了ManualResetEventSlim,它是ManualResetEvent的轻量级版本:

  • 更好的性能
  • 支持取消令牌
  • 混合同步模式(先自旋后阻塞)

使用场景

  • 启动/停止信号
  • 阶段性同步
  • 广播通知
  • 初始化完成信号

代码案例

案例1:启动门控制

using System;
using System.Threading;
using System.Threading.Tasks;

class StartGateExample
{
    private static readonly ManualResetEvent startGate = new ManualResetEvent(false);
    private static readonly ManualResetEvent finishLine = new ManualResetEvent(false);
    
    static void Main(string[] args)
    {
        Console.WriteLine("启动门控制示例");
        Console.WriteLine("所有工作线程将等待启动信号...");
        
        // 创建多个工作线程
        Task[] workers = new Task[5];
        for (int i = 0; i < workers.Length; i++)
        {
            int workerId = i;
            workers[i] = Task.Run(() => WorkerThread(workerId));
        }
        
        // 让工作线程先准备好
        Thread.Sleep(2000);
        
        Console.WriteLine("\n主线程: 3秒后发送启动信号...");
        Thread.Sleep(3000);
        
        Console.WriteLine("主线程: 发送启动信号!");
        startGate.Set(); // 释放所有等待的工作线程
        
        // 等待所有工作线程完成
        Task.WaitAll(workers);
        
        Console.WriteLine("主线程: 所有工作线程已完成");
        
        startGate.Dispose();
        finishLine.Dispose();
    }
    
    static void WorkerThread(int workerId)
    {
        Console.WriteLine($"工作线程{workerId}: 准备就绪,等待启动信号...");
        
        // 等待启动信号
        startGate.WaitOne();
        
        Console.WriteLine($"工作线程{workerId}: 收到启动信号,开始工作!");
        
        // 模拟工作
        Random random = new Random(workerId);
        int workTime = random.Next(2000, 5000);
        
        for (int i = 1; i <= 5; i++)
        {
            Thread.Sleep(workTime / 5);
            Console.WriteLine($"工作线程{workerId}: 工作进度 {i}/5");
        }
        
        Console.WriteLine($"工作线程{workerId}: 工作完成!");
    }
}

案例2:多阶段初始化

using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;

class MultiPhaseInitialization
{
    private static readonly ManualResetEvent configLoaded = new ManualResetEvent(false);
    private static readonly ManualResetEvent databaseConnected = new ManualResetEvent(false);
    private static readonly ManualResetEvent servicesInitialized = new ManualResetEvent(false);
    
    private static Dictionary<string, string> configuration = new Dictionary<string, string>();
    private static bool isDatabaseReady = false;
    private static List<string> initializedServices = new List<string>();
    
    static void Main(string[] args)
    {
        Console.WriteLine("多阶段初始化示例");
        
        // 启动初始化任务
        Task configTask = Task.Run(() => LoadConfiguration());
        Task dbTask = Task.Run(() => InitializeDatabase());
        Task serviceTask = Task.Run(() => InitializeServices());
        
        // 启动依赖于初始化的工作任务
        Task[] workerTasks = new Task[3];
        for (int i = 0; i < workerTasks.Length; i++)
        {
            int workerId = i;
            workerTasks[i] = Task.Run(() => WorkerTask(workerId));
        }
        
        // 等待所有任务完成
        Task.WaitAll(new[] { configTask, dbTask, serviceTask }.Concat(workerTasks).ToArray());
        
        Console.WriteLine("应用程序初始化完成");
        
        // 清理资源
        configLoaded.Dispose();
        databaseConnected.Dispose();
        servicesInitialized.Dispose();
    }
    
    static void LoadConfiguration()
    {
        Console.WriteLine("配置加载: 开始加载配置文件...");
        
        // 模拟配置加载
        Thread.Sleep(1500);
        
        configuration["DatabaseConnectionString"] = "Server=localhost;Database=MyApp;";
        configuration["ApiKey"] = "abc123def456";
        configuration["LogLevel"] = "Info";
        
        Console.WriteLine("配置加载: 配置文件加载完成");
        
        // 通知配置已加载
        configLoaded.Set();
    }
    
    static void InitializeDatabase()
    {
        Console.WriteLine("数据库初始化: 等待配置加载完成...");
        
        // 等待配置加载完成
        configLoaded.WaitOne();
        
        Console.WriteLine("数据库初始化: 开始连接数据库...");
        
        // 模拟数据库连接
        Thread.Sleep(2000);
        
        string connectionString = configuration["DatabaseConnectionString"];
        Console.WriteLine($"数据库初始化: 使用连接字符串: {connectionString}");
        
        isDatabaseReady = true;
        Console.WriteLine("数据库初始化: 数据库连接成功");
        
        // 通知数据库已就绪
        databaseConnected.Set();
    }
    
    static void InitializeServices()
    {
        Console.WriteLine("服务初始化: 等待配置和数据库就绪...");
        
        // 等待配置和数据库都就绪
        WaitHandle.WaitAll(new[] { configLoaded, databaseConnected });
        
        Console.WriteLine("服务初始化: 开始初始化服务...");
        
        // 模拟服务初始化
        string[] services = { "日志服务", "缓存服务", "邮件服务", "通知服务" };
        
        foreach (string service in services)
        {
            Thread.Sleep(800);
            initializedServices.Add(service);
            Console.WriteLine($"服务初始化: {service} 初始化完成");
        }
        
        Console.WriteLine("服务初始化: 所有服务初始化完成");
        
        // 通知服务已就绪
        servicesInitialized.Set();
    }
    
    static void WorkerTask(int workerId)
    {
        Console.WriteLine($"工作任务{workerId}: 等待所有初始化完成...");
        
        // 等待所有初始化阶段完成
        WaitHandle.WaitAll(new[] { configLoaded, databaseConnected, servicesInitialized });
        
        Console.WriteLine($"工作任务{workerId}: 初始化完成,开始工作");
        
        // 访问初始化的资源
        Console.WriteLine($"工作任务{workerId}: 访问配置 - API密钥: {configuration["ApiKey"]}");
        Console.WriteLine($"工作任务{workerId}: 数据库状态: {(isDatabaseReady ? "就绪" : "未就绪")}");
        Console.WriteLine($"工作任务{workerId}: 可用服务数量: {initializedServices.Count}");
        
        // 模拟工作
        for (int i = 1; i <= 3; i++)
        {
            Thread.Sleep(1000);
            Console.WriteLine($"工作任务{workerId}: 执行工作步骤 {i}/3");
        }
        
        Console.WriteLine($"工作任务{workerId}: 工作完成");
    }
}

案例3:生产者消费者广播模式

using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

class BroadcastProducerConsumer
{
    private static readonly ManualResetEventSlim dataReady = new ManualResetEventSlim(false);
    private static readonly ManualResetEventSlim processingComplete = new ManualResetEventSlim(false);
    
    private static ConcurrentBag<string> sharedData = new ConcurrentBag<string>();
    private static volatile bool isProducing = true;
    private static int activeConsumers = 0;
    
    static void Main(string[] args)
    {
        Console.WriteLine("生产者消费者广播模式示例");
        
        // 启动生产者
        Task producer = Task.Run(() => Producer());
        
        // 启动多个消费者
        Task[] consumers = new Task[4];
        for (int i = 0; i < consumers.Length; i++)
        {
            int consumerId = i;
            consumers[i] = Task.Run(() => Consumer(consumerId));
        }
        
        // 运行8秒后停止生产
        Thread.Sleep(8000);
        isProducing = false;
        
        // 等待生产者完成
        producer.Wait();
        
        // 发送最后一批数据处理信号
        dataReady.Set();
        
        // 等待所有消费者完成
        Task.WaitAll(consumers);
        
        Console.WriteLine("所有任务完成");
        
        dataReady.Dispose();
        processingComplete.Dispose();
    }
    
    static void Producer()
    {
        int batchNumber = 1;
        
        while (isProducing)
        {
            Console.WriteLine($"生产者: 准备第 {batchNumber} 批数据...");
            
            // 重置事件,准备新一批数据
            dataReady.Reset();
            processingComplete.Reset();
            
            // 生产一批数据
            for (int i = 1; i <= 5; i++)
            {
                string data = $"Batch{batchNumber}_Item{i}_{DateTime.Now:HH:mm:ss.fff}";
                sharedData.Add(data);
                Thread.Sleep(200); // 模拟生产时间
            }
            
            Console.WriteLine($"生产者: 第 {batchNumber} 批数据准备完成,共 {sharedData.Count} 项");
            
            // 通知所有消费者数据就绪
            dataReady.Set();
            
            // 等待所有消费者完成处理
            Console.WriteLine($"生产者: 等待消费者处理第 {batchNumber} 批数据...");
            processingComplete.Wait();
            
            Console.WriteLine($"生产者: 第 {batchNumber} 批数据处理完成\n");
            batchNumber++;
            
            Thread.Sleep(1000); // 准备下一批数据的间隔
        }
        
        Console.WriteLine("生产者: 停止生产");
    }
    
    static void Consumer(int consumerId)
    {
        Interlocked.Increment(ref activeConsumers);
        Console.WriteLine($"消费者{consumerId}: 启动");
        
        while (true)
        {
            Console.WriteLine($"消费者{consumerId}: 等待数据就绪...");
            
            // 等待数据就绪信号
            dataReady.Wait();
            
            if (!isProducing && sharedData.IsEmpty)
            {
                Console.WriteLine($"消费者{consumerId}: 生产结束且无数据,退出");
                break;
            }
            
            Console.WriteLine($"消费者{consumerId}: 开始处理数据");
            
            // 处理所有可用数据
            int processedCount = 0;
            while (sharedData.TryTake(out string data))
            {
                Console.WriteLine($"消费者{consumerId}: 处理 {data}");
                Thread.Sleep(300); // 模拟处理时间
                processedCount++;
            }
            
            Console.WriteLine($"消费者{consumerId}: 处理了 {processedCount} 项数据");
            
            // 检查是否所有消费者都完成了处理
            if (sharedData.IsEmpty)
            {
                // 使用线程安全的方式检查并通知
                lock (processingComplete)
                {
                    if (sharedData.IsEmpty && !processingComplete.IsSet)
                    {
                        Console.WriteLine($"消费者{consumerId}: 数据处理完成,发送完成信号");
                        processingComplete.Set();
                    }
                }
            }
        }
        
        Interlocked.Decrement(ref activeConsumers);
        Console.WriteLine($"消费者{consumerId}: 已退出");
    }
}

案例4:同步检查点

using System;
using System.Threading;
using System.Threading.Tasks;

class SynchronizationCheckpoint
{
    private static readonly ManualResetEvent[] checkpoints = new ManualResetEvent[3];
    private static readonly object lockObject = new object();
    private static int completedTasks = 0;
    
    static void Main(string[] args)
    {
        Console.WriteLine("同步检查点示例");
        
        // 初始化检查点
        for (int i = 0; i < checkpoints.Length; i++)
        {
            checkpoints[i] = new ManualResetEvent(false);
        }
        
        // 启动多个任务
        Task[] tasks = new Task[6];
        for (int i = 0; i < tasks.Length; i++)
        {
            int taskId = i;
            tasks[i] = Task.Run(() => PhaseBasedTask(taskId));
        }
        
        // 监控任务进度
        Task monitorTask = Task.Run(() => MonitorProgress());
        
        // 等待所有任务完成
        Task.WaitAll(tasks);
        
        Console.WriteLine("所有任务完成");
        
        // 清理资源
        foreach (var checkpoint in checkpoints)
        {
            checkpoint.Dispose();
        }
    }
    
    static void PhaseBasedTask(int taskId)
    {
        Console.WriteLine($"任务{taskId}: 开始执行");
        
        // 阶段1:初始化
        Console.WriteLine($"任务{taskId}: 阶段1 - 初始化");
        Thread.Sleep(new Random(taskId).Next(1000, 2000));
        Console.WriteLine($"任务{taskId}: 阶段1完成");
        
        ReportPhaseCompletion(taskId, 0);
        
        // 等待所有任务完成阶段1
        Console.WriteLine($"任务{taskId}: 等待其他任务完成阶段1...");
        checkpoints[0].WaitOne();
        Console.WriteLine($"任务{taskId}: 所有任务完成阶段1,继续阶段2");
        
        // 阶段2:数据处理
        Console.WriteLine($"任务{taskId}: 阶段2 - 数据处理");
        Thread.Sleep(new Random(taskId + 100).Next(1500, 2500));
        Console.WriteLine($"任务{taskId}: 阶段2完成");
        
        ReportPhaseCompletion(taskId, 1);
        
        // 等待所有任务完成阶段2
        Console.WriteLine($"任务{taskId}: 等待其他任务完成阶段2...");
        checkpoints[1].WaitOne();
        Console.WriteLine($"任务{taskId}: 所有任务完成阶段2,继续阶段3");
        
        // 阶段3:清理
        Console.WriteLine($"任务{taskId}: 阶段3 - 清理");
        Thread.Sleep(new Random(taskId + 200).Next(800, 1200));
        Console.WriteLine($"任务{taskId}: 阶段3完成");
        
        ReportPhaseCompletion(taskId, 2);
        
        // 等待所有任务完成阶段3
        Console.WriteLine($"任务{taskId}: 等待其他任务完成阶段3...");
        checkpoints[2].WaitOne();
        Console.WriteLine($"任务{taskId}: 所有阶段完成!");
    }
    
    static void ReportPhaseCompletion(int taskId, int phase)
    {
        lock (lockObject)
        {
            completedTasks++;
            Console.WriteLine($"检查点: 任务{taskId}完成阶段{phase + 1},已完成任务数: {completedTasks}");
            
            // 检查是否所有任务都完成了当前阶段
            if (completedTasks % 6 == 0) // 6个任务
            {
                int completedPhase = completedTasks / 6;
                Console.WriteLine($"检查点: 所有任务完成阶段{completedPhase},发送继续信号");
                checkpoints[completedPhase - 1].Set();
            }
        }
    }
    
    static void MonitorProgress()
    {
        Console.WriteLine("监控器: 开始监控任务进度");
        
        for (int phase = 0; phase < 3; phase++)
        {
            Console.WriteLine($"监控器: 等待阶段{phase + 1}完成...");
            checkpoints[phase].WaitOne();
            Console.WriteLine($"监控器: 阶段{phase + 1}完成,所有任务同步成功");
        }
        
        Console.WriteLine("监控器: 所有阶段完成");
    }
}

知识点总结

  1. ManualResetEvent的核心特性

    • 手动重置:需要显式调用Reset()
    • 广播通知:Set()会释放所有等待的线程
    • 状态保持:信号状态会一直保持直到手动重置
  2. 适用场景

    • 启动门控制:所有线程同时开始
    • 阶段性同步:等待所有任务完成某个阶段
    • 广播通知:一次通知多个等待者
    • 初始化完成信号
  3. ManualResetEventSlim优势

    • 更好的性能
    • 支持CancellationToken
    • 混合等待模式(先自旋后阻塞)
  4. 最佳实践

    • 使用WaitHandle.WaitAll等待多个事件
    • 合理使用Reset()控制信号状态
    • 考虑使用ManualResetEventSlim提升性能
    • 总是释放资源避免内存泄漏
  5. 与AutoResetEvent的选择

    • 需要广播通知时选择ManualResetEvent
    • 需要一对一通知时选择AutoResetEvent
    • 考虑线程唤醒的数量和时机
  6. 注意事项

    • 忘记Reset()可能导致所有后续WaitOne立即返回
    • 合理处理超时和取消操作
    • 避免在持有锁时调用WaitOne
Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐