MessageHandler.java 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174
  1. package com.zksy.telemetry.utils;
  2. import com.zksy.common.exception.InvalidMessageException;
  3. import com.zksy.telemetry.domain.TelemetryData;
  4. import com.zksy.telemetry.service.TelemetryDataService;
  5. import io.netty.buffer.ByteBuf;
  6. import io.netty.buffer.Unpooled;
  7. import io.netty.channel.ChannelHandler;
  8. import io.netty.channel.ChannelHandlerContext;
  9. import io.netty.channel.ChannelInboundHandlerAdapter;
  10. import io.netty.handler.timeout.ReadTimeoutException;
  11. import io.netty.util.ReferenceCountUtil;
  12. import lombok.extern.slf4j.Slf4j;
  13. import org.slf4j.Logger;
  14. import org.slf4j.LoggerFactory;
  15. import org.springframework.beans.factory.annotation.Autowired;
  16. import org.springframework.stereotype.Component;
  17. import java.util.Date;
  18. @ChannelHandler.Sharable
  19. @Slf4j
  20. @Component
  21. public class MessageHandler extends ChannelInboundHandlerAdapter {
  22. private static Logger logger = LoggerFactory.getLogger(MessageHandler.class);
  23. private final TelemetryDataService service;
  24. @Autowired
  25. private DeviceOfflineCheckTask deviceOfflineCheckTask;
  26. @Autowired
  27. public MessageHandler(TelemetryDataService telemetryDataService) {
  28. this.service = telemetryDataService;
  29. }
  30. @Override
  31. public void channelActive(ChannelHandlerContext ctx) throws Exception {
  32. super.channelActive(ctx);
  33. //sendDataToDevice(ctx);
  34. }
  35. @Override
  36. public void channelRead(ChannelHandlerContext ctx, Object msg) {
  37. ByteBuf in = (ByteBuf) msg;
  38. byte[] msgBytes = new byte[in.readableBytes()];
  39. in.readBytes(msgBytes);
  40. in.release();
  41. try {
  42. logger.debug("接收到的原始数据帧: {}", printHexBinary(msgBytes));
  43. // 提取帧类型
  44. byte frameType = msgBytes[6];
  45. logger.debug("帧类型: 0x{}", String.format("%02X", frameType));
  46. // 1. 完整协议校验
  47. DataParser.validateMessage(msgBytes);
  48. logger.info("数据帧校验通过");
  49. // 2. 数据解析
  50. TelemetryData resultData = DataParser.parseMessage(msgBytes);
  51. if (resultData.getSystemIdentifier() == null) {
  52. resultData.setSystemIdentifier("123456");
  53. }
  54. // 3. 根据帧类型处理
  55. if (frameType == 0x31) {
  56. if (msgBytes.length > 34) {
  57. byte flagByte = msgBytes[34];
  58. String binary = String.format("%8s", Integer.toBinaryString(flagByte & 0xFF)).replace(' ', '0');
  59. boolean isEnd = (flagByte & 0x08) != 0;
  60. logger.debug("结束位检测: 字节=0x{}, 二进制={}, 第5位(结束标识)={}",
  61. String.format("%02X", flagByte), binary, isEnd ? "1(结束)" : "0(未结束)");
  62. } else {
  63. logger.warn("数据长度不足,无法检测结束位");
  64. }
  65. // 入库
  66. service.save(resultData);
  67. // 更新设备最后接收时间,并标记设备在线
  68. DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(resultData.getSystemIdentifier(), new Date());
  69. deviceOfflineCheckTask.markDeviceOnline(resultData.getSystemIdentifier());
  70. logger.info("上报历史记录数据入库成功: {}", resultData);
  71. // 如果结束位=1,可以根据需要在此处理关闭逻辑
  72. if (Boolean.TRUE.equals(resultData.getIsLastPacket())) {
  73. logger.info("检测到结束标志帧,执行结束处理逻辑...");
  74. byte[] response = ProtocolUtils.buildCustomReplyFrame(msgBytes);
  75. ctx.writeAndFlush(Unpooled.copiedBuffer(response));
  76. Thread.sleep(1000);// 休眠1秒
  77. logger.debug("已发送结束应答帧上报自定义回应包: {}", printHexBinary(response));
  78. byte[] endResponse = ProtocolUtils.buildEndReplyFrame(msgBytes);
  79. ctx.writeAndFlush(Unpooled.copiedBuffer(endResponse));
  80. logger.debug("已发送结束应答帧: {}", printHexBinary(endResponse));
  81. }else{
  82. // 回复自定义应答帧
  83. byte[] response = ProtocolUtils.buildCustomReplyFrame(msgBytes);
  84. ctx.writeAndFlush(Unpooled.copiedBuffer(response));
  85. logger.debug("已发送上报自定义回应包: {}", printHexBinary(response));
  86. }
  87. } else if (frameType == 0x34) {
  88. /*logger.info("收到结束通讯帧,准备发送回应包");
  89. byte[] response = ProtocolUtils.buildShutdownAckPacket(msgBytes);
  90. ctx.writeAndFlush(Unpooled.copiedBuffer(response));
  91. logger.debug("已发送结束通讯回应包");*/
  92. } else {
  93. logger.warn("收到未知帧类型: 0x{}", String.format("%02X", frameType));
  94. }
  95. } catch (InvalidMessageException e) {
  96. logger.error("数据校验失败,不入库: {}", e.getMessage());
  97. sendErrorResponse(ctx, "数据校验失败");
  98. } catch (Exception e) {
  99. logger.error("数据解析或入库异常", e);
  100. sendErrorResponse(ctx, "数据处理异常");
  101. }
  102. }
  103. // 工具方法:发送错误响应
  104. private void sendErrorResponse(ChannelHandlerContext ctx, String msg) {
  105. ctx.writeAndFlush(Unpooled.copiedBuffer(("数据处理失败: " + msg).getBytes()));
  106. }
  107. private String printHexBinary(byte[] bytes) {
  108. StringBuilder sb = new StringBuilder();
  109. for (byte b : bytes) {
  110. sb.append(String.format("%02X ", b));
  111. }
  112. System.out.println("Received raw data: " + sb.toString());
  113. logger.info("Received raw data: " + sb.toString());
  114. return sb.toString();
  115. }
  116. @Override
  117. public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
  118. if (cause instanceof ReadTimeoutException) {
  119. logger.info("来自" + ctx.channel().remoteAddress() + "的连接超时断开");
  120. } else {
  121. cause.printStackTrace();
  122. logger.info("来自" + ctx.channel().remoteAddress() + "的连接异常断开");
  123. ctx.close();
  124. }
  125. }
  126. @Override
  127. public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
  128. ctx.flush();
  129. }
  130. @Override
  131. public void channelUnregistered(ChannelHandlerContext ctx) throws Exception {
  132. logger.info("来自" + ctx.channel().remoteAddress() + "的连接主动断开");
  133. ctx.fireChannelUnregistered();
  134. }
  135. // 主动向设备发送数据
  136. public void sendDataToDevice(ChannelHandlerContext ctx) {
  137. try {
  138. // 将字符串形式的十六进制数据转换为字节数组
  139. String hexData = "01 03 00 00 00 01 84 0A ";
  140. String[] hexArray = hexData.split(" ");
  141. byte[] dataBytes = new byte[hexArray.length];
  142. for (int i = 0; i < hexArray.length; i++) {
  143. dataBytes[i] = (byte) Integer.parseInt(hexArray[i], 16);
  144. }
  145. ByteBuf byteBuf = Unpooled.copiedBuffer(dataBytes);
  146. ctx.writeAndFlush(byteBuf);
  147. logger.info("已向设备发送数据: {}", printHexBinary(dataBytes));
  148. } catch (Exception e) {
  149. logger.error("发送数据失败: {}", e.getMessage());
  150. }
  151. }
  152. }