electrical_data_fetcher.go 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331
  1. package common
  2. import (
  3. "Tech/models"
  4. "bytes"
  5. "crypto/md5"
  6. "encoding/hex"
  7. "encoding/json"
  8. "fmt"
  9. "io"
  10. "net/http"
  11. "time"
  12. "github.com/beego/beego/v2/client/orm"
  13. beego "github.com/beego/beego/v2/server/web"
  14. "github.com/sirupsen/logrus"
  15. )
  16. // EquipmentData 表示设备数据
  17. type EquipmentData struct {
  18. DataID string `json:"dataId"`
  19. DataName string `json:"dataName"`
  20. Value string `json:"value"`
  21. ID string `json:"id"`
  22. EqpmtTypeID int `json:"eqpmtTypeID"`
  23. }
  24. // ElectricalDataFetcher 电调数据获取器
  25. type ElectricalDataFetcher struct {
  26. manager *models.ElectricModelManager
  27. electricalIP string
  28. electricalPort int
  29. lastDataHash string
  30. lastEquipments []map[string]interface{} // 记录上次广播的设备数据
  31. isRunning bool
  32. }
  33. var fetcher *ElectricalDataFetcher
  34. // InitElectricalDataFetcher 初始化电调数据获取器
  35. func InitElectricalDataFetcher() {
  36. fetcher = &ElectricalDataFetcher{
  37. manager: models.GetInstance(),
  38. isRunning: false,
  39. lastEquipments: nil, // 初始化为空,第一次会推送全量数据
  40. }
  41. // 从配置文件读取电调服务配置
  42. electricalIP, err := beego.AppConfig.String("ElectricalIP")
  43. if err != nil {
  44. LogError("读取电调服务IP配置失败", err)
  45. return
  46. }
  47. electricalPort, err := beego.AppConfig.Int("ElectricalPort")
  48. if err != nil {
  49. LogError("读取电调服务端口配置失败", err)
  50. return
  51. }
  52. fetcher.electricalIP = electricalIP
  53. fetcher.electricalPort = electricalPort
  54. // 启动WebSocket管理器
  55. GlobalWSManager.Start()
  56. // 启动定时获取数据任务
  57. fetcher.Start()
  58. }
  59. // Start 启动定时获取数据任务
  60. func (f *ElectricalDataFetcher) Start() {
  61. if f.isRunning {
  62. return
  63. }
  64. f.isRunning = true
  65. // 获取全量数据更新间隔配置,默认5秒
  66. intervalSeconds := 5
  67. if interval, err := beego.AppConfig.Int("EGetAllDataInterval"); err == nil {
  68. intervalSeconds = interval
  69. }
  70. ticker := time.NewTicker(time.Duration(intervalSeconds) * time.Second)
  71. go func() {
  72. // 初始延迟1秒执行,避免系统启动时的资源竞争
  73. time.Sleep(1 * time.Second)
  74. for range ticker.C {
  75. f.fetchAndBroadcastData()
  76. }
  77. }()
  78. LogInfo(fmt.Sprintf("电调数据定时获取任务已启动,间隔 %d 秒", intervalSeconds))
  79. }
  80. // fetchAndBroadcastData 获取并广播数据
  81. func (f *ElectricalDataFetcher) fetchAndBroadcastData() {
  82. defer func() {
  83. if r := recover(); r != nil {
  84. LogError("获取电调数据时发生异常", nil, logrus.Fields{
  85. "panic": r,
  86. })
  87. }
  88. }()
  89. // 获取最新的课程运行记录以获取线路号
  90. o := orm.NewOrm()
  91. var latestCourse models.CourseRunLog
  92. err := o.QueryTable(new(models.CourseRunLog)).
  93. OrderBy("-RunLogID").
  94. One(&latestCourse)
  95. lineID := 0
  96. if err == nil && err != orm.ErrNoRows {
  97. // 根据运行课程ID查询课程信息以获取线路号
  98. var courseInfo models.CourseInfo
  99. err = o.QueryTable(new(models.CourseInfo)).
  100. Filter("CourseID", latestCourse.RunCourseID).
  101. One(&courseInfo)
  102. if err == nil {
  103. lineID = courseInfo.LineNo
  104. }
  105. }
  106. // 构造电调服务地址
  107. electricalAddr := fmt.Sprintf("http://%s:%d/ServerCtrl/AllEqpmtList", f.electricalIP, f.electricalPort)
  108. // 准备请求数据
  109. type RequestData struct {
  110. LineID int `json:"lineID"`
  111. StationID int `json:"StationID"`
  112. ViewID int `json:"ViewID"`
  113. }
  114. reqData := RequestData{
  115. LineID: lineID,
  116. StationID: 0,
  117. ViewID: 0,
  118. }
  119. // 将请求数据转换为JSON
  120. reqBody, err := json.Marshal(reqData)
  121. if err != nil {
  122. LogError("构造请求数据失败", err)
  123. return
  124. }
  125. // 调用电调服务获取所有设备列表接口
  126. resp, err := http.Post(electricalAddr, "application/json", bytes.NewBuffer(reqBody))
  127. if err != nil {
  128. LogError("调用电调服务获取所有设备列表接口失败", err)
  129. return
  130. }
  131. defer resp.Body.Close()
  132. // 检查响应状态
  133. if resp.StatusCode != http.StatusOK {
  134. bodyBytes, _ := io.ReadAll(resp.Body)
  135. LogError("电调服务返回异常状态码", fmt.Errorf("状态码: %d, 响应: %s", resp.StatusCode, string(bodyBytes)))
  136. return
  137. }
  138. // 读取并解析响应体
  139. bodyBytes, err := io.ReadAll(resp.Body)
  140. if err != nil {
  141. LogError("读取电调服务响应失败", err)
  142. return
  143. }
  144. // 解析JSON响应
  145. var result map[string]interface{}
  146. if err := json.Unmarshal(bodyBytes, &result); err != nil {
  147. LogError("解析电调服务JSON响应失败", err)
  148. return
  149. }
  150. // 检查是否有data字段
  151. data, exists := result["data"]
  152. if !exists {
  153. LogError("电调服务响应中没有data字段", nil)
  154. return
  155. }
  156. // 解析data中的设备列表
  157. var equipments []models.DiandiaoModel
  158. if equipBytes, err := json.Marshal(data); err == nil {
  159. if err := json.Unmarshal(equipBytes, &equipments); err != nil {
  160. LogError("解析设备列表数据失败", err)
  161. return
  162. }
  163. }
  164. // 更新模型管理器的总包设备数据
  165. f.manager.UpdateBatchData(equipments)
  166. // 转换为标准设备数据格式
  167. equipmentsData := f.manager.ConvertToEquipmentData(equipments)
  168. // 加入测试代码:生成随机值修改第一个设备数据
  169. // if len(equipmentsData) > 0 {
  170. // // 生成0-100之间的随机值
  171. // randomValue := time.Now().Unix() % 101
  172. // equipmentsData[0]["value"] = fmt.Sprintf("%d", randomValue)
  173. // LogDebug(fmt.Sprintf("测试:修改第一个设备数据的值为 %d", randomValue), nil)
  174. // }
  175. // 比较数据差异,获取变化的设备
  176. changedEquipments := f.compareWithLastData(equipmentsData)
  177. // 如果没有变化,则跳过广播
  178. if len(changedEquipments) == 0 {
  179. return
  180. }
  181. // 更新上次广播的数据
  182. f.lastEquipments = equipmentsData
  183. // 广播变化的设备数据到所有WebSocket连接
  184. GlobalWSManager.Broadcast(changedEquipments)
  185. LogDebug(fmt.Sprintf("成功获取并广播变化设备数据,变化数量: %d", len(changedEquipments)), logrus.Fields{
  186. "lineID": lineID,
  187. "totalCount": len(equipmentsData),
  188. "changedCount": len(changedEquipments),
  189. })
  190. }
  191. // compareWithLastData 比较当前数据与上次数据的差异
  192. func (f *ElectricalDataFetcher) compareWithLastData(currentData []map[string]interface{}) []map[string]interface{} {
  193. // 如果是第一次(没有上次数据),返回所有数据
  194. if f.lastEquipments == nil {
  195. return currentData
  196. }
  197. // 构建上次数据的映射
  198. lastDataMap := make(map[string]map[string]interface{})
  199. for _, equip := range f.lastEquipments {
  200. if id, ok := equip["id"].(string); ok {
  201. lastDataMap[id] = equip
  202. }
  203. }
  204. // 构建当前数据的映射
  205. currentDataMap := make(map[string]map[string]interface{})
  206. for _, equip := range currentData {
  207. if id, ok := equip["id"].(string); ok {
  208. currentDataMap[id] = equip
  209. }
  210. }
  211. // 查找变化的设备
  212. var changedEquipments []map[string]interface{}
  213. allIDs := make(map[string]bool)
  214. // 检查上次数据中的设备
  215. for id, lastEquip := range lastDataMap {
  216. allIDs[id] = true
  217. currentEquip, exists := currentDataMap[id]
  218. if !exists {
  219. // 设备被删除了,不推送删除的设备
  220. continue
  221. }
  222. // 比较值是否变化
  223. if lastEquip["value"] != currentEquip["value"] {
  224. changedEquipments = append(changedEquipments, currentEquip)
  225. }
  226. }
  227. // 检查新增加的设备
  228. for id, currentEquip := range currentDataMap {
  229. if !allIDs[id] {
  230. changedEquipments = append(changedEquipments, currentEquip)
  231. }
  232. }
  233. return changedEquipments
  234. }
  235. // calculateDataHash 计算数据哈希
  236. func (f *ElectricalDataFetcher) calculateDataHash(data []map[string]interface{}) string {
  237. // 将数据序列化为JSON
  238. jsonData, err := json.Marshal(data)
  239. if err != nil {
  240. return ""
  241. }
  242. // 计算MD5哈希
  243. hash := md5.Sum(jsonData)
  244. return hex.EncodeToString(hash[:])
  245. }
  246. // GetCurrentData 获取当前缓存的设备数据
  247. func (f *ElectricalDataFetcher) GetCurrentData() ([]map[string]interface{}, error) {
  248. // 获取所有设备数据
  249. equipments := f.manager.GetAllEquipmentData()
  250. // 转换为标准设备数据格式
  251. equipmentsData := f.manager.ConvertToEquipmentData(equipments)
  252. return equipmentsData, nil
  253. }
  254. // SendCurrentDataToClient 发送当前数据到指定客户端
  255. func SendCurrentDataToClient(clientID string) {
  256. if fetcher == nil {
  257. LogError("电调数据获取器未初始化", nil)
  258. return
  259. }
  260. // 获取当前数据
  261. equipments, err := fetcher.GetCurrentData()
  262. if err != nil {
  263. LogError("获取当前设备数据失败", err)
  264. return
  265. }
  266. // 发送到指定客户端(直接发送设备数组)
  267. if err := GlobalWSManager.BroadcastToClient(clientID, equipments); err != nil {
  268. LogError("发送设备数据到客户端失败", err, logrus.Fields{
  269. "clientID": clientID,
  270. })
  271. return
  272. }
  273. LogInfo(fmt.Sprintf("成功发送全量设备数据到客户端 %s,设备数量: %d", clientID, len(equipments)))
  274. }