10 使用手动重置事件释放多个线程
·
使用手动重置事件释放多个线程
知识点
什么是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("监控器: 所有阶段完成");
}
}
知识点总结
-
ManualResetEvent的核心特性:
- 手动重置:需要显式调用Reset()
- 广播通知:Set()会释放所有等待的线程
- 状态保持:信号状态会一直保持直到手动重置
-
适用场景:
- 启动门控制:所有线程同时开始
- 阶段性同步:等待所有任务完成某个阶段
- 广播通知:一次通知多个等待者
- 初始化完成信号
-
ManualResetEventSlim优势:
- 更好的性能
- 支持CancellationToken
- 混合等待模式(先自旋后阻塞)
-
最佳实践:
- 使用WaitHandle.WaitAll等待多个事件
- 合理使用Reset()控制信号状态
- 考虑使用ManualResetEventSlim提升性能
- 总是释放资源避免内存泄漏
-
与AutoResetEvent的选择:
- 需要广播通知时选择ManualResetEvent
- 需要一对一通知时选择AutoResetEvent
- 考虑线程唤醒的数量和时机
-
注意事项:
- 忘记Reset()可能导致所有后续WaitOne立即返回
- 合理处理超时和取消操作
- 避免在持有锁时调用WaitOne
更多推荐



所有评论(0)