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

在Azure Function中订阅Azure VM上Mosquitto MQTT主题的问题

在.NET 6 Azure Function中订阅Mosquitto MQTT Broker的两种方案

方案一:直接订阅MQTT主题(无需中间件)

原有的WebJobs扩展库因.NET版本兼容问题无法使用,你可以借助通用MQTT客户端库(比如MQTTnet)直接在Azure Function中实现订阅逻辑,推荐结合定时器触发或使用专用/弹性高级计划维持长连接:

  • 安装依赖:通过NuGet安装MQTTnet和MQTTnet.Extensions.ManagedClient包
  • 初始化客户端:配置Azure VM上Mosquitto的公网IP/域名、端口、认证信息,创建托管MQTT客户端以自动处理重连
  • 订阅与处理:连接成功后订阅目标主题,编写消息处理逻辑
  • 注意:消耗计划下Azure Function可能因闲置回收,可添加定时器触发定期检查连接状态并重连

示例代码片段:

using MQTTnet;
using MQTTnet.Client;
using MQTTnet.Extensions.ManagedClient;
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
using System.Text;

public class MqttSubscriptionFunction
{
    private readonly IManagedMqttClient _mqttClient;
    private readonly ILogger<MqttSubscriptionFunction> _logger;

    public MqttSubscriptionFunction(ILogger<MqttSubscriptionFunction> logger)
    {
        _logger = logger;
        var factory = new MqttFactory();
        _mqttClient = factory.CreateManagedMqttClient();

        // 配置MQTT客户端参数
        var clientOptions = new MqttClientOptionsBuilder()
            .WithTcpServer("你的VM公网IP", 1883) // Mosquitto默认端口,TLS启用时用8883
            .WithCredentials("mqtt用户名", "mqtt密码") // 若Broker启用认证需配置
            .WithClientId($"AzureFuncClient_{Guid.NewGuid()}")
            .Build();

        var managedOptions = new ManagedMqttClientOptionsBuilder()
            .WithAutoReconnectDelay(TimeSpan.FromSeconds(5))
            .WithClientOptions(clientOptions)
            .Build();

        // 绑定消息接收事件
        _mqttClient.ApplicationMessageReceivedAsync += args =>
        {
            var payload = Encoding.UTF8.GetString(args.ApplicationMessage.Payload);
            _logger.LogInformation("收到MQTT消息:主题={Topic}, 内容={Payload}", args.ApplicationMessage.Topic, payload);
            // 此处添加业务处理逻辑
            return Task.CompletedTask;
        };

        // 启动客户端并订阅主题
        _mqttClient.StartAsync(managedOptions).Wait();
        _mqttClient.SubscribeAsync(new MqttTopicFilterBuilder().WithTopic("你的目标主题").Build()).Wait();
    }

    // 定时器触发维持连接(消耗计划下可选)
    [Function("MqttKeepAlive")]
    public async Task Run([TimerTrigger("*/5 * * * *")] TimerInfo timer)
    {
        if (!_mqttClient.IsConnected)
        {
            _logger.LogWarning("MQTT连接断开,尝试重连");
            await _mqttClient.StartAsync();
        }
        _logger.LogInformation("MQTT连接正常,当前时间: {Time}", DateTime.Now);
    }
}

方案二:通过Kafka中转(适合复杂场景)

如果业务需要消息持久化、多消费者负载均衡,或已有Kafka集群,可通过中转方式实现:

  • 部署Kafka:使用Azure Event Hubs(兼容Kafka协议)或自建Kafka集群
  • 消息转发:利用Mosquitto内置桥接功能,或编写轻量服务将MQTT主题消息转发至Kafka Topic
  • Kafka触发Function:创建Kafka触发的Azure Function,配置集群地址与目标Topic,处理消息
  • 优势:借助Kafka的持久化能力避免消息丢失,适配高并发、多下游消费场景

总结

  • 简单场景优先选方案一,实现成本低、部署快
  • 复杂业务场景(需持久化、负载均衡)选方案二,提升整体可靠性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 03:50:41