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

如何管理MQTTnet客户端生命周期?避免多线程使用时释放客户端

如何避免在另一个线程正在使用MQTTnet客户端时释放它?

这可能适用于任何IDisposable,但对于ManagedMqttClient,还需要担心异步调用前的IsConnected等检查。

说明:我们使用的是MQTTnet v3.0.16版本,接受“升级至最新版后使用方案X”这类答案

我接手了一个使用ManagedMqttClient的应用,最初在用户修改代理设置时会替换/释放该客户端,代码如下:

using MQTTnet;
using MQTTnet.Client.Disconnecting;
using MQTTnet.Client.Options;
using MQTTnet.Extensions.ManagedClient;
using System;
using System.Threading.Tasks;

internal class OriginalApproach
{
    private IManagedMqttClient _mqttClient;
    private static MqttFactory _factory;

    public OriginalApproach()
    {
        _mqttClient.DisconnectedHandler = new MqttClientDisconnectedHandlerDelegate(MqttClientDisconnectedEventArgs => OnDisconnect(MqttClientDisconnectedEventArgs));
    }

    //Called if the user changes settings that affect the way we connect
    //to the broker.
    public async void OnSettingsChange()
    {
        if (_mqttClient != null && _mqttClient.IsConnected)
        {
            StopAsync();
            return;
        }

        //Disposal isn't the only thread safety issue
        if (_mqttClient != null && _mqttClient.IsStarted)
        {
            await Reconnect(TimeSpan.FromSeconds(2));
        }
    }

    public async void StopAsync()
    {
        if (_mqttClient != null)
        {
            await _mqttClient.StopAsync();
            await Task.Delay(System.TimeSpan.FromSeconds(2));
        }
    }

    public async void OnDisconnect(MqttClientDisconnectedEventArgs e)
    {
        await Reconnect(TimeSpan.FromSeconds(5));
    }

    public async Task Reconnect(TimeSpan delay)
    {
        StopAsync();
        await Task.Delay(delay);
        Connect();
    }

    public async void Connect()
    {
        await CreateManagedClient();

        try
        {
            if (!_mqttClient.IsConnected && !_mqttClient.IsStarted)
            {
                StartAsync();
            }
        }
        catch (MQTTnet.Exceptions.MqttCommunicationException ex) { /* ... */  }
        catch (MQTTnet.Exceptions.MqttProtocolViolationException ex) { /* ... */  }
    }

    public async Task<bool> CreateManagedClient()
    {
        try
        {
            if (_mqttClient != null)
                _mqttClient.Dispose();

            _factory = new MqttFactory();
            _mqttClient = _factory.CreateManagedMqttClient();
            await Task.Delay(System.TimeSpan.FromSeconds(2));
        }
        catch (Exception e)
        {
            _mqttClient.Dispose();
            _mqttClient = null;
            return false;
        }
        return true;
    }

    public async void StartAsync()
    {
        MqttApplicationMessage mess = new MqttApplicationMessage();

        mess.Payload = BuildDeathCertificate();
        mess.Topic = "...";

        MqttClientOptionsBuilder clientOptionsBuilder = new MqttClientOptionsBuilder();

        IMqttClientOptions options = clientOptionsBuilder.WithTcpServer("Broker Address", 1234)
                .WithClientId("ABCD")
                .WithCleanSession(true)
                .WithWillMessage(mess)
                .WithKeepAlivePeriod(new System.TimeSpan(1234))
                .WithCommunicationTimeout(new System.TimeSpan(int.MaxValue))
                .Build();

        var managedClientOptions = new ManagedMqttClientOptionsBuilder()
            .WithClientOptions(options)
            .Build();

        if (!_mqttClient.IsStarted && !_mqttClient.IsConnected)
        {
            try
            {
                await _mqttClient.StartAsync(managedClientOptions);
            }
            catch (Exception e) { /* ... */  }
        }
    }

    byte[] BuildDeathCertificate()
    {
        return new byte[1234];
    }

    public async void PublishMessage(byte[] payloadBytes)
    {
        var message = new MqttApplicationMessageBuilder()
            .WithTopic("...")
            .WithPayload(payloadBytes)
            .WithExactlyOnceQoS()
            .WithRetainFlag(false)
            .Build();

        try
        {
            await _mqttClient.PublishAsync(message);
        }
        catch (NullReferenceException e) { /* ... */  }
    }
}

显然这里存在诸多线程安全问题,多种场景下会抛出ObjectDisposed异常。

尝试的方案1:单实例客户端

我尝试在应用生命周期内使用单个ManagedMqttClient,代码如下:

internal class SingleClientTest
{
    private IManagedMqttClient _mqttClient;
    public SingleClientTest()
    {
        var factory = new MqttFactory();

        //Used for lifetime of application
        _mqttClient = factory.CreateManagedMqttClient();
    }

    public async void Connect()
    {
        //No longer calling CreateManagedClient() here

        try
        {
            if (!_mqttClient.IsConnected && !_mqttClient.IsStarted)
            {
                StartAsync();
            }
        }
        catch (MQTTnet.Exceptions.MqttCommunicationException ex) { /* ... */  }
        catch (MQTTnet.Exceptions.MqttProtocolViolationException ex) { /* ... */  }
    }

    //The other methods are mostly unchanged
}

整体而言,这解决了ObjectDisposed问题,但并未解决异步调用前调用IsConnected的线程安全问题。而且考虑到MqttFactory的存在,重用单个客户端感觉像是一种权宜之计。此外,我遇到了类似场景:尽管IsStarted为false,但调用StartAsync()仍抛出“托管客户端已启动”的异常。

尝试的方案2:使用lock

我还尝试在客户端调用周围添加lock,但由于死锁风险,无法在await调用周围使用lock。

我查阅了MQTTnet的示例、维基文档、部分Issue,并浏览了部分代码,但尚未在库中找到额外的并发机制。

正在探索的方案

我正在探索几种方案(可能组合使用):

  • 在所有客户端调用周围使用SemaphoreSlim——它似乎可以解决await调用的问题,但不确定是否会引入新的时序问题,且考虑到我们使用.NET Framework,其使用存在风险
  • 使用MqttClient而非ManagedMqttClient。相关讨论表明MqttClient是首选。我是否应该改用它?在应用生命周期内使用单个MqttClient(修改代理设置时使用DisconnectAsync()/ConnectAsync())是否合理?(这仍未解决_mqttClient.IsConnected这类检查的问题)
  • 在每个客户端对象调用周围添加try/catch捕获ObjectDisposed异常,并按如下方式替换客户端:
var oldClient = _mqttClient
_mqttClient = _factory.CreateManagedMqttClient();
oldClient?.Dispose();

同样,这也未解决_mqttClient.IsConnected这类检查的问题。

我想知道业界普遍认可的解决方案是什么。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 21:00:56