异步代码看似同步执行求助:C# Task未并行运行问题
问题根源:同步阻塞的
WaitOne()让任务无法并行 嘿,我一眼就揪出问题所在了——你的Process方法里的@event.WaitOne()是同步阻塞调用,这直接导致所有任务只能排队执行,根本没法真正并行!
为什么会这样?
虽然Process方法签名是Task<TaskInfo>,看起来是异步的,但里面没有真正的异步等待(await),反而用了WaitOne()这种会死死卡住当前线程的同步方法。当你在循环里调用i.Processor.Process(i.TaskId)时,每一次调用都会立刻执行到WaitOne(),然后把当前线程按住不动,直到对应的RabbitMQ响应回来、事件被触发,才会继续下一个循环迭代。这就难怪所有任务看起来都是同步执行的——它们确实是在排队等线程!
怎么修复?
把同步的AutoResetEvent换成**TaskCompletionSource<T>**,这是.NET里把基于事件的异步模式转换成Task异步模式的标准做法,能让你的方法真正异步起来,不会阻塞线程。
修改后的核心代码示例:
首先重构Process方法,用TaskCompletionSource替代事件:
public Task<TaskInfo> Process(int taskId) { var tcs = new TaskCompletionSource<TaskInfo>(); var taskInfo = new TaskInfo(); // 务必创建新实例,避免多任务共享对象的线程安全问题 var stopwatch = Stopwatch.StartNew(); // 把计时器改成局部变量,不要用成员变量 try { taskId = taskId + 1; using (var bus = RabbitHutch.CreateBus(xxDev)) { var replyTo = Guid.NewGuid().ToString(); var messageQueue = bus.Advanced.QueueDeclare(replyTo, autoDelete: true); // 消费消息时,直接用TaskCompletionSource标记任务完成 bus.Advanced.Consume(messageQueue, (payload, properties, info) => { var file = Format(outputFile, properties.CorrelationId); taskInfo.OutputFile = file; Console.WriteLine($"Output written to {file}, TaskId {taskId}"); File.WriteAllBytes(file, payload); var remaining = Interlocked.Decrement(ref outstandingRequests); if (remaining == 0) { stopwatch.Stop(); taskInfo.TimeTaken = stopwatch.Elapsed; tcs.SetResult(taskInfo); // 告诉Task:我完成了! } return Task.CompletedTask; }); // 准备并发送消息 taskInfo.InputFile = inputFile; var html = await File.ReadAllTextAsync(inputFile); // 把同步IO换成异步IO,进一步释放线程 taskInfo.Html = html; var message = PrepareMessage(new RenderRequest() { Html = Encoding.UTF8.GetBytes(html), Options = new RenderRequestOptions() { PageSize = "A4", ImageQuality = 70, PageLoadRetryAttempts = 3 } }); var correlation = Guid.NewGuid().ToString(); Console.WriteLine($"CorrelationId: {correlation}, TaskId {taskId}"); var props = new MessageProperties { CorrelationId = correlation, ReplyTo = replyTo, Expiration = "6000" }; Publish(bus, props, message); taskInfo.CorrelationId = Guid.Parse(correlation); } return tcs.Task; // 返回等待完成的Task } catch (Exception e) { stopwatch.Stop(); taskInfo.OutputFile = Empty; taskInfo.TimeTaken = stopwatch.Elapsed; tcs.SetException(e); // 传递异常给调用方 return tcs.Task; } }
额外注意事项:
- 别用成员变量存储
taskInfo和stopwatch,每个任务都要创建自己的实例,否则会出现多线程竞争的问题。 - 把
File.ReadAllText换成await File.ReadAllTextAsync,避免同步IO占用线程池线程。 - 确保
outstandingRequests的初始化正确,每次启动一批任务时要重置为任务总数,否则Interlocked.Decrement会逻辑出错。
为什么这样改能解决问题?
TaskCompletionSource让我们把“等待事件触发”转换成异步的await操作,不会阻塞当前线程。这样,循环里的每个任务都会立刻启动并发送RabbitMQ消息,然后把自己的Task加入列表,线程可以立刻去处理下一个任务,所有任务就能真正并行执行了。最后await Task.WhenAll会等待所有任务完成。
内容的提问来源于stack exchange,提问作者TheKeyboarder
相关产品推荐
相关产品推荐

