package lc import ( "encoding/json" "fmt" "github.com/redis/go-redis/v9" "github.com/sirupsen/logrus" "regexp" "strings" "time" ) // 启动定时任务:动态调整轮询间隔(未拿到车牌号3秒轮询,拿到则5秒轮询) func (s *SentinelServer) StartFetchingLatestPlate() { go func() { for { // 1. 执行获取最新车牌号逻辑 latestData, err := s.GetLatestPlateData(1) if err != nil { //logrus.Warn("获取最新车牌号失败", "错误", err) time.Sleep(3 * time.Second) } else if len(latestData) == 0 { logrus.Info("未获取到车牌号数据") time.Sleep(3 * time.Second) } else { // 2. 成功获取到车牌号,处理业务逻辑 latestPlateNo := latestData[0].PlateNo logrus.Info("最新车牌号", "plate_no", latestPlateNo) s.Notify(latestPlateNo, 2, 0) time.Sleep(12 * time.Second) } } }() } type RedisProcess struct{} type PlateData struct { PlateNo string `json:"plate_no"` // 车牌号 PlateColor string `json:"plate_color"` // 车牌颜色 DetectConf string `json:"detect_conf"` // 检测置信度 ColorConf string `json:"color_conf"` // 颜色置信度 RecAvg string `json:"rec_avg"` // 识别平均置信度 Timestamp string `json:"timestamp"` // 秒级时间戳 TimestampMs string `json:"timestamp_ms"` // 毫秒级时间戳 Datetime string `json:"datetime"` // 格式化时间 Source string `json:"source"` // 数据来源 } func (s *SentinelServer) GetLatestPlateData(limit int) ([]PlateData, error) { if s.rdb == nil { return nil, fmt.Errorf("Redis客户端未初始化") } // 步骤1:从ZSet获取最新的limit个timestamp(按Score倒序) zsetKey := fmt.Sprintf("%s:sorted", RedisKeyPrefix) // ZREVRANGE:倒序取0到limit-1的member(最新的N个) timestamps, err := s.rdb.ZRevRange(s.ctx, zsetKey, 0, int64(limit-1)).Result() if err != nil { return nil, fmt.Errorf("读取ZSet失败: %w", err) } if len(timestamps) == 0 { return nil, fmt.Errorf("Redis中无车牌数据") } // 步骤2:批量从Hash读取对应数据 hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix) pipe := s.rdb.Pipeline() for _, ts := range timestamps { pipe.HGet(s.ctx, hashKey, ts) } results, err := pipe.Exec(s.ctx) if err != nil { return nil, fmt.Errorf("批量读取Hash失败: %w", err) } // 步骤3:解析Python格式的字符串为Go结构体 var plateDataList []PlateData for i, res := range results { if res == nil { continue } // 获取Hash的value(Python的str(dict)字符串) pyStr, err := res.(*redis.StringCmd).Result() if err != nil { logrus.Warn("解析单条数据失败", "timestamp", timestamps[i], "err", err) continue } // 转换Python字符串为标准JSON,再解析 plateData, err := parsePythonDictStr(pyStr) if err != nil { logrus.Warn("转换Python字符串失败", "str", pyStr, "err", err) continue } plateDataList = append(plateDataList, plateData) } return plateDataList, nil } // 2. 读取指定时间戳的车牌数据(精准读取) func (s *SentinelServer) GetPlateDataByTimestamp(timestamp string) (*PlateData, error) { if s.rdb == nil { return nil, fmt.Errorf("Redis客户端未初始化") } hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix) pyStr, err := s.rdb.HGet(s.ctx, hashKey, timestamp).Result() if err != nil { return nil, fmt.Errorf("读取指定时间戳数据失败: %w", err) } // 解析为结构体 plateData, err := parsePythonDictStr(pyStr) if err != nil { return nil, fmt.Errorf("解析数据失败: %w", err) } return &plateData, nil } // ========== 辅助函数:解析Python的dict字符串为Go结构体 ========== // Python的str(dict)是单引号,需替换为双引号,且处理特殊字符 func parsePythonDictStr(pyStr string) (PlateData, error) { var plateData PlateData // 1. 替换单引号为双引号(Python -> JSON) jsonStr := strings.ReplaceAll(pyStr, "'", "\"") // 2. 处理可能的空格/换行(可选,视Python字符串格式) jsonStr = regexp.MustCompile(`\s+`).ReplaceAllString(jsonStr, "") // 3. 解析为JSON err := json.Unmarshal([]byte(jsonStr), &plateData) if err != nil { return plateData, fmt.Errorf("JSON解析失败: %w, raw_str: %s", err, jsonStr) } return plateData, nil }