06-第6章-数据流与控制流设计
·
第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
更多推荐
所有评论(0)