mbdevhandler.go 8.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292
  1. package main
  2. import (
  3. "runtime"
  4. "runtime/debug"
  5. "sync"
  6. "time"
  7. "github.com/go-redis/redis/v7"
  8. jsoniter "github.com/json-iterator/go"
  9. "github.com/sirupsen/logrus"
  10. "lc/common/models"
  11. "lc/common/mqtt"
  12. "lc/common/protocol"
  13. "lc/common/util"
  14. )
  15. var _modbusDeviceHandlerOnce sync.Once
  16. var _modbusDeviceHandlerSingle *ModbusDeviceHandler
  17. func GetModbusDeviceHandler() *ModbusDeviceHandler {
  18. _modbusDeviceHandlerOnce.Do(func() {
  19. _modbusDeviceHandlerSingle = &ModbusDeviceHandler{
  20. queue: util.NewQueue(10000),
  21. }
  22. })
  23. return _modbusDeviceHandlerSingle
  24. }
  25. type ModbusDeviceHandler struct {
  26. queue *util.MlQueue
  27. }
  28. func (o *ModbusDeviceHandler) SubscribeTopics() {
  29. //环境传感器
  30. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_ENVIRONMENT, protocol.TP_MODBUS_CONTROL_ACK), mqtt.AtMostOnce, o.HandlerData)
  31. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_ENVIRONMENT, protocol.TP_MODBUS_DATA), mqtt.AtMostOnce, o.HandlerData)
  32. //液位计
  33. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LIQUID, protocol.TP_MODBUS_CONTROL_ACK), mqtt.AtMostOnce, o.HandlerData)
  34. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LIQUID, protocol.TP_MODBUS_DATA), mqtt.AtMostOnce, o.HandlerData)
  35. //路况传感器
  36. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_ROAD_COND, protocol.TP_MODBUS_CONTROL_ACK), mqtt.AtMostOnce, o.HandlerData)
  37. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_ROAD_COND, protocol.TP_MODBUS_DATA), mqtt.AtMostOnce, o.HandlerData)
  38. }
  39. func (o *ModbusDeviceHandler) HandlerData(m mqtt.Message) {
  40. for {
  41. ok, cnt := o.queue.Put(&m)
  42. if ok {
  43. break
  44. } else {
  45. logrus.Errorf("ModbusDeviceHandler.HandlerData:查询队列失败,队列消息数量:%d", cnt)
  46. runtime.Gosched()
  47. }
  48. }
  49. }
  50. func (o *ModbusDeviceHandler) Handler(args ...interface{}) interface{} {
  51. defer func() {
  52. if err := recover(); err != nil {
  53. time.Sleep(time.Second)
  54. gopool.Add(o.Handler, args)
  55. logrus.Errorf("ModbusDeviceHandler.Handler发生异常:%s", string(debug.Stack()))
  56. }
  57. }()
  58. timer := time.NewTicker(1 * time.Minute)
  59. for {
  60. select {
  61. case <-timer.C: //每隔5分钟执行一次 状态更新
  62. mapMbDevData.Range(func(key, value interface{}) bool {
  63. p, ok := value.(*MbDevData)
  64. if ok {
  65. p.UpdateState()
  66. }
  67. return true
  68. })
  69. default:
  70. msg, ok, quantity := o.queue.Get()
  71. if !ok {
  72. time.Sleep(10 * time.Millisecond)
  73. continue
  74. } else if quantity > 1000 {
  75. logrus.Warnf("ModbusDeviceHandler.Handler:数据队列累积过多,请注意优化,当前队列条数:%d", quantity)
  76. }
  77. m, ok := msg.(*mqtt.Message)
  78. if !ok {
  79. continue
  80. }
  81. Tenant, DevType, DID, topic, err := ParseTopic(m.Topic())
  82. if err != nil {
  83. continue
  84. }
  85. switch topic {
  86. case protocol.TP_MODBUS_DATA:
  87. var obj protocol.Pack_UploadData
  88. if err := obj.DeCode(m.PayloadString()); err == nil {
  89. var pMbDevData *MbDevData = nil
  90. mgr, ok := mapMbDevData.Load(obj.Id)
  91. if ok {
  92. pMbDevData = mgr.(*MbDevData)
  93. } else {
  94. pMbDevData = NewMbDevData(Tenant, DevType, obj.Gid, DID)
  95. mapMbDevData.Store(obj.Id, pMbDevData)
  96. }
  97. if pMbDevData != nil {
  98. pMbDevData.HandleData(&obj)
  99. }
  100. }
  101. case protocol.TP_MODBUS_CONTROL_ACK:
  102. var obj protocol.Pack_Ack
  103. if err := obj.DeCode(m.PayloadString()); err == nil {
  104. o := models.DeviceCmdRecord{
  105. ID: obj.Seq,
  106. State: uint(obj.Data.State),
  107. Resp: obj.Data.Error,
  108. }
  109. if err := o.Update(); err != nil {
  110. logrus.Errorf("收到网关[%s]的响应[seq:%d],主题:%s,但更新数据库失败[%s]", obj.Id, obj.Seq, m.Topic(), err.Error())
  111. }
  112. logrus.Debugf("ModbusDeviceHandler.Handler:收到网关[%s]发布的主题:%s,内容:%s", obj.Id, m.Topic(), m.PayloadString())
  113. }
  114. default:
  115. logrus.Warnf("ModbusDeviceHandler.Handler:收到暂不支持的主题:%s", topic)
  116. }
  117. }
  118. }
  119. }
  120. const (
  121. DevStatusPrefix string = "dev_stat_"
  122. DevDataPrefix string = "dev_data_"
  123. ONLINE string = "online"
  124. TLast string = "tlast"
  125. TIME string = "time"
  126. )
  127. var json = jsoniter.ConfigFastest
  128. var mapMbDevData sync.Map
  129. type MbDevData struct {
  130. Lock sync.Mutex
  131. Tenant string //基础数据,租户ID
  132. Devtype string //基础数据,设备类型
  133. GID string //基础数据,网关编码
  134. DID string //基础数据,设备ID
  135. TID uint16 //基础数据,物模型ID
  136. LastDataTime time.Time //实时数据,最新数据时间
  137. Data map[uint16]float64 //实时数据,最新数据
  138. LastStateTime time.Time //实时数据,最新状态时间
  139. State uint8 //实时数据,0在线,1离线
  140. ErrCnt uint //错误计数
  141. NextHourTime time.Time //下次保存实时数据时间
  142. }
  143. func NewMbDevData(tenant, devtype, gid, did string) *MbDevData {
  144. LastHour := protocol.ToBJTime(util.BeginningOfHour().Add(1 * time.Hour))
  145. return &MbDevData{
  146. Tenant: tenant,
  147. Devtype: devtype,
  148. DID: did,
  149. GID: gid,
  150. Data: make(map[uint16]float64),
  151. State: 0xff,
  152. NextHourTime: LastHour,
  153. }
  154. }
  155. func (o *MbDevData) handleStateChange(t time.Time) {
  156. //最新状态
  157. state := uint8(0)
  158. if o.ErrCnt >= 10 {
  159. state = 1
  160. } else if o.ErrCnt == 0 {
  161. state = 0
  162. } else {
  163. return
  164. }
  165. //状态处理
  166. if o.LastStateTime.IsZero() || o.State == 0xff {
  167. t0, s0, err := getState(o.DID)
  168. if err != nil {
  169. o.State = state
  170. o.LastStateTime = t
  171. return
  172. }
  173. o.State = s0
  174. o.LastStateTime = t0
  175. }
  176. if o.State == 0 && state == 1 { //在线->离线
  177. GetEventMgr().PushEvent(&EventObject{ID: o.DID, EventType: models.ET_OFFLINE, Time: t})
  178. } else if o.State == 1 && state == 0 { //离线->在线
  179. GetEventMgr().PushEvent(&EventObject{ID: o.DID, EventType: models.ET_ONLINE, Time: t})
  180. }
  181. o.State = state
  182. o.LastStateTime = t
  183. }
  184. func (o *MbDevData) checkSaveData() {
  185. if o.NextHourTime.Before(o.LastDataTime) {
  186. //判断小时是否一样,不一样则以数据时间为准
  187. if !util.New(o.NextHourTime).BeginningOfHour().Equal(util.New(o.LastDataTime).BeginningOfHour()) {
  188. o.NextHourTime = util.New(o.LastDataTime).BeginningOfHour()
  189. }
  190. if len(o.Data) > 0 {
  191. var datas []models.DeviceHourData
  192. for k, v := range o.Data {
  193. o := models.DeviceHourData{ID: o.DID, Sid: k, Val: float32(v), Time: o.NextHourTime, CreatedAt: time.Now()}
  194. datas = append(datas, o)
  195. }
  196. if err := models.MultiInsertDeviceHourData(datas); err != nil {
  197. logrus.Errorf("小时[%s]数据插入数据库失败:%s", o.NextHourTime.Format("2006-01-02 15:04:05"), err.Error())
  198. }
  199. }
  200. o.NextHourTime = o.NextHourTime.Add(time.Hour)
  201. o.Data = make(map[uint16]float64)
  202. }
  203. }
  204. func (o *MbDevData) UpdateState() {
  205. if o.LastStateTime.IsZero() {
  206. t0, s0, err := getState(o.DID)
  207. if err == nil {
  208. o.State = s0
  209. o.LastStateTime = t0
  210. } else {
  211. err := redisCltRawData.HGet(DeviceAlarmId, o.DID).Err()
  212. if err == redis.Nil {
  213. o.State = protocol.FAILED
  214. o.LastStateTime = util.MlNow()
  215. GetEventMgr().PushEvent(&EventObject{ID: o.DID, EventType: models.ET_OFFLINE, Time: o.LastStateTime})
  216. }
  217. }
  218. }
  219. if o.State == protocol.FAILED ||
  220. (!o.LastDataTime.IsZero() && util.MlNow().Sub(o.LastDataTime).Minutes() < OfflineInterval) {
  221. return
  222. }
  223. //如果之前一直是在线状态的,则置为离线;若之前是离线状态的,则不修改状态
  224. if o.State == protocol.SUCCESS {
  225. o.State = protocol.FAILED
  226. o.LastStateTime = util.MlNow()
  227. GetEventMgr().PushEvent(&EventObject{ID: o.DID, EventType: models.ET_OFFLINE, Time: o.LastStateTime})
  228. cacheState(o.DID, o.LastStateTime.Format("2006-01-02 15:04:05"), o.State)
  229. //检查是否要缓存数据
  230. o.checkSaveData()
  231. }
  232. }
  233. func (o *MbDevData) HandleData(data *protocol.Pack_UploadData) {
  234. t, err := util.MlParseTime(data.Time)
  235. if err != nil {
  236. logrus.Errorf("时间[%s]解析错误:%s", data.Time, err.Error())
  237. return
  238. }
  239. o.Lock.Lock()
  240. defer o.Lock.Unlock()
  241. o.GID = data.Gid
  242. o.TID = data.Data.Tid //基础数据,物模型ID
  243. if len(data.Data.Data) > 0 {
  244. for k, v := range data.Data.Data {
  245. o.Data[k] = v
  246. }
  247. o.LastDataTime = t
  248. cacheData(o.DID, t, data.Data.Data)
  249. o.ErrCnt = 0
  250. } else {
  251. o.ErrCnt++
  252. }
  253. //需要告警,则推入告警管理器
  254. bv := BizValue{ID: o.DID, Time: t, Tid: o.TID, Data: data.Data.Data}
  255. GetBizAlarmMgr().PushData(&bv)
  256. //先处理状态变化,再存入最新状态
  257. o.handleStateChange(t)
  258. cacheState(o.DID, data.Time, o.State)
  259. o.checkSaveData()
  260. }