redis_process.go 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193
  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. // Redis基础配置(极简版,仅用于初始化连接)
  12. const (
  13. RedisInitDevice = "init_device_data" // 仅用于标记Redis初始化完成的空Key
  14. RedisKeyPrefix = "plate_results"
  15. )
  16. // 启动定时任务:动态调整轮询间隔(未拿到车牌号3秒轮询,拿到则5秒轮询)
  17. func (s *SentinelServer) StartFetchingLatestPlate() {
  18. go func() {
  19. for {
  20. // 1. 执行获取最新车牌号逻辑
  21. latestData, err := s.GetLatestPlateData(1)
  22. if err != nil {
  23. //logrus.Warn("获取最新车牌号失败", "错误", err)
  24. time.Sleep(3 * time.Second)
  25. } else if len(latestData) == 0 {
  26. logrus.Info("未获取到车牌号数据")
  27. time.Sleep(3 * time.Second)
  28. } else {
  29. // 2. 成功获取到车牌号,处理业务逻辑
  30. latestPlateNo := latestData[0].PlateNo
  31. logrus.Info("最新车牌号", "plate_no", latestPlateNo)
  32. s.Notify(latestPlateNo, 2, 0)
  33. time.Sleep(12 * time.Second)
  34. }
  35. }
  36. }()
  37. }
  38. type RedisProcess struct{}
  39. type PlateData struct {
  40. PlateNo string `json:"plate_no"` // 车牌号
  41. PlateColor string `json:"plate_color"` // 车牌颜色
  42. DetectConf string `json:"detect_conf"` // 检测置信度
  43. ColorConf string `json:"color_conf"` // 颜色置信度
  44. RecAvg string `json:"rec_avg"` // 识别平均置信度
  45. Timestamp string `json:"timestamp"` // 秒级时间戳
  46. TimestampMs string `json:"timestamp_ms"` // 毫秒级时间戳
  47. Datetime string `json:"datetime"` // 格式化时间
  48. Source string `json:"source"` // 数据来源
  49. }
  50. func (s *SentinelServer) GetLatestPlateData(limit int) ([]PlateData, error) {
  51. if s.rdb == nil {
  52. return nil, fmt.Errorf("Redis客户端未初始化")
  53. }
  54. // 步骤1:从ZSet获取最新的limit个timestamp(按Score倒序)
  55. zsetKey := fmt.Sprintf("%s:sorted", RedisKeyPrefix)
  56. // ZREVRANGE:倒序取0到limit-1的member(最新的N个)
  57. timestamps, err := s.rdb.ZRevRange(s.ctx, zsetKey, 0, int64(limit-1)).Result()
  58. if err != nil {
  59. return nil, fmt.Errorf("读取ZSet失败: %w", err)
  60. }
  61. if len(timestamps) == 0 {
  62. return nil, fmt.Errorf("Redis中无车牌数据")
  63. }
  64. // 步骤2:批量从Hash读取对应数据
  65. hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix)
  66. pipe := s.rdb.Pipeline()
  67. for _, ts := range timestamps {
  68. pipe.HGet(s.ctx, hashKey, ts)
  69. }
  70. results, err := pipe.Exec(s.ctx)
  71. if err != nil {
  72. return nil, fmt.Errorf("批量读取Hash失败: %w", err)
  73. }
  74. // 步骤3:解析Python格式的字符串为Go结构体
  75. var plateDataList []PlateData
  76. for i, res := range results {
  77. if res == nil {
  78. continue
  79. }
  80. // 获取Hash的value(Python的str(dict)字符串)
  81. pyStr, err := res.(*redis.StringCmd).Result()
  82. if err != nil {
  83. logrus.Warn("解析单条数据失败", "timestamp", timestamps[i], "err", err)
  84. continue
  85. }
  86. // 转换Python字符串为标准JSON,再解析
  87. plateData, err := parsePythonDictStr(pyStr)
  88. if err != nil {
  89. logrus.Warn("转换Python字符串失败", "str", pyStr, "err", err)
  90. continue
  91. }
  92. plateDataList = append(plateDataList, plateData)
  93. }
  94. return plateDataList, nil
  95. }
  96. // 2. 读取指定时间戳的车牌数据(精准读取)
  97. func (s *SentinelServer) GetPlateDataByTimestamp(timestamp string) (*PlateData, error) {
  98. if s.rdb == nil {
  99. return nil, fmt.Errorf("Redis客户端未初始化")
  100. }
  101. hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix)
  102. pyStr, err := s.rdb.HGet(s.ctx, hashKey, timestamp).Result()
  103. if err != nil {
  104. return nil, fmt.Errorf("读取指定时间戳数据失败: %w", err)
  105. }
  106. // 解析为结构体
  107. plateData, err := parsePythonDictStr(pyStr)
  108. if err != nil {
  109. return nil, fmt.Errorf("解析数据失败: %w", err)
  110. }
  111. return &plateData, nil
  112. }
  113. // ========== 辅助函数:解析Python的dict字符串为Go结构体 ==========
  114. // Python的str(dict)是单引号,需替换为双引号,且处理特殊字符
  115. func parsePythonDictStr(pyStr string) (PlateData, error) {
  116. var plateData PlateData
  117. // 1. 替换单引号为双引号(Python -> JSON)
  118. jsonStr := strings.ReplaceAll(pyStr, "'", "\"")
  119. // 2. 处理可能的空格/换行(可选,视Python字符串格式)
  120. jsonStr = regexp.MustCompile(`\s+`).ReplaceAllString(jsonStr, "")
  121. // 3. 解析为JSON
  122. err := json.Unmarshal([]byte(jsonStr), &plateData)
  123. if err != nil {
  124. return plateData, fmt.Errorf("JSON解析失败: %w, raw_str: %s", err, jsonStr)
  125. }
  126. return plateData, nil
  127. }
  128. // ========== 新增:存入 DeviceData 结构体到Redis ==========
  129. func (s *SentinelServer) SetDeviceDataToRedis(data DeviceData) error {
  130. if s.rdb == nil {
  131. return fmt.Errorf("redis客户端未初始化")
  132. }
  133. // 1. 将结构体序列化为JSON字节流
  134. jsonData, err := json.Marshal(data)
  135. if err != nil {
  136. return fmt.Errorf("结构体序列化失败: %w", err)
  137. }
  138. // 2. 存入Redis (永久有效,重启不丢失)
  139. err = s.rdb.Set(s.ctx, RedisInitDevice, jsonData, 0).Err()
  140. if err != nil {
  141. return fmt.Errorf("redis存入失败: %w", err)
  142. }
  143. logrus.Info("✅ 设备配置已存入Redis", "data", data)
  144. return nil
  145. }
  146. // ========== 新增:从Redis读取 DeviceData 结构体 ==========
  147. func (s *SentinelServer) GetDeviceDataFromRedis() (DeviceData, error) {
  148. // 默认值
  149. data := DeviceData{
  150. LowSpeed: 10,
  151. OverSpeed: 60,
  152. Brightness: 4,
  153. Volume: 4,
  154. NormalVoice: "注意来车",
  155. OverSpeedVoice: "您已超速",
  156. }
  157. if s.rdb == nil {
  158. return data, fmt.Errorf("redis客户端未初始化")
  159. }
  160. // 1. 从Redis读取JSON字符串
  161. jsonData, err := s.rdb.Get(s.ctx, RedisInitDevice).Result()
  162. // 特殊处理:redis中无此key时,返回默认空结构体,不报错
  163. if err == redis.Nil {
  164. logrus.Warn("⚠️ Redis中暂无设备配置数据,返回默认值")
  165. go s.SetDeviceDataToRedis(data)
  166. return data, nil
  167. }
  168. if err != nil {
  169. return data, fmt.Errorf("redis读取失败: %w", err)
  170. }
  171. // 2. JSON反序列化为结构体
  172. err = json.Unmarshal([]byte(jsonData), &data)
  173. if err != nil {
  174. return data, fmt.Errorf("结构体反序列化失败: %w", err)
  175. }
  176. return data, nil
  177. }