You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 07:57:51