redis_process.go 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135
  1. package lc
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "github.com/redis/go-redis/v9"
  6. "github.com/sirupsen/logrus"
  7. "regexp"
  8. "strings"
  9. "time"
  10. )
  11. // 启动定时任务:动态调整轮询间隔(未拿到车牌号3秒轮询,拿到则5秒轮询)
  12. func (s *SentinelServer) StartFetchingLatestPlate() {
  13. go func() {
  14. for {
  15. // 1. 执行获取最新车牌号逻辑
  16. latestData, err := s.GetLatestPlateData(1)
  17. if err != nil {
  18. //logrus.Warn("获取最新车牌号失败", "错误", err)
  19. time.Sleep(3 * time.Second)
  20. } else if len(latestData) == 0 {
  21. logrus.Info("未获取到车牌号数据")
  22. time.Sleep(3 * time.Second)
  23. } else {
  24. // 2. 成功获取到车牌号,处理业务逻辑
  25. latestPlateNo := latestData[0].PlateNo
  26. logrus.Info("最新车牌号", "plate_no", latestPlateNo)
  27. s.Notify(latestPlateNo, 2, 0)
  28. time.Sleep(12 * time.Second)
  29. }
  30. }
  31. }()
  32. }
  33. type RedisProcess struct{}
  34. type PlateData struct {
  35. PlateNo string `json:"plate_no"` // 车牌号
  36. PlateColor string `json:"plate_color"` // 车牌颜色
  37. DetectConf string `json:"detect_conf"` // 检测置信度
  38. ColorConf string `json:"color_conf"` // 颜色置信度
  39. RecAvg string `json:"rec_avg"` // 识别平均置信度
  40. Timestamp string `json:"timestamp"` // 秒级时间戳
  41. TimestampMs string `json:"timestamp_ms"` // 毫秒级时间戳
  42. Datetime string `json:"datetime"` // 格式化时间
  43. Source string `json:"source"` // 数据来源
  44. }
  45. func (s *SentinelServer) GetLatestPlateData(limit int) ([]PlateData, error) {
  46. if s.rdb == nil {
  47. return nil, fmt.Errorf("Redis客户端未初始化")
  48. }
  49. // 步骤1:从ZSet获取最新的limit个timestamp(按Score倒序)
  50. zsetKey := fmt.Sprintf("%s:sorted", RedisKeyPrefix)
  51. // ZREVRANGE:倒序取0到limit-1的member(最新的N个)
  52. timestamps, err := s.rdb.ZRevRange(s.ctx, zsetKey, 0, int64(limit-1)).Result()
  53. if err != nil {
  54. return nil, fmt.Errorf("读取ZSet失败: %w", err)
  55. }
  56. if len(timestamps) == 0 {
  57. return nil, fmt.Errorf("Redis中无车牌数据")
  58. }
  59. // 步骤2:批量从Hash读取对应数据
  60. hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix)
  61. pipe := s.rdb.Pipeline()
  62. for _, ts := range timestamps {
  63. pipe.HGet(s.ctx, hashKey, ts)
  64. }
  65. results, err := pipe.Exec(s.ctx)
  66. if err != nil {
  67. return nil, fmt.Errorf("批量读取Hash失败: %w", err)
  68. }
  69. // 步骤3:解析Python格式的字符串为Go结构体
  70. var plateDataList []PlateData
  71. for i, res := range results {
  72. if res == nil {
  73. continue
  74. }
  75. // 获取Hash的value(Python的str(dict)字符串)
  76. pyStr, err := res.(*redis.StringCmd).Result()
  77. if err != nil {
  78. logrus.Warn("解析单条数据失败", "timestamp", timestamps[i], "err", err)
  79. continue
  80. }
  81. // 转换Python字符串为标准JSON,再解析
  82. plateData, err := parsePythonDictStr(pyStr)
  83. if err != nil {
  84. logrus.Warn("转换Python字符串失败", "str", pyStr, "err", err)
  85. continue
  86. }
  87. plateDataList = append(plateDataList, plateData)
  88. }
  89. return plateDataList, nil
  90. }
  91. // 2. 读取指定时间戳的车牌数据(精准读取)
  92. func (s *SentinelServer) GetPlateDataByTimestamp(timestamp string) (*PlateData, error) {
  93. if s.rdb == nil {
  94. return nil, fmt.Errorf("Redis客户端未初始化")
  95. }
  96. hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix)
  97. pyStr, err := s.rdb.HGet(s.ctx, hashKey, timestamp).Result()
  98. if err != nil {
  99. return nil, fmt.Errorf("读取指定时间戳数据失败: %w", err)
  100. }
  101. // 解析为结构体
  102. plateData, err := parsePythonDictStr(pyStr)
  103. if err != nil {
  104. return nil, fmt.Errorf("解析数据失败: %w", err)
  105. }
  106. return &plateData, nil
  107. }
  108. // ========== 辅助函数:解析Python的dict字符串为Go结构体 ==========
  109. // Python的str(dict)是单引号,需替换为双引号,且处理特殊字符
  110. func parsePythonDictStr(pyStr string) (PlateData, error) {
  111. var plateData PlateData
  112. // 1. 替换单引号为双引号(Python -> JSON)
  113. jsonStr := strings.ReplaceAll(pyStr, "'", "\"")
  114. // 2. 处理可能的空格/换行(可选,视Python字符串格式)
  115. jsonStr = regexp.MustCompile(`\s+`).ReplaceAllString(jsonStr, "")
  116. // 3. 解析为JSON
  117. err := json.Unmarshal([]byte(jsonStr), &plateData)
  118. if err != nil {
  119. return plateData, fmt.Errorf("JSON解析失败: %w, raw_str: %s", err, jsonStr)
  120. }
  121. return plateData, nil
  122. }