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

如何让ASP.NET Core的BackgroundService无需Console.ReadLine()持续接收MQTT消息

在ASP.NET Core 6中用BackgroundService稳定运行MQTTNet客户端

我在ASP.NET Core 6应用里实现了一个继承自BackgroundService的MqttClientService,用来运行MQTTNet客户端,处理收到的MQTT消息并回复成功标识。目前用Console.ReadLine()维持服务运行,感觉这是临时方案,有没有更好的办法让BackgroundService持续处理消息而不频繁重启?

另外,找到过ASP.NET Core搭配MQTTNet 3的示例,但它采用接口实现处理器,和当前版本的异步事件模式不匹配(参考MQTTNet升级指南)。


原相关代码

Services/MqttClientService.cs

using MQTTnet;
using MQTTnet.Client;
using System.Text;

namespace MqttClientAspNetCore.Services
{
    public class MqttClientService : BackgroundService
    {
        protected override async Task ExecuteAsync(CancellationToken stoppingToken)
        {
            while (!stoppingToken.IsCancellationRequested)
            {
                await Handle_Received_Application_Message();
            }
        }

        public static async Task Handle_Received_Application_Message()
        {

            var mqttFactory = new MqttFactory();

            using (var mqttClient = mqttFactory.CreateMqttClient())
            {
                var mqttClientOptions = new MqttClientOptionsBuilder()
                    .WithTcpServer("test.mosquitto.org")
                    .Build();

                // Setup message handling before connecting so that queued messages
                // are also handled properly. 
                mqttClient.ApplicationMessageReceivedAsync += e =>
                {
                    Console.WriteLine("### RECEIVED APPLICATION MESSAGE ###");
                    Console.WriteLine($"+ Payload = {Encoding.UTF8.GetString(e.ApplicationMessage.Payload)}");

                    // Publish successful message in response
                    var applicationMessage = new MqttApplicationMessageBuilder()
                        .WithTopic("keipalatest/1/resp")
                        .WithPayload("OK")
                        .Build();

                    mqttClient.PublishAsync(applicationMessage, CancellationToken.None);

                    Console.WriteLine("MQTT application message is published.");

                    return Task.CompletedTask;
                };

                await mqttClient.ConnectAsync(mqttClientOptions, CancellationToken.None);

                var mqttSubscribeOptions = mqttFactory.CreateSubscribeOptionsBuilder()
                    .WithTopicFilter(f =>
                    {
                        f.WithTopic("keipalatest/1/post");
                        f.WithAtLeastOnceQoS();
                    })
                    .Build();

                await mqttClient.SubscribeAsync(mqttSubscribeOptions, CancellationToken.None);

                Console.WriteLine("MQTT client subscribed to topic.");
                // The line below feels like a hack to keep background service from restarting
                Console.ReadLine();
            }
        }
    }
}

Program.cs

using MqttClientAspNetCore.Services;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddHostedService<MqttClientService>();

var app = builder.Build();

// To check if web server is still responsive
app.MapGet("/", () =>
{
    return "Hello World";
});


app.Run();

优化方案

核心优化方向

  1. 移除Console.ReadLine(),改用取消令牌实现优雅等待
  2. 复用MQTT客户端实例,避免重复创建和连接
  3. 添加自动重连逻辑,处理客户端意外断开场景
  4. 正确异步处理消息事件,避免异步操作未完成的问题

修改后的代码

Services/MqttClientService.cs

using MQTTnet;
using MQTTnet.Client;
using System.Text;

namespace MqttClientAspNetCore.Services
{
    public class MqttClientService : BackgroundService
    {
        private readonly IMqttClient _mqttClient;
        private readonly MqttClientOptions _mqttClientOptions;

        public MqttClientService()
        {
            var mqttFactory = new MqttFactory();
            _mqttClient = mqttFactory.CreateMqttClient();
            
            // 初始化客户端配置
            _mqttClientOptions = new MqttClientOptionsBuilder()
                .WithTcpServer("test.mosquitto.org")
                .Build();
            
            // 注册事件回调
            _mqttClient.ApplicationMessageReceivedAsync += OnApplicationMessageReceivedAsync;
            _mqttClient.DisconnectedAsync += OnDisconnectedAsync;
        }

        protected override async Task ExecuteAsync(CancellationToken stoppingToken)
        {
            // 循环维持连接,直到服务停止
            while (!stoppingToken.IsCancellationRequested)
            {
                try
                {
                    if (!_mqttClient.IsConnected)
                    {
                        await _mqttClient.ConnectAsync(_mqttClientOptions, stoppingToken);
                        await SubscribeToTopicsAsync(stoppingToken);
                    }

                    // 等待取消信号,避免循环空转
                    await Task.Delay(Timeout.Infinite, stoppingToken);
                }
                catch (OperationCanceledException)
                {
                    // 服务停止时退出循环
                    break;
                }
                catch (Exception ex)
                {
                    Console.WriteLine($"MQTT客户端异常: {ex.Message}");
                    // 异常后延迟重试,避免频繁连接
                    await Task.Delay(5000, stoppingToken);
                }
            }

            // 服务停止时优雅断开连接
            if (_mqttClient.IsConnected)
            {
                await _mqttClient.DisconnectAsync(stoppingToken: stoppingToken);
            }
        }

        private async Task OnApplicationMessageReceivedAsync(MqttApplicationMessageReceivedEventArgs e)
        {
            Console.WriteLine("### 收到应用消息 ###");
            Console.WriteLine($"+ 负载 = {Encoding.UTF8.GetString(e.ApplicationMessage.Payload)}");

            // 回复成功消息
            var responseMessage = new MqttApplicationMessageBuilder()
                .WithTopic("keipalatest/1/resp")
                .WithPayload("OK")
                .Build();

            await _mqttClient.PublishAsync(responseMessage, CancellationToken.None);
            Console.WriteLine("MQTT回复消息已发布");
        }

        private async Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs e)
        {
            Console.WriteLine("MQTT客户端已断开连接,将自动重试");
        }

        private async Task SubscribeToTopicsAsync(CancellationToken stoppingToken)
        {
            var mqttFactory = new MqttFactory();
            var subscribeOptions = mqttFactory.CreateSubscribeOptionsBuilder()
                .WithTopicFilter(f =>
                {
                    f.WithTopic("keipalatest/1/post");
                    f.WithAtLeastOnceQoS();
                })
                .Build();

            await _mqttClient.SubscribeAsync(subscribeOptions, stoppingToken);
            Console.WriteLine("MQTT客户端已订阅主题");
        }
    }
}

Program.cs(无需修改)

using MqttClientAspNetCore.Services;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddHostedService<MqttClientService>();

var app = builder.Build();

// 检查Web服务器是否响应
app.MapGet("/", () =>
{
    return "Hello World";
});

app.Run();

关键优化说明

  • 复用客户端实例:将MQTT客户端作为类成员变量,避免每次循环创建新对象,减少资源消耗和连接开销
  • 优雅等待机制:用await Task.Delay(Timeout.Infinite, stoppingToken)替代Console.ReadLine(),在服务收到停止信号时能立即响应并退出
  • 自动重连逻辑:通过ExecuteAsync循环检查连接状态,断开后自动重试连接和订阅,保证服务连续性
  • 异步事件正确处理:消息处理事件中await PublishAsync,确保异步操作完成,避免潜在的线程安全问题
  • 优雅关闭:服务停止时主动断开MQTT客户端连接,保证资源正确释放

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 09:35:22