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

.NET 6控制台应用RabbitMQ消费者移出Program.cs后失效求助

RabbitMQ RPC消费者后台服务实现方案

一、适配后台服务的RpcServer类实现

直接用IHostedService规范实现,确保连接和消费者的生命周期被正确管理:

using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System.Text;

public class RpcServer : IHostedService
{
    // 类级别持有连接和通道,避免被GC回收
    private IConnection _connection;
    private IModel _channel;
    private const string QueueName = "rpc_queue";

    public Task StartAsync(CancellationToken cancellationToken)
    {
        var factory = new ConnectionFactory() { HostName = "localhost" };
        _connection = factory.CreateConnection();
        _channel = _connection.CreateModel();
        
        // 声明队列,和官方教程逻辑一致
        _channel.QueueDeclare(queue: QueueName, durable: false, exclusive: false, autoDelete: false, arguments: null);
        _channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);

        var consumer = new EventingBasicConsumer(_channel);
        consumer.Received += (model, ea) =>
        {
            string response = string.Empty;
            var body = ea.Body.ToArray();
            var props = ea.BasicProperties;
            var replyProps = _channel.CreateBasicProperties();
            replyProps.CorrelationId = props.CorrelationId;

            try
            {
                var message = Encoding.UTF8.GetString(body);
                int n = int.Parse(message);
                Console.WriteLine($" [.] 计算Fib({n})");
                response = Fib(n).ToString();
            }
            catch (Exception e)
            {
                Console.WriteLine($" [x] 出错:{e.Message}");
                response = string.Empty;
            }
            finally
            {
                // 返回响应并手动确认消息
                var responseBytes = Encoding.UTF8.GetBytes(response);
                _channel.BasicPublish(exchange: string.Empty, routingKey: props.ReplyTo, basicProperties: replyProps, body: responseBytes);
                _channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
            }
        };

        _channel.BasicConsume(queue: QueueName, autoAck: false, consumer: consumer);

        return Task.CompletedTask;
    }

    // 服务停止时释放资源
    public Task StopAsync(CancellationToken cancellationToken)
    {
        _channel?.Close();
        _connection?.Close();
        return Task.CompletedTask;
    }

    // 官方教程中的斐波那契计算逻辑
    private static int Fib(int n)
    {
        if (n is 0 or 1)
            return n;
        return Fib(n - 1) + Fib(n - 2);
    }
}

二、Program.cs中托管启动后台服务

用.NET的Host机制托管服务,确保程序不会启动后立即退出:

using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;

var host = Host.CreateDefaultBuilder(args)
    .ConfigureServices(services =>
    {
        // 注册RpcServer为后台服务
        services.AddHostedService<RpcServer>();
    })
    .Build();

// 启动并保持服务运行
await host.RunAsync();

三、踩过的坑提醒

  • 别把连接/通道写成局部变量:如果在Start方法里把_connection或_channel声明成局部变量,方法执行完就会被GC回收,直接导致消费者失效。
  • AutoAck必须设为false:一定要手动调用BasicAck确认消息,否则消息处理失败时会丢失,或者消费者崩溃后消息无法重新投递。
  • 必须用Host托管:不能直接调用new RpcServer().StartAsync()就完事,控制台程序会执行完代码直接退出,消费者根本没机会处理消息。
  • 一定要加异常捕获:在Received事件回调里捕获所有异常,避免单条消息处理失败导致整个消费者进程崩溃。

内容的提问来源于stack exchange,提问作者barak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 03:42:40