sfsdb EdgeX 适配器 中的MQTT架构与数据流程
·
MQTT架构详细流程图
┌──────────┐ 1. 发布原始数据 ┌────────────┐ 2. 转发数据 ┌──────────┐
│ │───────────────────────>│ │──────────────────>│ │
│ 设备 │ │ MQTT Broker│ │ EdgeX │
│ │<───────────────────────│ │<──────────────────│ │
└──────────┘ 6. 发送控制命令 └────────────┘ 5. 发布命令 └──────────┘
↑
│ 3. 发布标准化事件
│
↓
┌──────────┐ 7. 存储数据 ┌────────────┐ 4. 转发事件 ┌──────────┐
│ │<────────────────────│ │<──────────────────│ │
│ sfsDb │ │ 适配器 │ │ MQTT Broker│
│ │ │ │ │ │
└──────────┘ └────────────┘ └────────────┘
详细流程说明
- 设备 → MQTT Broker:设备通过MQTT协议发布原始传感器数据到MQTT broker
- MQTT Broker → EdgeX:MQTT broker将设备数据转发给EdgeX(EdgeX作为MQTT客户端订阅了相关主题)
- EdgeX → MQTT Broker:EdgeX对数据进行标准化处理后,将处理后的事件发布到MQTT broker
- MQTT Broker → 适配器:MQTT broker将EdgeX发布的事件转发给适配器(适配器作为MQTT客户端订阅了相关主题)
- 适配器 → sfsDb:适配器解析和处理EdgeX事件,将数据存储到sfsDb
- EdgeX → MQTT Broker:EdgeX可以通过MQTT broker向设备发送控制命令
- 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进行通信。
这种架构的优势
- 松耦合:组件之间通过MQTT broker通信,不需要直接连接
- 可扩展性:可以轻松添加新的设备、服务或适配器
- 可靠性:MQTT协议支持消息的可靠传递
- 灵活性:支持多种设备和协议
- 可管理性:集中式的消息管理和监控
这种基于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")
}
代码证据总结
从以上代码可以清楚地看到:
- 适配器是MQTT客户端:通过
mqtt.NewClient创建MQTT客户端,连接到MQTT broker - 适配器订阅EdgeX消息:通过
mqttClient.Subscribe订阅edgex/events/core/#主题 - 适配器处理MQTT消息:通过
messageHandler处理收到的消息 - 数据存储到sfsDb:通过
database.BatchInsertWithRetry将数据批量存储到sfsDb - EdgeX消息格式:定义了标准的EdgeX消息格式,包括消息信封、事件和读数
这些代码证据证实了MQTT架构的角色和数据流向:
- 设备 → MQTT broker → EdgeX → MQTT broker → 适配器 → sfsDb
更多推荐


所有评论(0)