| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394 |
- package com.zksy.manhole.utils;
- import com.zksy.common.exception.InvalidMessageException;
- import com.zksy.manhole.domain.ManholeData;
- import com.zksy.manhole.service.ManholeDataService;
- import io.netty.buffer.ByteBuf;
- import io.netty.buffer.Unpooled;
- import io.netty.channel.ChannelHandler;
- import io.netty.channel.ChannelHandlerContext;
- import io.netty.channel.ChannelInboundHandlerAdapter;
- import io.netty.handler.timeout.ReadTimeoutException;
- import lombok.extern.slf4j.Slf4j;
- import org.slf4j.Logger;
- import org.slf4j.LoggerFactory;
- import org.springframework.stereotype.Component;
- import java.util.Date;
- @ChannelHandler.Sharable
- @Slf4j
- @Component
- public class MessageHandler extends ChannelInboundHandlerAdapter {
- private static Logger logger = LoggerFactory.getLogger(MessageHandler.class);
- private final ManholeDataService messageParseResultService;
- public MessageHandler(ManholeDataService messageParseResultService) {
- this.messageParseResultService = messageParseResultService;
- }
- @Override
- public void channelRead(ChannelHandlerContext ctx, Object msg) {
- try {
- ByteBuf msgByteBuf = (ByteBuf) msg;
- if (msgByteBuf == null || !msgByteBuf.isReadable()) {
- logger.warn("接收到无效的消息");
- return;
- }
- logger.info("接收到 {} 字节的数据,来自: {}", msgByteBuf.readableBytes(), ctx.channel().remoteAddress());
- byte[] msgBytes = new byte[msgByteBuf.readableBytes()];
- msgByteBuf.readBytes(msgBytes);
- String result = printHexBinary(msgBytes);
- String msgString = result.replaceAll("\\s", "");
- String CRCString = msgString.substring(msgString.length() - 4);
- String bodyResult = msgString.substring(0, msgString.length() - 4);
- String codeString = DataCheckUtil.crc16(bodyResult);
- System.out.println("校验码:" + CRCString + "-----------生成的校验码:" + codeString);
- if (!CRCString.equals(codeString)) {
- throw new InvalidMessageException("数据校验不成功");
- } else {
- ManholeData resultData = DataParser.parseMessage(msgString);
- // 更新设备最后一次接收数据的时间
- String imeiCardNumber = resultData.getImeiCardNumber();
- DeviceOfflineCheckTask.deviceLastReceiveTimeMap.put(imeiCardNumber, new Date());
- messageParseResultService.saveManholeData(resultData);
- }
- } catch (InvalidMessageException e) {
- logger.error("数据入库失败: {}", e.getMessage());
- ctx.writeAndFlush(Unpooled.copiedBuffer("数据入库失败".getBytes()));
- }
- }
- private String printHexBinary(byte[] bytes) {
- StringBuilder sb = new StringBuilder();
- for (byte b : bytes) {
- sb.append(String.format("%02X ", b));
- }
- System.out.println("Received raw data: " + sb.toString());
- logger.info("Received raw data: " + sb.toString());
- return sb.toString();
- }
- @Override
- public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
- if (cause instanceof ReadTimeoutException) {
- logger.info("来自" + ctx.channel().remoteAddress() + "的连接超时断开");
- } else {
- cause.printStackTrace();
- logger.info("来自" + ctx.channel().remoteAddress() + "的连接异常断开");
- ctx.close();
- }
- }
- @Override
- public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
- ctx.flush();
- }
- @Override
- public void channelUnregistered(ChannelHandlerContext ctx) throws Exception {
- logger.info("来自" + ctx.channel().remoteAddress() + "的连接主动断开");
- ctx.fireChannelUnregistered();
- }
- }
|