第6章:数据流与控制流设计

6.1 主数据流

6.1.1 完整数据流路径

EdgeX Foundry
    ↓ (MQTT)
MQTT Client (mqtt/client.go)
    ↓
EdgeX Message Processor (edgex/processor.go)
    ↓
数据验证与转换
    ↓
批量收集
    ↓ [可选:数据库写入失败]
Queue (queue/queue.go) ←──────────────┐
    ↓                                  │
Database Batch Insert (database/database.go)
    ↓ [成功]
Monitor 更新指标
    ↓
Analyzer 数据分析(可选)

6.1.2 MQTT 消息处理流

// mqtt/client.go 中的消息处理流程
func (c *Client) messageHandler(client mqtt.Client, msg mqtt.Message) {
    // 1. 接收消息
    payload := msg.Payload()
    
    // 2. 解析 EdgeX 消息
    event, err := edgex.ProcessMessage(payload)
    if err != nil {
        log.Printf("Failed to process EdgeX message: %v", err)
        return
    }
    
    // 3. 处理每个读数
    for _, reading := range event.Readings {
        // 4. 准备数据
        data := map[string]any{
            "id":         reading.ID,
            "deviceName": event.DeviceName,
            "reading":    reading.ResourceName,
            "value":      common.ParseValue(reading.Value),
            "valueType":  reading.ValueType,
            "baseType":   reading.BaseType,
            "timestamp":  reading.Origin,
            "metadata":   metadataStr,
        }
        
        // 5. 加入批量队列
        c.batchMessages = append(c.batchMessages, data)
    }
    
    // 6. 检查是否达到批量大小
    if len(c.batchMessages) >= c.batchSize {
        c.flushBatch()
    }
}

6.1.3 批量刷新流程

// mqtt/client.go 中的批量刷新
func (c *Client) flushBatch() {
    if len(c.batchMessages) == 0 {
        return
    }
    
    // 1. 复制消息切片
    messages := make([]*map[string]any, len(c.batchMessages))
    for i := range c.batchMessages {
        messages[i] = &c.batchMessages[i]
    }
    
    // 2. 清空批量队列
    c.batchMessages = c.batchMessages[:0]
    
    // 3. 批量插入数据库
    _, err := database.Table.BatchInsertNoInc(messages)
    if err != nil {
        log.Printf("Failed to batch insert: %v", err)
        // 4. 写入失败,加入队列
        if err := c.dataQueue.Enqueue(messages); err != nil {
            log.Printf("Failed to enqueue: %v", err)
        }
        return
    }
    
    // 5. 更新监控
    if c.monitor != nil {
        c.monitor.IncrementMessagesReceived(len(messages))
    }
    
    // 6. 数据分析(可选)
    if c.analyzer != nil {
        c.analyzer.Analyze(messages)
    }
}

6.2 查询数据流

6.2.1 HTTP 查询流程

HTTP Request (GET /api/readings)
    ↓
AuthMiddleware (认证检查)
    ↓
DeviceNameMiddleware (格式化设备名)
    ↓
handleQueryReadings (server/server.go)
    ↓
database.QueryRecords (database/database.go)
    ↓
构建查询范围 (deviceName + timestamp)
    ↓
Table.SearchRange (sfsDb)
    ↓
迭代器获取记录
    ↓
JSON 编码响应
    ↓
HTTP Response

6.2.2 查询实现细节

// server/server.go:159-192
func (s *Server) handleQueryReadings(w http.ResponseWriter, r *http.Request) {
    if s.Monitor != nil {
        s.Monitor.IncrementHTTPRequests()
    }
    
    w.Header().Set("Content-Type", "application/json")
    
    deviceName := r.URL.Query().Get("deviceName")
    startTime := r.URL.Query().Get("startTime")
    endTime := r.URL.Query().Get("endTime")
    
    readings, err := database.QueryRecords(database.Table, deviceName, startTime, endTime)
    if err != nil {
        w.WriteHeader(http.StatusInternalServerError)
        json.NewEncoder(w).Encode(map[string]string{"error": err.Error()})
        return
    }
    defer readings.Release()
    
    readingsMap := make([]map[string]any, len(readings))
    for i, reading := range readings {
        readingsMap[i] = reading
    }
    
    json.NewEncoder(w).Encode(map[string]interface{}{
        "count":    len(readings),
        "readings": readingsMap,
    })
}

6.3 队列处理流

6.3.1 队列处理流程

后台 Goroutine 启动
    ↓
循环检查队列大小
    ↓
队列为空?→ 是 → 等待 5 秒
    ↓ 否
Dequeue 取出数据
    ↓
调用处理函数
    ↓
成功?→ 是 → 继续循环
    ↓ 否
重新 Enqueue
    ↓
等待 5 秒
    ↓
继续循环

6.3.2 队列处理实现

// queue/queue.go:144-185
func (q *Queue) ProcessQueue(processFunc func(interface{}) error) {
    go func() {
        for {
            size, err := q.Size()
            if err != nil {
                log.Printf("Failed to get queue size: %v", err)
                time.Sleep(5 * time.Second)
                continue
            }
            
            if size == 0 {
                time.Sleep(5 * time.Second)
                continue
            }
            
            data, err := q.Dequeue()
            if err != nil {
                log.Printf("Failed to dequeue data: %v", err)
                time.Sleep(5 * time.Second)
                continue
            }
            
            if data == nil {
                time.Sleep(5 * time.Second)
                continue
            }
            
            if err := processFunc(data); err != nil {
                log.Printf("Failed to process queue data: %v", err)
                if err := q.Enqueue(data); err != nil {
                    log.Printf("Failed to re-enqueue data: %v", err)
                }
                time.Sleep(5 * time.Second)
            }
        }
    }()
}

6.3.3 主程序中的队列处理

// main.go:105-112
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)
})

6.4 控制流设计

6.4.1 启动控制流

// main.go 中的启动顺序
func main() {
    // 1. 加载配置
    appConfig, err = config.Load()
    
    // 2. 初始化监控
    monitorInstance = monitor.NewMonitor()
    
    // 3. 初始化告警
    alertNotifier = alert.NewNotifier(appConfig)
    monitorInstance.SetNotifier(alertNotifier)
    alertNotifier.Start()
    
    // 4. 初始化数据库
    database.Init(appConfig.DBPath, ...)
    
    // 5. 初始化队列
    dataQueue, err = queue.NewQueue("./data_queue")
    
    // 6. 初始化 MQTT
    mqttClient, err = mqtt.NewClient(...)
    mqttClient.Subscribe()
    
    // 7. 启动队列处理
    dataQueue.ProcessQueue(...)
    
    // 8. 初始化其他模块
    agentInstance, err = agent.NewAgent(...)
    retentionManager = retention.NewRetentionManager(...)
    syncManager, err = sync.NewSyncManager(...)
    resourceMonitor = resource.NewResourceMonitor(...)
    
    // 9. 启动 HTTP 服务器
    serverInstance := server.NewServer(...)
    serverInstance.Start()
    
    // 10. 等待信号
    quit := make(chan os.Signal, 1)
    signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
    <-quit
}

6.4.2 关闭控制流

// main.go:166-205
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
<-quit
log.Println("Shutting down adapter...")

// 按依赖逆序停止
if agentInstance != nil {
    agentInstance.Stop()
}
if retentionManager != nil {
    retentionManager.Stop()
}
if alertNotifier != nil {
    alertNotifier.Stop()
}
if syncManager != nil {
    syncManager.Stop()
}
if resourceMonitor != nil {
    resourceMonitor.Stop()
}

time.Sleep(5 * time.Second)
log.Println("Adapter exited")

6.5 并发控制流

6.5.1 Goroutine 管理

main Goroutine
    ├─ MQTT 消息处理 Goroutine
    ├─ 队列处理 Goroutine
    ├─ HTTP 服务器 Goroutine
    ├─ Agent Goroutine
    ├─ Retention Manager Goroutine
    ├─ Alert Notifier Goroutine
    ├─ Sync Manager Goroutine
    └─ Resource Monitor Goroutine

6.5.2 信道通信

// 信号处理信道(main.go:167-169)
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
<-quit

// 互斥锁保护共享资源(queue/queue.go:34)
type Queue struct {
    queueDir string
    mutex    sync.Mutex
}

6.6 实战练习

练习 6.1:数据流追踪

为数据流添加追踪日志,追踪一条消息的完整路径。

练习 6.2:队列优化

优化队列处理流程,提高吞吐量。

练习 6.3:并发控制

为某个模块添加更细粒度的并发控制。

6.7 本章小结

本章详细介绍了 sfsEdgeStore 的数据流和控制流设计:

  • 主数据流路径(从 EdgeX 到数据库)
  • 查询数据流路径(从 HTTP 请求到响应)
  • 队列处理流程
  • 启动和关闭控制流
  • 并发控制和 Goroutine 管理

理解这些流程对于排查问题和优化性能至关重要。


本书版本:1.0.0
最后更新:2026-03-08
sfsEdgeStore - 让边缘数据存储更简单!🚀
技术栈 - Go语言、sfsDb与EdgeX Foundry。纯golang工业物联网边缘计算技术栈
项目地址GitHub

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐