MessageHandler.java 3.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394
  1. package com.zksy.manhole.utils;
  2. import com.zksy.common.exception.InvalidMessageException;
  3. import com.zksy.manhole.domain.ManholeData;
  4. import com.zksy.manhole.service.ManholeDataService;
  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 lombok.extern.slf4j.Slf4j;
  12. import org.slf4j.Logger;
  13. import org.slf4j.LoggerFactory;
  14. import org.springframework.stereotype.Component;
  15. import java.util.Date;
  16. @ChannelHandler.Sharable
  17. @Slf4j
  18. @Component
  19. public class MessageHandler extends ChannelInboundHandlerAdapter {
  20. private static Logger logger = LoggerFactory.getLogger(MessageHandler.class);
  21. private final ManholeDataService messageParseResultService;
  22. public MessageHandler(ManholeDataService messageParseResultService) {
  23. this.messageParseResultService = messageParseResultService;
  24. }
  25. @Override
  26. public void channelRead(ChannelHandlerContext ctx, Object msg) {
  27. try {
  28. ByteBuf msgByteBuf = (ByteBuf) msg;
  29. if (msgByteBuf == null || !msgByteBuf.isReadable()) {
  30. logger.warn("接收到无效的消息");
  31. return;
  32. }
  33. logger.info("接收到 {} 字节的数据,来自: {}", msgByteBuf.readableBytes(), ctx.channel().remoteAddress());
  34. byte[] msgBytes = new byte[msgByteBuf.readableBytes()];
  35. msgByteBuf.readBytes(msgBytes);
  36. String result = printHexBinary(msgBytes);
  37. String msgString = result.replaceAll("\\s", "");
  38. String CRCString = msgString.substring(msgString.length() - 4);
  39. String bodyResult = msgString.substring(0, msgString.length() - 4);
  40. String codeString = DataCheckUtil.crc16(bodyResult);
  41. System.out.println("校验码:" + CRCString + "-----------生成的校验码:" + codeString);
  42. if (!CRCString.equals(codeString)) {
  43. throw new InvalidMessageException("数据校验不成功");
  44. } else {
  45. ManholeData resultData = DataParser.parseMessage(msgString);
  46. // 更新设备最后一次接收数据的时间
  47. String imeiCardNumber = resultData.getImeiCardNumber();
  48. DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(imeiCardNumber, new Date());
  49. messageParseResultService.saveManholeData(resultData);
  50. }
  51. } catch (InvalidMessageException e) {
  52. logger.error("数据入库失败: {}", e.getMessage());
  53. ctx.writeAndFlush(Unpooled.copiedBuffer("数据入库失败".getBytes()));
  54. }
  55. }
  56. private String printHexBinary(byte[] bytes) {
  57. StringBuilder sb = new StringBuilder();
  58. for (byte b : bytes) {
  59. sb.append(String.format("%02X ", b));
  60. }
  61. System.out.println("Received raw data: " + sb.toString());
  62. logger.info("Received raw data: " + sb.toString());
  63. return sb.toString();
  64. }
  65. @Override
  66. public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
  67. if (cause instanceof ReadTimeoutException) {
  68. logger.info("来自" + ctx.channel().remoteAddress() + "的连接超时断开");
  69. } else {
  70. cause.printStackTrace();
  71. logger.info("来自" + ctx.channel().remoteAddress() + "的连接异常断开");
  72. ctx.close();
  73. }
  74. }
  75. @Override
  76. public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
  77. ctx.flush();
  78. }
  79. @Override
  80. public void channelUnregistered(ChannelHandlerContext ctx) throws Exception {
  81. logger.info("来自" + ctx.channel().remoteAddress() + "的连接主动断开");
  82. ctx.fireChannelUnregistered();
  83. }
  84. }