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

C# .NET Framework 4.7.2无法连接AWS MQTT Broker问题求助

问题

我正在开发基于.NET Framework 4.7.2的Windows Forms项目,尝试使用MQTTnet(版本3.0.16)连接AWS MQTT Broker,但多次尝试均未成功。调试时程序能识别并加载证书,却始终无法建立连接,每次都会返回“Broker doesn't respond to {endpoint}:{port}”的提示。移除.WithTls(parameters)配置后,可正常连接本地Broker。相关代码如下:

using System;
using System.Collections.Generic;
using System.IO;
using System.Security.Authentication;
using System.Security.Cryptography.X509Certificates;
using System.Text;
using System.Text.RegularExpressions;   //for Regex Class
using System.Windows.Forms;
using MQTTnet;          //version 3.0.16
using MQTTnet.Client;
using MQTTnet.Client.Connecting;
using MQTTnet.Client.Disconnecting;
using MQTTnet.Client.Options;
using MQTTnet.Client.Subscribing;
using MQTTnet.Extensions.ManagedClient;
using MQTTnet.Formatter;
using MQTTnet.Server;

namespace MQTTprj
{
    public partial class Form1 : Form
    {
        public Regex ipRegex = new Regex(@"^\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}$");
        public Regex hostnameRegex = new Regex(@"^(([a-zA-Z0-9]|[a-zA-Z0-9][a-zA-Z0-9\-_]*[a-zA-Z0-9])\.)*([a-zA-Z0-9]|[a-zA-Z0-9][a-zA-Z0-9\-_]*[a-zA-Z0-9])$");
        //log section string 
        private StringBuilder logMessages = new StringBuilder();
        private IMqttClient _mqttClient;
        

        //check if the string is such as "example.com"
        bool check_hostname(string endpoint)
        {
            return hostnameRegex.IsMatch(endpoint);
        }
        //check if the string is such as "192.168.0.1"
        bool check_ip(string endpoint)
        {
            return ipRegex.IsMatch(endpoint);
        }

        public Form1()
        {
            InitializeComponent();
        }


        private void btn_connect_Click(object sender, EventArgs e)
        {
            try
            {
                if (_mqttClient != null && _mqttClient.IsConnected)
                {
                    AddLog($"Client already connected to {this.endpoint.Text}:{int.Parse(this.port.Text)}");

                }
                else
                {
                    
                    //clear log
                    logMessages.Clear();
                    txt_log.Text = "";
                    //check if the port is a number and the address are valid
                    if (int.TryParse(port.Text, out int port_num) && (check_hostname(endpoint.Text) || check_ip(endpoint.Text)))
                    {
                        if (port_num < 1 || port_num > 65535)
                        {
                            throw new ArgumentOutOfRangeException("Input Error: port must be inside (1 - 65535) range");
                        }
                        var factory = new MqttFactory();
                        _mqttClient = factory.CreateMqttClient();

                        var endpoint = this.endpoint.Text;
                        var port = int.Parse(this.port.Text);
                        var clientCert = X509Certificate.CreateFromCertFile("AmazonRootCA1.pem");
                        var caCert = new X509Certificate2("certificate.pem.crt", "private.pem.key");

                          IMqttClientOptions options = new MqttClientOptionsBuilder()
                                .WithTcpServer(endpoint, port)
                                .WithCredentials("test", "testcl")
                                .WithProtocolVersion(MqttProtocolVersion.V311)
                                .WithTls(new MqttClientOptionsBuilderTlsParameters
                                {
                                    UseTls = true,
                                    SslProtocol = SslProtocols.Tls12,
                                    AllowUntrustedCertificates = false,
                                    Certificates = new List<X509Certificate>()
                                      {
                                          clientCert,
                                          caCert
                                      }
                                })
                                // .WithClientId(Guid.NewGuid().ToString())
                                .WithClientId("Merlin")
                                .Build();
                      
                        if (!_mqttClient.IsConnected)
                        {
                            try
                            {
                                _mqttClient.ConnectAsync(options).GetAwaiter().GetResult();
                                if (_mqttClient.IsConnected)
                                {
                                    AddLog($"Client successfully connected to {endpoint}:{port}");
                                }
                            }
                            catch (Exception)
                            {
                                AddLog($"Broker doesn't respond to {endpoint}:{port} ");
                            }

                        }
                        else
                        {
                            AddLog($"Client already connected to {endpoint}:{port}");
                        }
                    }
                    else
                    {
                        throw new FormatException("Input Error: the endpoint or port are not valid, check if they are correct and retry.");
                    }
                }


            }
            catch (FormatException ex)
            {
                logMessages.Append(ex.Message + "\r\n");
                txt_log.Text = logMessages.ToString();

            }
            catch (ArgumentOutOfRangeException ex)
            {
                logMessages.Append(ex.Message + "\r\n");
                txt_log.Text = logMessages.ToString();
            }
            catch (Exception ex)
            {
                logMessages.Append(ex.Message + "\r\n");
                txt_log.Text = logMessages.ToString();
            }


        }

        private void btn_disconnect_Click(object sender, EventArgs e)
        {
            logMessages.Clear();
            txt_log.Text = "";
            if (_mqttClient.IsConnected)
            {
                _mqttClient?.DisconnectAsync().GetAwaiter().GetResult();
                if (!_mqttClient.IsConnected)
                {
                    AddLog("Disconnected");
                }
            }
            else
            {
                AddLog("No Client was started");
            }


        }

        private void btn_publish_Click(object sender, EventArgs e)
        {
            if (_mqttClient.IsConnected)
            {
                var topic = this.topic.Text;
                var message = this.message.Text;

                var mqttMessage = new MqttApplicationMessageBuilder()
                    .WithTopic(topic)
                    .WithPayload(message)
                    .WithExactlyOnceQoS()
                    .WithRetainFlag()
                    .Build();

                _mqttClient.PublishAsync(mqttMessage).GetAwaiter().GetResult();

                AddLog($"Published message to topic '{topic}': {message}");
            }
            else
            {
                AddLog("Before connect a client to a broker");
            }
        }
        private void AddLog(string message)
        {
            this.txt_log.AppendText($"[{DateTime.Now}] {message}{Environment.NewLine}");
        }
        private void Form1_FormClosing(object sender, FormClosingEventArgs e)
        {
            _mqttClient?.DisconnectAsync().GetAwaiter().GetResult();
        }


        private void btn_subscribe_Click(object sender, EventArgs e)
        {

            try
            {
                if (_mqttClient.IsConnected && _mqttClient != null)
                {
                    //setting up the handler
                    _mqttClient.UseApplicationMessageReceivedHandler(session =>
                    {
                        var json = Encoding.UTF8.GetString(session.ApplicationMessage.Payload);
                        BeginInvoke(new Action(() =>
                        {
                            txt_log.AppendText("Received JSON data: " + json + Environment.NewLine);
                        }));

                    });
                    //set up topic filter
                    var topic = this.topic.Text;
                    var topicfilter = new MqttTopicFilterBuilder().WithTopic(topic).Build();
                    //verifing the subscribe was successful
                    MQTTnet.Client.Subscribing.MqttClientSubscribeResult sub_res = _mqttClient.SubscribeAsync(topicfilter).GetAwaiter().GetResult();
                    if (sub_res.Items[0].ResultCode == MqttClientSubscribeResultCode.GrantedQoS1 || sub_res.Items[0].ResultCode == MqttClientSubscribeResultCode.GrantedQoS2 || sub_res.Items[0].ResultCode == MqttClientSubscribeResultCode.GrantedQoS0)
                    {

                        AddLog("Successfully subscribed to topic. " + topic);

                    }
                    else
                    {
                        AddLog("Failed to subscribe to topic. " + topic);
                    }

                }
                else
                {
                    //throw new NullReferenceException("Before specify a client to a broker");
                    AddLog("Before connect a client to a broker");
                }
            }
            catch (NullReferenceException)
            {
                AddLog("Before connect a client to a broker");
            }
        }
    }
}

应用界面布局


解决建议

  • 修正证书配置逻辑
    当前代码把根证书和客户端证书都放到了Certificates列表,不符合MQTTnet的TLS参数规范。根证书应放入CertificateAuthorityCertificates用于服务端验证,客户端证书(含私钥)放入ClientCertificates用于客户端身份认证。修复示例:

    var rootCa = new X509Certificate2("AmazonRootCA1.pem");
    var clientCert = new X509Certificate2("certificate.pem.crt", "private.pem.key");
    
    var tlsParams = new MqttClientOptionsBuilderTlsParameters
    {
        UseTls = true,
        SslProtocol = SslProtocols.Tls12,
        AllowUntrustedCertificates = false,
        CertificateAuthorityCertificates = new List<X509Certificate> { rootCa },
        ClientCertificates = new List<X509Certificate> { clientCert }
    };
    
  • 优化异常捕获
    目前的通用异常捕获掩盖了真实错误原因,建议输出完整异常信息(包括内部异常):

    catch (Exception ex)
    {
        AddLog($"连接失败: {ex.Message}");
        if (ex.InnerException != null)
        {
            AddLog($"内部错误详情: {ex.InnerException.Message}");
        }
    }
    

    这样能定位具体问题,比如证书链验证失败、端口不通、权限不足等。

  • 验证AWS MQTT配置

    1. 确认AWS端点格式正确(应为xxxx-ats.iot.region.amazonaws.com),端口使用TLS对应的8883;
    2. 检查IAM策略是否允许当前客户端ID、用户名密码(或证书)进行连接;
    3. 测试网络连通性:用telnet <endpoint> 8883验证端口是否可达;
    4. 如果使用证书认证,可移除.WithCredentials("test", "testcl")配置,AWS MQTT通常仅需证书即可完成身份验证。
  • 升级MQTTnet版本
    尝试升级到3.0.x分支的最新版本(如3.0.25),修复旧版本可能存在的TLS连接bug。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 02:09:56