MQTT架构详细流程图

┌──────────┐     1. 发布原始数据     ┌────────────┐     2. 转发数据     ┌──────────┐
│          │───────────────────────>│            │──────────────────>│          │
│  设备    │                       │ MQTT Broker│                   │  EdgeX   │
│          │<───────────────────────│            │<──────────────────│          │
└──────────┘     6. 发送控制命令     └────────────┘     5. 发布命令     └──────────┘
                                           ↑
                                           │ 3. 发布标准化事件
                                           │
                                           ↓
┌──────────┐     7. 存储数据     ┌────────────┐     4. 转发事件     ┌──────────┐
│          │<────────────────────│            │<──────────────────│          │
│  sfsDb   │                    │  适配器    │                   │ MQTT Broker│
│          │                    │            │                   │            │
└──────────┘                    └────────────┘                   └────────────┘

详细流程说明

  1. 设备 → MQTT Broker:设备通过MQTT协议发布原始传感器数据到MQTT broker
  2. MQTT Broker → EdgeX:MQTT broker将设备数据转发给EdgeX(EdgeX作为MQTT客户端订阅了相关主题)
  3. EdgeX → MQTT Broker:EdgeX对数据进行标准化处理后,将处理后的事件发布到MQTT broker
  4. MQTT Broker → 适配器:MQTT broker将EdgeX发布的事件转发给适配器(适配器作为MQTT客户端订阅了相关主题)
  5. 适配器 → sfsDb:适配器解析和处理EdgeX事件,将数据存储到sfsDb
  6. EdgeX → MQTT Broker:EdgeX可以通过MQTT broker向设备发送控制命令
  7. MQTT Broker → 设备:MQTT broker将控制命令转发给设备

关键角色说明

1. 设备

  • 角色:MQTT客户端(发布者)
  • 功能:收集传感器数据并发布到MQTT broker

2. MQTT Broker

  • 角色:消息代理(中心枢纽)
  • 功能:接收、存储和转发MQTT消息

3. EdgeX

  • 角色:MQTT客户端(订阅者 + 发布者) 既是消息的消费者,也是消息的发布者
  • 功能
    • 订阅设备数据
    • 处理和标准化数据
    • 发布标准化事件
    • 发送设备控制命令

4. 适配器

  • 角色:MQTT客户端(订阅者)
  • 功能
    • 订阅EdgeX事件
    • 解析和处理事件
    • 存储数据到sfsDb

5. sfsDb

  • 角色:数据存储
  • 功能:持久化存储设备数据

这种架构设计使得系统各组件之间实现了松耦合,每个组件只需要与MQTT broker通信,而不需要直接与其他组件通信,提高了系统的可扩展性和可靠性。
完全正确!EdgeX和适配器都是MQTT客户端,它们通过MQTT broker进行通信。

这种架构的优势

  1. 松耦合:组件之间通过MQTT broker通信,不需要直接连接
  2. 可扩展性:可以轻松添加新的设备、服务或适配器
  3. 可靠性:MQTT协议支持消息的可靠传递
  4. 灵活性:支持多种设备和协议
  5. 可管理性:集中式的消息管理和监控

这种基于MQTT的架构是物联网和边缘计算系统中非常常见的设计模式,它提供了一种可靠、灵活的方式来处理设备数据。

我可以提供代码证据来支持MQTT架构的角色和流程:

1. 适配器作为MQTT客户端

mqtt/client.go 文件中,适配器创建并使用MQTT客户端:

// NewClient 创建新的MQTT客户端
func NewClient(cfg *config.Config, dataQueue *queue.Queue) (*Client, error) {
    client, err := initMQTT(cfg)
    if err != nil {
        return nil, err
    }

    return &Client{
        client:    client,
        config:    cfg,
        dataQueue: dataQueue,
    }, nil
}

// initMQTT 初始化MQTT客户端
func initMQTT(cfg *config.Config) (mqtt.Client, error) {
    opts := mqtt.NewClientOptions()
    opts.AddBroker(cfg.MQTTBroker) // 连接到MQTT broker
    opts.SetClientID(cfg.ClientID)
    opts.SetCleanSession(true)
    opts.SetAutoReconnect(true)
    opts.SetMaxReconnectInterval(time.Second * 30)

    client := mqtt.NewClient(opts)
    token := client.Connect()
    token.Wait()
    if token.Error() != nil {
        return nil, fmt.Errorf("failed to connect to MQTT broker: %v", token.Error())
    }

    log.Printf("Connected to MQTT broker: %s", cfg.MQTTBroker)
    return client, nil
}

2. 适配器订阅MQTT消息

mqtt/client.go 文件中,适配器订阅EdgeX发布的消息:

// Subscribe 订阅EdgeX消息
func (c *Client) Subscribe() error {
    token := c.client.Subscribe(c.config.MQTTTopic, 1, c.messageHandler())
    token.Wait()
    if token.Error() != nil {
        return fmt.Errorf("failed to subscribe to topic %s: %v", c.config.MQTTTopic, token.Error())
    }

    log.Printf("Subscribed to topic: %s", c.config.MQTTTopic)
    return nil
}

3. 适配器处理MQTT消息

mqtt/client.go 文件中,适配器处理收到的MQTT消息:

// messageHandler 消息处理函数
func (c *Client) messageHandler() mqtt.MessageHandler {
    return func(client mqtt.Client, msg mqtt.Message) {
        log.Printf("Received message on topic: %s", msg.Topic())

        // 使用edgex包处理消息
        event, err := edgex.ProcessMessage(msg.Payload())
        if err != nil {
            log.Printf("Failed to process message: %v", err)
            return
        }

        // 如果消息类型不是event,event会为nil
        if event == nil {
            return
        }

        // 收集所有读数,准备批量插入
        var records []*map[string]any

        // 处理每个读数
        for _, reading := range event.Readings {
            // 准备数据
            metadataStr := ""
            if reading.Metadata != nil {
                metadataStr = string(reading.Metadata)
            }

            // 解析值的类型
            value := common.ParseValue(reading.Value)

            data := map[string]any{
                "id":         reading.ID,
                "deviceName": event.DeviceName, // 设备名称已经在ProcessMessage中格式化
                "reading":    reading.ResourceName,
                "value":      value,
                "valueType":  reading.ValueType,
                "baseType":   reading.BaseType,
                "timestamp":  reading.Origin, // 纳秒级时间戳,类型为 int64
                "metadata":   metadataStr,
            }

            records = append(records, &data)
        }

        // 批量存储到 sfsDb
        if len(records) > 0 {
            // 使用重试机制插入数据
            err := database.BatchInsertWithRetry(database.Table, records, 3, 2*time.Second)
            if err != nil {
                log.Printf("Failed to batch store data after retries: %v", err)
                // 将数据加入队列,以便后续处理
                if err := c.dataQueue.Enqueue(records); err != nil {
                    log.Printf("Failed to enqueue data: %v", err)
                } else {
                    log.Printf("Enqueued %d readings for later processing", len(records))
                }
            } else {
                log.Printf("Batch stored %d readings from %s", len(records), event.DeviceName)
            }
        }
    }
}

4. 配置中的MQTT broker设置

config/config.go 文件中,配置了MQTT broker的地址和订阅主题:

// Load 加载配置
func Load() (*Config, error) {
    // 1. 设置默认配置
    cfg := &Config{
        DBPath:     "./edgex_data",
        MQTTBroker: "tcp://localhost:1883", // MQTT broker地址
        MQTTTopic:  "edgex/events/core/#", // 订阅EdgeX的核心事件
        ClientID:   generateClientID(),
        HTTPPort:   "8081", // 默认HTTP端口
    }

    // 2. 尝试从EdgeX配置中心加载
    if err := loadFromConfigCenter(cfg); err != nil {
        log.Printf("Failed to load config from EdgeX config center: %v", err)
        log.Println("Falling back to local config file")

        // 3. 从配置文件加载
        if err := loadFromFile(cfg); err != nil {
            log.Printf("Failed to load config from file: %v", err)
            log.Println("Using default config")
        }
    }

    // 4. 从环境变量加载(优先级最高)
    loadFromEnv(cfg)

    return cfg, nil
}

5. EdgeX消息格式

edgex/models.go 文件中,定义了EdgeX消息的格式:

// EdgeXMessage 消息结构(符合 MessageEnvelope 格式)
type EdgeXMessage struct {
    CorrelationID string          `json:"correlationId,omitempty"`
    MessageType   string          `json:"messageType,omitempty"`
    Origin        int64           `json:"origin,omitempty"`
    Payload       json.RawMessage `json:"payload"`
}

// EdgeXEvent 事件结构
type EdgeXEvent struct {
    ID          string         `json:"id"`
    DeviceName  string         `json:"deviceName"`
    Readings    []EdgeXReading `json:"readings"`
    Origin      int64          `json:"origin"`
    ProfileName string         `json:"profileName,omitempty"`
    SourceName  string         `json:"sourceName,omitempty"`
}

// EdgeXReading 读数结构
type EdgeXReading struct {
    ID           string          `json:"id"`
    ResourceName string          `json:"resourceName"`
    Value        string          `json:"value"`
    ValueType    string          `json:"valueType,omitempty"`
    Origin       int64           `json:"origin"`
    ProfileName  string          `json:"profileName,omitempty"`
    DeviceName   string          `json:"deviceName,omitempty"`
    BaseType     string          `json:"baseType,omitempty"`
    Metadata     json.RawMessage `json:"metadata,omitempty"`
}

6. 处理EdgeX消息

edgex/processor.go 文件中,处理EdgeX消息:

// ProcessMessage 处理EdgeX消息
func ProcessMessage(payload []byte) (*EdgeXEvent, error) {
    var edgexMsg EdgeXMessage
    if err := json.Unmarshal(payload, &edgexMsg); err != nil {
        return nil, err
    }

    // 检查 MessageType 是否为 "event"
    if edgexMsg.MessageType != "event" && edgexMsg.MessageType != "Event" {
        log.Printf("Ignoring message with type: %s", edgexMsg.MessageType)
        return nil, nil
    }

    // 解析 payload 中的事件
    var event EdgeXEvent
    if err := json.Unmarshal(edgexMsg.Payload, &event); err != nil {
        return nil, err
    }

    // 从源头格式化设备名称,确保长度为64字符
    event.DeviceName = common.FormatDeviceName(event.DeviceName)

    return &event, nil
}

7. 主程序中的MQTT客户端使用

main.go 文件中,创建和使用MQTT客户端:

func main() {
    // 加载配置
    var err error
    appConfig, err = config.Load()
    if err != nil {
        log.Fatalf("Failed to load config: %v", err)
    }

    // 连接 sfsDb
    if err := database.Init(appConfig.DBPath); err != nil {
        log.Fatalf("Failed to initialize database: %v", err)
    }

    // 初始化数据队列
    dataQueue, err = queue.NewQueue("./data_queue")
    if err != nil {
        log.Fatalf("Failed to initialize data queue: %v", err)
    }

    // 创建MQTT客户端
    mqttClient, err := mqtt.NewClient(appConfig, dataQueue)
    if err != nil {
        log.Fatalf("Failed to initialize MQTT: %v", err)
    }
    defer mqttClient.Disconnect()

    // 订阅EdgeX消息
    if err := mqttClient.Subscribe(); err != nil {
        log.Fatalf("Failed to subscribe to EdgeX messages: %v", err)
    }

    log.Println("sfsDb EdgeX adapter started successfully")

    // 启动队列处理 goroutine,处理可能存在添加失败的数据
    dataQueue.ProcessQueue(func(data interface{}) error {
        records, ok := data.([]*map[string]any)
        if !ok {
            return fmt.Errorf("invalid data type in queue")
        }

        // 使用重试机制插入数据
        return database.BatchInsertWithRetry(database.Table, records, 3, 2*time.Second)
    })

    // 启动 HTTP 服务器
    serverInstance := server.NewServer(database.Table, appConfig)
    if err := serverInstance.Start(); err != nil {
        log.Fatalf("Failed to start HTTP server: %v", err)
    }

    // 等待中断信号以优雅地关闭服务器
    quit := make(chan os.Signal, 1)                      // 创建一个信号通道,用于接收中断信号
    signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) // 注册信号通知,当收到SIGINT或SIGTERM时,将信号发送到quit通道
    <-quit                                               // 阻塞直到收到中断信号。
    log.Println("Shutting down adapter...")
    // 给服务器 5 秒的时间来完成正在处理的请求
    time.Sleep(5 * time.Second)

    log.Println("Adapter exited")
}

代码证据总结

从以上代码可以清楚地看到:

  1. 适配器是MQTT客户端:通过 mqtt.NewClient 创建MQTT客户端,连接到MQTT broker
  2. 适配器订阅EdgeX消息:通过 mqttClient.Subscribe 订阅 edgex/events/core/# 主题
  3. 适配器处理MQTT消息:通过 messageHandler 处理收到的消息
  4. 数据存储到sfsDb:通过 database.BatchInsertWithRetry 将数据批量存储到sfsDb
  5. EdgeX消息格式:定义了标准的EdgeX消息格式,包括消息信封、事件和读数

这些代码证据证实了MQTT架构的角色和数据流向:

  • 设备 → MQTT broker → EdgeX → MQTT broker → 适配器 → sfsDb
Logo

智能硬件社区聚焦AI智能硬件技术生态,汇聚嵌入式AI、物联网硬件开发者,打造交流分享平台,同步全国赛事资讯、开展 OPC 核心人才招募,助力技术落地与开发者成长。

更多推荐