基于Channels的发布订阅模式异常:输入消息未被消费
问题分析与修复方案
你的代码存在几个关键问题,导致输入消息无法在OrderProcessor中输出:
1. 托管服务未启动
通过AddHostedService<OrderProcessor>注册的托管服务,必须依托IHost启动才能触发StartAsync方法。原代码仅构建了ServiceProvider但未启动Host,导致OrderProcessor的消息监听逻辑从未执行。
2. 异步方法误用async void
OrderEntryPoint中的SendMessageToOrderProcessor使用async void,这种写法在控制台应用中无法跟踪异步任务状态,可能引发未处理异常或操作未完成就继续执行后续逻辑,应改为async Task。
3. Main方法未支持异步
要调用异步方法,需将Main改为异步形式,否则无法正确等待异步操作完成。
4. StartAsync任务处理不当
原代码用Task.Factory.StartNew包裹异步lambda,会返回嵌套的Task<Task>,外层任务会提前完成,无法正确跟踪后台循环的生命周期,应直接返回异步任务。
修正后的完整代码
Program.cs
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using System.Threading.Channels; namespace PubSubWithChannelInConsole { internal class Program { private static OrderEntryPoint? _orderEntryPoint; static async Task Main(string[] args) { var hostTask = SetupEngineAsync(); while (true) { Console.WriteLine("Say something"); var whatsSaid = Console.ReadLine(); if (!string.IsNullOrEmpty(whatsSaid)) { await _orderEntryPoint!.SendMessageToOrderProcessor(whatsSaid); } } } private static async Task SetupEngineAsync() { var host = Host.CreateDefaultBuilder() .ConfigureServices(services => { services.AddSingleton(Channel.CreateBounded<string>(100)); services.AddHostedService<OrderProcessor>(); services.AddSingleton<OrderEntryPoint>(); }) .Build(); await host.StartAsync(); _orderEntryPoint = host.Services.GetRequiredService<OrderEntryPoint>(); await host.WaitForShutdownAsync(); } } }
OrderProcessor.cs
using Microsoft.Extensions.Hosting; using System.Threading.Channels; using System.Threading; namespace PubSubWithChannelInConsole { public class OrderProcessor : BackgroundService { private readonly Channel<string> _channel; public OrderProcessor(Channel<string> channel) { _channel = channel; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await Console.Out.WriteLineAsync("HELLO! 开始监听消息..."); while (!stoppingToken.IsCancellationRequested) { var response = await _channel.Reader.ReadAsync(stoppingToken); Console.WriteLine($"You said {response}"); } } } }
OrderEntryPoint.cs
using System.Threading.Channels; using System.Threading.Tasks; namespace PubSubWithChannelInConsole { public interface IOrderEntryPoint { Task SendMessageToOrderProcessor(string message); } public class OrderEntryPoint : IOrderEntryPoint { private readonly Channel<string> _channel; public OrderEntryPoint(Channel<string> channel) { _channel = channel; } public async Task SendMessageToOrderProcessor(string message) { await _channel.Writer.WriteAsync(message); } } }
关键修改说明
- 使用
Host.CreateDefaultBuilder构建并启动Host,确保OrderProcessor的ExecuteAsync(继承BackgroundService后的标准方法)被执行。 - 将
OrderProcessor改为继承BackgroundService,这是托管服务的推荐实现方式,简化生命周期管理。 - 把
SendMessageToOrderProcessor改为async Task,并在Main中await调用,确保消息写入操作完成。 - 用
while (true)替代goto,代码更易读且符合C#规范。 - 使用
stoppingToken监听Host停止信号,确保后台循环能优雅终止。
内容的提问来源于stack exchange,提问作者Ibanez1408
相关产品推荐
相关产品推荐

