From 8f224dbf1a1bd1a65dead7ceda8dd0a3fa567115 Mon Sep 17 00:00:00 2001
From: ‘liusuyi’ <1951119284@qq.com>
Date: Thu, 28 Dec 2023 15:57:09 +0800
Subject: [PATCH] 优化雷达tcp客户端
---
/dev/null | 25 --------
src/main/java/com/ard/utils/netty/tcp/ClientInitialize.java | 52 +++++++----------
src/main/java/com/ard/utils/netty/tcp/ClientHandler.java | 67 +++++++++++----------
src/main/java/com/ard/alarm/radar/controller/RadarController.java | 31 +++++-----
src/main/java/com/ard/utils/netty/tcp/MessageHandler.java | 2
5 files changed, 73 insertions(+), 104 deletions(-)
diff --git a/src/main/java/com/ard/alarm/radar/controller/RadarController.java b/src/main/java/com/ard/alarm/radar/controller/RadarController.java
index 5cd3e6c..ad455e3 100644
--- a/src/main/java/com/ard/alarm/radar/controller/RadarController.java
+++ b/src/main/java/com/ard/alarm/radar/controller/RadarController.java
@@ -41,39 +41,38 @@
if (ardEquipRadar == null) {
return AjaxResult.error("雷达不存在");
}
- Channel channel = (Channel)ClientInitialize.SuccessConnectMap.get(ardEquipRadar.getId());
- if (channel==null)
- {
+ Channel channel = ClientInitialize.SucChannelMap.get(ardEquipRadar.getIp() + ":" + ardEquipRadar.getPort());
+ if (channel == null) {
return AjaxResult.error("雷达未连接");
}
Double longitude = ardEquipRadar.getLongitude();//雷达经度
Double latitude = ardEquipRadar.getLatitude();//雷达纬度
Double altitude = ardEquipRadar.getAltitude();//雷达高度
//计算水平角度
- float p = (float)GisUtils.getNorthAngle(longitude, latitude, targetPosition[0], targetPosition[1]);
+ float p = (float) GisUtils.getNorthAngle(longitude, latitude, targetPosition[0], targetPosition[1]);
//计算垂直角度
double[] radarPosition = new double[2];
radarPosition[0] = longitude;
radarPosition[1] = latitude;
double distance = GisUtils.getDistance(radarPosition, targetPosition);
- float angleInRadians = (float)Math.atan(distance / altitude);
- float t = (90-(float)Math.toDegrees(angleInRadians))*-1;
- log.debug("distance:"+distance);
- log.debug("p:"+p);
- log.debug("t:"+t);
+ float angleInRadians = (float) Math.atan(distance / altitude);
+ float t = (90 - (float) Math.toDegrees(angleInRadians)) * -1;
+ log.debug("distance:" + distance);
+ log.debug("p:" + p);
+ log.debug("t:" + t);
//发送告警前端的角度提示
byte[] header = {0x01, 0x02, 0x01};//包头
byte[] payloadHeader = {0x10, 0x03, 0x20, 0x00};//负载头
- byte[] distanceBytes = ByteUtils.decimalToBytes((int)distance);
- distanceBytes=ByteUtils.toLittleEndian(distanceBytes);
+ byte[] distanceBytes = ByteUtils.decimalToBytes((int) distance);
+ distanceBytes = ByteUtils.toLittleEndian(distanceBytes);
byte[] pBytes = ByteUtils.floatToBytes(p);
- pBytes=ByteUtils.toLittleEndian(pBytes);
+ pBytes = ByteUtils.toLittleEndian(pBytes);
byte[] tBytes = ByteUtils.floatToBytes(t);
- tBytes=ByteUtils.toLittleEndian(tBytes);
- byte[] resBytes=new byte[20];
- byte[] payloadBody = ByteUtils.appendArrays(distanceBytes,pBytes,tBytes,resBytes);//负载
+ tBytes = ByteUtils.toLittleEndian(tBytes);
+ byte[] resBytes = new byte[20];
+ byte[] payloadBody = ByteUtils.appendArrays(distanceBytes, pBytes, tBytes, resBytes);//负载
- byte[] payload = ByteUtils.appendArrays(payloadHeader,payloadBody);//负载头+负载
+ byte[] payload = ByteUtils.appendArrays(payloadHeader, payloadBody);//负载头+负载
byte[] payloadCrc32 = ByteUtils.parseCrc32(payload);//负载头+负载的crc32校验
byte[] footer = {0x01, 0x02, 0x00};//包尾
byte[] data = ByteUtils.appendArrays(header, payload, payloadCrc32, footer);
diff --git a/src/main/java/com/ard/utils/netty/tcp/BootNettyChannelInboundHandlerAdapter.java b/src/main/java/com/ard/utils/netty/tcp/BootNettyChannelInboundHandlerAdapter.java
deleted file mode 100644
index 5e49218..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/BootNettyChannelInboundHandlerAdapter.java
+++ /dev/null
@@ -1,439 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.alibaba.fastjson2.JSON;
-import com.ard.alarm.radar.domain.ArdAlarmRadar;
-import com.ard.alarm.radar.domain.ArdEquipRadar;
-import com.ard.alarm.radar.domain.RadarAlarmData;
-import com.ard.utils.mqtt.MqttProducer;
-import com.ard.utils.util.ByteUtils;
-import com.ard.utils.util.GisUtils;
-import io.netty.buffer.ByteBuf;
-import io.netty.buffer.EmptyByteBuf;
-import io.netty.channel.*;
-import io.netty.util.CharsetUtil;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.lang3.StringUtils;
-
-import javax.xml.bind.DatatypeConverter;
-import java.io.IOException;
-import java.net.InetSocketAddress;
-import java.text.SimpleDateFormat;
-import java.util.*;
-
-import static com.ard.utils.util.ByteUtils.byteToBitString;
-import static com.ard.utils.util.ByteUtils.toLittleEndian;
-
-@Slf4j(topic = "netty")
-public class BootNettyChannelInboundHandlerAdapter extends SimpleChannelInboundHandler<ByteBuf> {
-
- /**
- * 从服务端收到新的数据时,这个方法会在收到消息时被调用
- */
- @Override
- protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) {
- InetSocketAddress inSocket = (InetSocketAddress) ctx.channel().remoteAddress();
- String host = inSocket.getAddress().getHostAddress();
- int port = inSocket.getPort();
- //log.info("收到来自 {}:{} 的数据{}", host, port,msg.toString(CharsetUtil.UTF_8));
- ArdEquipRadar ardEquipRadar = BootNettyClientChannelCache.getRadar(host + ":" + port);
- if (ardEquipRadar != null) {
- //// 创建缓冲中字节数的字节数组
- //byte[] byteArray = new byte[msg.readableBytes()];
- //// 写入数组
- //msg.readBytes(byteArray);
- //// 处理接收到的消息
- //byte[] bytes = MessageParsing.receiveCompletePacket(byteArray);
- //if (bytes != null) {
- // processData(ardEquipRadar, bytes);
- //}
- }
- }
-
- /**
- * 从服务端收到新的数据、读取完成时调用
- */
- @Override
- public void channelReadComplete(ChannelHandlerContext ctx) throws IOException {
- //System.out.println("channelReadComplete");
- ctx.flush();
- }
-
- /**
- * 当出现 Throwable 对象才会被调用,即当 Netty 由于 IO 错误或者处理器在处理事件时抛出的异常时
- */
- @Override
- public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws IOException {
-// System.out.println("exceptionCaught");
- cause.printStackTrace();
- ctx.close();//抛出异常,断开与客户端的连接
- }
-
- /**
- * 客户端与服务端第一次建立连接时 执行
- */
- @Override
- public void channelActive(ChannelHandlerContext ctx) throws Exception {
- super.channelActive(ctx);
- // 客户端与服务端 建立连接
- InetSocketAddress inSocket = (InetSocketAddress) ctx.channel().remoteAddress();
- String host = inSocket.getAddress().getHostAddress();
- int port = inSocket.getPort();
- log.debug("连接成功:【" + host + ":" + port + "】");
- }
-
- /**
- * 客户端与服务端 断连时 执行
- */
- @Override
- public void channelInactive(ChannelHandlerContext ctx) throws Exception {
- super.channelInactive(ctx);
- InetSocketAddress ipSocket = (InetSocketAddress) ctx.channel().remoteAddress();
- int port = ipSocket.getPort();
- String host = ipSocket.getHostString();
- log.error("与设备" + host + ":" + port + "连接断开!");
- // 重连
- ArdEquipRadar ardEquipRadar = BootNettyClientChannelCache.getRadar(host + ":" + port);
- if (ardEquipRadar != null) {
- BootNettyClientThread thread = new BootNettyClientThread(ardEquipRadar);
- thread.start();
- }
- }
-
- /**
- * 解析报警数据
- */
- public void processData(ArdEquipRadar ardEquipRadarbyte, byte[] data) {
- try {
- String radarId = ardEquipRadarbyte.getId();
- String radarName = ardEquipRadarbyte.getName();
- Double radarLongitude = ardEquipRadarbyte.getLongitude();
- Double radarLagitude = ardEquipRadarbyte.getLatitude();
- Double radarAltitude = ardEquipRadarbyte.getAltitude();
- //region crc校验-目前仅用于显示校验结果
- Boolean crc32Check = MessageParsing.CRC32Check(data);
- if (!crc32Check) {
- log.debug("CRC32校验不通过");
- } else {
- //log.debug("CRC32校验通过");
- }
- //endregion
- //log.info("原始数据:" + DatatypeConverter.printHexBinary(data));
- //log.info("雷达信息:" + host + "【port】" + port + "【X】" + longitude + "【Y】" + lagitude + "【Z】" + altitude);
- data = MessageParsing.transferData(data);//去掉包头和包尾、校验及转义
- //region 负载头解析
- byte[] type = Arrays.copyOfRange(data, 0, 1);//命令类型
- // log.info("命令类型:" + DatatypeConverter.printHexBinary(type));
- byte[] cmdId = Arrays.copyOfRange(data, 1, 2);//命令ID
- String cmdIdStr = DatatypeConverter.printHexBinary(cmdId);
- //log.info("命令ID:" + DatatypeConverter.printHexBinary(cmdId));
- byte[] payloadSize = Arrays.copyOfRange(data, 2, 4);//有效负载大小
- payloadSize = toLittleEndian(payloadSize);
- //log.info("payloadSize:" + DatatypeConverter.printHexBinary(payloadSize));
- int payloadSizeToDecimal = ByteUtils.bytesToDecimal(payloadSize);
- // log.info("有效负载大小(转整型):" + payloadSizeToDecimal);
- //endregion
- List<ArdAlarmRadar> radarAlarmInfos = new ArrayList<>();
- ArdAlarmRadar radarFollowInfo = null;
- //抽油机状态雷达推送集合
- List<ArdAlarmRadar> well = new ArrayList<>();
- String alarmTime = "";
- Integer targetNum = 0;
- log.debug("Processing radar data 【" + radarName + "】数据-->命令ID:" + cmdIdStr + "二进制:" + byteToBitString(cmdId[0]));
- //雷达移动防火报警
- if (Arrays.equals(cmdId, new byte[]{0x01})) {
- //region 告警信息反馈
- byte[] dwTim = Arrays.copyOfRange(data, 4, 8);
- dwTim = toLittleEndian(dwTim);
- SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
- long l = ByteUtils.bytesToDecimal(dwTim);
- alarmTime = sdf.format(l * 1000);
- // log.info("周视图像的出现时间(转date):" + alarmTime);
-
- byte[] wTargetNum = Arrays.copyOfRange(data, 8, 10);
- wTargetNum = toLittleEndian(wTargetNum);
- targetNum = ByteUtils.bytesToDecimal(wTargetNum);
- if (targetNum == 0) {
- return;
- }
- //log.debug("目标总点数(转整型):" + targetNum);
-
- //解析NET_TARGET_UNIT(64是NET_TARGET_HEAD的字节数)
- int uintSize = (payloadSizeToDecimal - 64) / targetNum;
- // log.info("单条报警大小:" + uintSize);
-
- for (int i = 0; i < targetNum; i++) {
-
- Integer index = 68 + uintSize * i;
- byte[] dwID = Arrays.copyOfRange(data, index, index + 4);
- // log.info("目标ID:" + DatatypeConverter.printHexBinary(cmdId));
- dwID = toLittleEndian(dwID);
- int targetId = ByteUtils.bytesToDecimal(dwID);
- // log.info("目标ID号:" + targetId);
-
- byte[] iDistance = Arrays.copyOfRange(data, index + 8, index + 12);
- iDistance = toLittleEndian(iDistance);
- double Distance = ByteUtils.bytesToDecimal(iDistance);
- //log.debug("目标当前直线距离(m):" + Distance);
-
- //region 不需要的字段
-// byte[] dwGSum = Arrays.copyOfRange(data, index + 4, index + 8);
-// dwGSum = toLittleEndian(dwGSum);
-// int GSum = byteArrayToDecimal(dwGSum);
-// log.info("目标当前像素灰度和:" + GSum);
-// byte[] iTw = Arrays.copyOfRange(data, index + 12, index + 16);
-// iTw = toLittleEndian(iTw);
-// int Tw = byteArrayToDecimal(iTw);
-// log.info("目标当前的像素宽度:" + Tw);
-//
-// byte[] iTh = Arrays.copyOfRange(data, index + 16, index + 20);
-// iTh = toLittleEndian(iTh);
-// int Th = byteArrayToDecimal(iTh);
-// log.info("目标当前的像素高度:" + Th);
-//
-// byte[] wPxlArea = Arrays.copyOfRange(data, index + 20, index + 22);
-// wPxlArea = toLittleEndian(wPxlArea);
-// int PxlArea = byteArrayToDecimal(wPxlArea);
-// log.info("目标当前像素面积:" + PxlArea);
-//
-// byte[] cTrkNum = Arrays.copyOfRange(data, index + 22, index + 23);
-// cTrkNum = toLittleEndian(cTrkNum);
-// int TrkNum = byteArrayToDecimal(cTrkNum);
-// log.info("轨迹点数:" + TrkNum);
-
-// byte[] sVx = Arrays.copyOfRange(data, index + 24, index + 26);
-// sVx = toLittleEndian(sVx);
-// int Vx = byteArrayToDecimal(sVx);
-// log.info("目标当前速度矢量(像素距离)X:" + Vx);
-//
-// byte[] sVy = Arrays.copyOfRange(data, index + 26, index + 28);
-// sVy = toLittleEndian(sVy);
-// int Vy = byteArrayToDecimal(sVy);
-// log.info("目标当前速度矢量(像素距离)Y:" + Vy);
-//
-// byte[] sAreaNo = Arrays.copyOfRange(data, index + 28, index + 30);
-// sAreaNo = toLittleEndian(sAreaNo);
-// int AreaNo = byteArrayToDecimal(sAreaNo);
-// log.info("目标归属的告警区域号:" + AreaNo);
-//
-// byte[] cGrp = Arrays.copyOfRange(data, index + 30, index + 31);
-// cGrp = toLittleEndian(cGrp);
-// int Grp = byteArrayToDecimal(cGrp);
-// log.info("所属组:" + Grp);
- //endregion
- String alarmType = "";
- byte[] cStat = Arrays.copyOfRange(data, index + 23, index + 24);
- log.info("原始状态:" + byteToBitString(cStat[0]));
- // cStat = toLittleEndian(cStat);
- // 提取第4位至第6位的值
- int extractedValue = (cStat[0] >> 4) & 0b00001111;
- // 判断提取的值
- if (extractedValue == 0b0000) {
- alarmType = "运动目标检测";
- } else if (extractedValue == 0b0001) {
- alarmType = "热源检测";
- }
- // log.info("报警类型:" + alarmType);
- byte[] szName = Arrays.copyOfRange(data, index + 64, index + 96);
- String alarmPointName = ByteUtils.bytesToStringZh(szName);
- // log.info("所属告警区域名称:" + alarmPointName);
- byte[] afTx = Arrays.copyOfRange(data, index + 96, index + 100);
- afTx = toLittleEndian(afTx);
- float fTx = ByteUtils.bytesToFloat(afTx);
- // log.info("水平角度:" + fTx);
- byte[] afTy = Arrays.copyOfRange(data, index + 112, index + 116);
- afTy = toLittleEndian(afTy);
- float fTy = ByteUtils.bytesToFloat(afTy);
- //log.debug("垂直角度:" + fTy);
- // 将角度转换为弧度
- double thetaRadians = Math.toRadians(fTy + 90);
- // 使用正弦函数计算对边长度
- Distance = Math.sin(thetaRadians) * Distance;
- //log.debug("目标投影距离(m):" + Distance);
-
- Double[] radarXY = {radarLongitude, radarLagitude};
- Double[] alarmXY = GisUtils.CalculateCoordinates(radarXY, Distance, (double) fTx);
- log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】" + targetId + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime + "【name】" + alarmPointName);
- ArdAlarmRadar ardAlarmRadar = new ArdAlarmRadar();
- ardAlarmRadar.setTargetId(targetId);
- ardAlarmRadar.setName(alarmPointName);
- ardAlarmRadar.setLongitude(alarmXY[0]);
- ardAlarmRadar.setLatitude(alarmXY[1]);
- ardAlarmRadar.setAlarmType(alarmType);
- radarAlarmInfos.add(ardAlarmRadar);
-
- int bit1 = (cStat[0] >> 1) & 0x1;
- //目标的B1=1 锁定
- if (bit1 == 1) {
- radarFollowInfo = ardAlarmRadar;
- //将追踪锁定的报警对象属性复制给radarFollowInfo对象
- //BeanUtils.copyProperties(ardAlarmRadar, radarFollowInfo);
- }
- }
- //endregion
- if (StringUtils.isEmpty(alarmTime)) {
- return;
- }
- if (targetNum == 0) {
- return;
- }
- RadarAlarmData radarAlarmData = new RadarAlarmData();
- radarAlarmData.setRadarId(radarId);
- radarAlarmData.setRadarName(radarName);
- radarAlarmData.setAlarmTime(alarmTime);
- radarAlarmData.setArdAlarmRadars(radarAlarmInfos);
- MqttProducer.publish(2, false, "radar", JSON.toJSONString(radarAlarmData));
- if (radarFollowInfo != null) {
- //当前雷达扫描存在引导跟踪数据,只保留最后一次锁定的数据
- MqttProducer.publish(2, false, "radarFollowGuide", JSON.toJSONString(radarFollowInfo));
- }
- //抽油机状态MQTT队列
- radarAlarmData.setArdAlarmRadars(well);
- MqttProducer.publish(2, false, "radarWellData", JSON.toJSONString(radarAlarmData));
-
- }
- //抽油机AI状态反馈
- if (Arrays.equals(cmdId, new byte[]{0x04})) {
- //region抽油机AI状态反馈
- String hexString = DatatypeConverter.printHexBinary(data);
- //log.info(hexString);
-
- byte[] dwTim = Arrays.copyOfRange(data, 4, 8);
- dwTim = toLittleEndian(dwTim);
- SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
- long l = ByteUtils.bytesToDecimal(dwTim);
- alarmTime = sdf.format(l * 1000);
- //log.info("周视图像的出现时间(转date):" + alarmTime);
-
- byte[] wTargetNum = Arrays.copyOfRange(data, 8, 10);
- wTargetNum = toLittleEndian(wTargetNum);
- targetNum = ByteUtils.bytesToDecimal(wTargetNum);
- //log.debug("目标总点数(转整型):" + targetNum);
- if (targetNum == 0) {
- return;
- }
- //解析NET_TARGET_UNIT(64是NET_TARGET_HEAD的字节数)
- int uintSize = (payloadSizeToDecimal - 64) / targetNum;
- //log.info("单条报警大小:" + uintSize);
- for (int i = 0; i < targetNum; i++) {
- Integer index = 68 + uintSize * i;
- byte[] dwID = Arrays.copyOfRange(data, index, index + 4);
- //log.info("目标ID:" + DatatypeConverter.printHexBinary(dwID));
- dwID = toLittleEndian(dwID);
- int targetId = ByteUtils.bytesToDecimal(dwID);
- //log.info("目标ID号:" + targetId);
- //region 不需要的字段
- byte[] iTw = Arrays.copyOfRange(data, index + 4, index + 8);
- iTw = toLittleEndian(iTw);
- int Tw = ByteUtils.bytesToDecimal(iTw);
- // log.info("目标当前的像素宽度:" + Tw);
-
- byte[] iTh = Arrays.copyOfRange(data, index + 8, index + 12);
- iTh = toLittleEndian(iTh);
- int Th = ByteUtils.bytesToDecimal(iTh);
- //log.info("目标当前的像素高度:" + Th);
-
- byte[] fTx = Arrays.copyOfRange(data, index + 12, index + 16);
- fTx = toLittleEndian(fTx);
- float fTxAngle = ByteUtils.bytesToFloat(fTx);
- //log.debug("p角度:" + fTxAngle);
- byte[] fTy = Arrays.copyOfRange(data, index + 16, index + 20);
- fTy = toLittleEndian(fTy);
- float fTyAngle = ByteUtils.bytesToFloat(fTy);
- //log.debug("t角度:" + fTyAngle);
-
- byte[] sAreaNo = Arrays.copyOfRange(data, index + 20, index + 22);
- sAreaNo = toLittleEndian(sAreaNo);
- int AreaNo = ByteUtils.bytesToDecimal(sAreaNo);
- //log.debug("目标归属的告警区域号:" + AreaNo);
-
- byte[] cGrp = Arrays.copyOfRange(data, index + 22, index + 23);
- cGrp = toLittleEndian(cGrp);
- int Grp = ByteUtils.bytesToDecimal(cGrp);
- //log.info("所属组:" + Grp);
- //endregion
- String alarmType;
- //抽油机状态变量
- String wellType;
- byte[] cStat = Arrays.copyOfRange(data, index + 23, index + 24);
- cStat = toLittleEndian(cStat);
- //String binaryString = String.format("%8s", Integer.toBinaryString(cStat[0] & 0xFF)).replace(' ', '0');
- //log.info("目标当前状态:" + binaryString);
- // 提取第0位值
- // 使用位运算操作判断第0位是否为1
- boolean isB0 = (cStat[0] & 0x01) == 0x00;
- // 判断提取的值
- if (isB0) {
- alarmType = "雷达抽油机停机";
- byte[] szName = Arrays.copyOfRange(data, index + 32, index + 64);
- //log.info("所属告警区域名称:" + DatatypeConverter.printHexBinary(szName));
- String alarmPointName = ByteUtils.bytesToStringZh(szName);
- // log.info("所属告警区域名称:" + alarmPointName);
- //log.debug("报警信息:"+ "【radarName】" + radarName + "【targetId】" + targetId + "【name】" + alarmPointName + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime);
- ArdAlarmRadar ardAlarmRadar = new ArdAlarmRadar();
- ardAlarmRadar.setTargetId(targetId);
- ardAlarmRadar.setName(alarmPointName);
- ardAlarmRadar.setAlarmType(alarmType);
- radarAlarmInfos.add(ardAlarmRadar);
- wellType = "停机";
- } else {
- wellType = "运行";
- }
- //抽油机状态集合中装入数据
- byte[] szName = Arrays.copyOfRange(data, index + 32, index + 64);
- String alarmPointName = ByteUtils.bytesToStringZh(szName);
- log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】" + targetId + "【alarmType】抽油机状态报警" + "【alarmTime】" + alarmTime + "【name】" + alarmPointName + "【alarmState】" + wellType);
- ArdAlarmRadar wellAlarm = new ArdAlarmRadar();
- wellAlarm.setTargetId(targetId);
- wellAlarm.setName(alarmPointName);
- wellAlarm.setAlarmType(wellType);
- well.add(wellAlarm);
- }
- //endregion
- if (StringUtils.isEmpty(alarmTime)) {
- return;
- }
- if (targetNum == 0) {
- return;
- }
- RadarAlarmData radarAlarmData = new RadarAlarmData();
- radarAlarmData.setRadarId(radarId);
- radarAlarmData.setRadarName(radarName);
- radarAlarmData.setAlarmTime(alarmTime);
- radarAlarmData.setArdAlarmRadars(radarAlarmInfos);
- MqttProducer.publish(2, false, "radar", JSON.toJSONString(radarAlarmData));
- //抽油机状态MQTT队列
- radarAlarmData.setArdAlarmRadars(well);
- MqttProducer.publish(2, false, "radarWellData", JSON.toJSONString(radarAlarmData));
- }
- //强制引导
- if (Arrays.equals(cmdId, new byte[]{0x02})) {
- //region 告警前端发送的强制引导信息
- byte[] iDistance = Arrays.copyOfRange(data, 4, 8);
- iDistance = toLittleEndian(iDistance);
- long distance = ByteUtils.bytesToDecimal(iDistance);
- log.info("目标当前距离(m):" + distance);
- byte[] fTx = Arrays.copyOfRange(data, 8, 12);
- fTx = toLittleEndian(fTx);
- float tx = ByteUtils.bytesToFloat(fTx);
- log.debug("方位:" + tx);
- byte[] fTy = Arrays.copyOfRange(data, 12, 16);
- fTy = toLittleEndian(fTy);
- float ty = ByteUtils.bytesToFloat(fTy);
- if (ty < 0) {
- ty += 360;
- }
- log.debug("俯仰:" + ty);
- Map<String, Object> forceGuideMap = new HashMap<>();
- forceGuideMap.put("distance", distance);
- forceGuideMap.put("p", tx);
- forceGuideMap.put("t", ty);
- forceGuideMap.put("radarId", radarId);
- log.debug("强制引导信息" + forceGuideMap);
- //endregion
- MqttProducer.publish(2, false, "radarForceGuide", JSON.toJSONString(forceGuideMap));
- }
- } catch (Exception ex) {
- log.error("雷达报文解析异常:" + ex.getMessage());
- }
- }
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/BootNettyChannelInitializer.java b/src/main/java/com/ard/utils/netty/tcp/BootNettyChannelInitializer.java
deleted file mode 100644
index d87b9ff..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/BootNettyChannelInitializer.java
+++ /dev/null
@@ -1,20 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import io.netty.channel.ChannelInitializer;
-import io.netty.channel.socket.SocketChannel;
-import io.netty.handler.codec.string.StringDecoder;
-import io.netty.handler.codec.string.StringEncoder;
-import io.netty.util.CharsetUtil;
-
-public class BootNettyChannelInitializer extends ChannelInitializer<SocketChannel> {
-
- @Override
- protected void initChannel(SocketChannel ch){
- //ch.pipeline().addLast("encoder", new StringEncoder(CharsetUtil.UTF_8));
- //ch.pipeline().addLast("decoder", new StringDecoder(CharsetUtil.UTF_8));
- /**
- * 自定义ChannelInboundHandlerAdapter
- */
- ch.pipeline().addLast(new BootNettyChannelInboundHandlerAdapter());
- }
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/BootNettyClient.java b/src/main/java/com/ard/utils/netty/tcp/BootNettyClient.java
deleted file mode 100644
index 6a44e88..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/BootNettyClient.java
+++ /dev/null
@@ -1,110 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.ard.alarm.radar.domain.ArdEquipRadar;
-import com.ard.alarm.radar.service.IArdEquipRadarService;
-import com.ard.utils.netty.config.NettyTcpConfiguration;
-import io.netty.bootstrap.Bootstrap;
-import io.netty.channel.*;
-import io.netty.channel.nio.NioEventLoopGroup;
-import io.netty.channel.socket.nio.NioSocketChannel;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.boot.ApplicationArguments;
-import org.springframework.boot.ApplicationRunner;
-import org.springframework.stereotype.Component;
-
-import javax.annotation.Resource;
-import java.nio.channels.SocketChannel;
-import java.util.List;
-
-@Slf4j(topic = "netty")
-@Component
-public class BootNettyClient {
- @Resource
- IArdEquipRadarService ardEquipRadarService;
- @Resource
- NettyTcpConfiguration nettyTcpConfig;
-
-
- static Integer waitTimes = 1;
- static EventLoopGroup eventLoopGroup = new NioEventLoopGroup();
- /**
- * 初始化Bootstrap
- */
- public static final Bootstrap getBootstrap(EventLoopGroup group) {
- if (null == group) {
- group = eventLoopGroup;
- }
- Bootstrap bootstrap = new Bootstrap();
- bootstrap.group(group).channel(NioSocketChannel.class)
- .option(ChannelOption.TCP_NODELAY, true)
- .option(ChannelOption.SO_KEEPALIVE, true)
- .handler(new BootNettyChannelInitializer());
- return bootstrap;
- }
-
- public void connect(ArdEquipRadar radar) throws Exception {
- String host = radar.getIp();
- int port=radar.getPort();
- log.debug("正在进行连接:【" + host+":"+port+"】");
- eventLoopGroup.shutdownGracefully();
- eventLoopGroup = new NioEventLoopGroup();
- Bootstrap bootstrap = getBootstrap(null);
-
- try {
- bootstrap.remoteAddress(host, port);
- // 异步连接tcp服务端
- ChannelFuture future = bootstrap.connect().addListener((ChannelFuture futureListener) -> {
- final EventLoop eventLoop = futureListener.channel().eventLoop();
- if (futureListener.isSuccess()) {
- BootNettyClientChannel bootNettyClientChannel = new BootNettyClientChannel();
- Channel channel = futureListener.channel();
- String id = futureListener.channel().id().toString();
-// String id = host;
- bootNettyClientChannel.setChannel(channel);
- bootNettyClientChannel.setCode("clientId:" + id);
- BootNettyClientChannelCache.save("clientId:" + id, bootNettyClientChannel);
- BootNettyClientChannelCache.save(host+":"+port,radar);
- log.debug("netty client start success=" + id);
- } else {
-// System.err.println("连接失败," + waitTimes.toString() + "秒后重新连接:" + host);
- try {
- Thread.sleep(waitTimes * 1000);
- } finally {
- connect(radar);
- }
- }
- });
- future.channel().closeFuture().sync();
- } catch (Exception e) {
- System.err.println("连接异常," + waitTimes.toString() + "秒后重新连接:" + host);
- try {
- Thread.sleep(waitTimes * 1000);
- } finally {
- connect(radar);
- }
- e.printStackTrace();
- } finally {
- /**
- * 退出,释放资源
- */
- eventLoopGroup.shutdownGracefully().sync();
- }
-
- }
- /**
- * 初始化方法
- */
- public void run(ApplicationArguments args) {
- if (!nettyTcpConfig.getEnabled()) {
- return;
- }
- List<ArdEquipRadar> ardEquipRadars = ardEquipRadarService.selectArdEquipRadarList(new ArdEquipRadar());
- for (ArdEquipRadar ardEquipRadar : ardEquipRadars) {
- String host = ardEquipRadar.getIp();
- Integer port = Integer.valueOf(ardEquipRadar.getPort());
- log.debug("TCP client try to connect radar【:" + host + ":" + port+"】");
- BootNettyClientThread thread = new BootNettyClientThread(ardEquipRadar);
- thread.start();
- }
- }
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/BootNettyClientChannel.java b/src/main/java/com/ard/utils/netty/tcp/BootNettyClientChannel.java
deleted file mode 100644
index f8c9dba..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/BootNettyClientChannel.java
+++ /dev/null
@@ -1,16 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import io.netty.channel.Channel;
-import lombok.Data;
-
-@Data
-public class BootNettyClientChannel {
-
- // 连接客户端唯一的code
- private String code;
-
- // 客户端最新发送的消息内容
- private String last_data;
-
- private transient volatile Channel channel;
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/BootNettyClientChannelCache.java b/src/main/java/com/ard/utils/netty/tcp/BootNettyClientChannelCache.java
deleted file mode 100644
index 4d31833..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/BootNettyClientChannelCache.java
+++ /dev/null
@@ -1,41 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.ard.alarm.radar.domain.ArdEquipRadar;
-
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
-public class BootNettyClientChannelCache {
- public static volatile Map<String, ArdEquipRadar> radarMapCache = new ConcurrentHashMap<String, ArdEquipRadar>();
- public static volatile Map<String, BootNettyClientChannel> channelMapCache = new ConcurrentHashMap<String, BootNettyClientChannel>();
-
- public static void add(String code, BootNettyClientChannel channel){
- channelMapCache.put(code,channel);
- }
- public static void addRadar(String code, ArdEquipRadar radar){
- radarMapCache.put(code,radar);
- }
- public static BootNettyClientChannel get(String code){
- return channelMapCache.get(code);
- }
- public static ArdEquipRadar getRadar(String code){
- return radarMapCache.get(code);
- }
- public static void remove(String code){
- channelMapCache.remove(code);
- }
- public static void removeRadar(String code){
- radarMapCache.remove(code);
- }
- public static void save(String code, BootNettyClientChannel channel) {
- if(channelMapCache.get(code) == null) {
- add(code,channel);
- }
- }
- public static void save(String code, ArdEquipRadar radar) {
- if(radarMapCache.get(code) == null) {
- addRadar(code,radar);
- }
- }
-
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/BootNettyClientThread.java b/src/main/java/com/ard/utils/netty/tcp/BootNettyClientThread.java
deleted file mode 100644
index 8aa85ff..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/BootNettyClientThread.java
+++ /dev/null
@@ -1,20 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.ard.alarm.radar.domain.ArdEquipRadar;
-
-public class BootNettyClientThread extends Thread {
-
- private final ArdEquipRadar ardEquipRadar;
-
- public BootNettyClientThread(ArdEquipRadar ardEquipRadar){
- this.ardEquipRadar = ardEquipRadar;
- }
-
- public void run() {
- try {
- new BootNettyClient().connect(ardEquipRadar);
- } catch (Exception e) {
- throw new RuntimeException(e);
- }
- }
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/ClientHandler.java b/src/main/java/com/ard/utils/netty/tcp/ClientHandler.java
index c59dd32..6e12a61 100644
--- a/src/main/java/com/ard/utils/netty/tcp/ClientHandler.java
+++ b/src/main/java/com/ard/utils/netty/tcp/ClientHandler.java
@@ -8,7 +8,6 @@
import com.ard.utils.util.GisUtils;
import com.ard.utils.mqtt.MqttProducer;
import io.netty.buffer.ByteBuf;
-import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelId;
import io.netty.channel.SimpleChannelInboundHandler;
@@ -17,7 +16,6 @@
import javax.xml.bind.DatatypeConverter;
import java.net.InetSocketAddress;
-import java.net.SocketAddress;
import java.text.SimpleDateFormat;
import java.util.*;
import java.util.concurrent.ScheduledFuture;
@@ -58,19 +56,19 @@
*/
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
- ChannelId id = ctx.channel().id();
InetSocketAddress ipSocket = (InetSocketAddress) ctx.channel().remoteAddress();
- int port = ipSocket.getPort();
- String host = ipSocket.getHostString();
- log.error("与设备" + host + ":" + port + "连接断开!");
- ArdEquipRadar ardEquipRadar = ClientInitialize.tureConnectMap.get(host+ ":" + port);
+ String ipPort = ipSocket.getHostString() + ":" + ipSocket.getPort();
+ log.error("与设备" + ipPort + "连接断开!");
// 连接断开后的最后处理
ctx.pipeline().remove(this);
ctx.deregister();
ctx.close();
-
// 将失败信息插入Set集合
- ClientInitialize.falseConnectSet.add(ardEquipRadar);
+ ArdEquipRadar radar = ClientInitialize.trueConnectMap.get(ipPort);
+ if (radar != null) {
+ ClientInitialize.falseConnectSet.add(radar);
+ ClientInitialize.trueConnectMap.remove(ipPort);
+ }
super.channelInactive(ctx);
}
@@ -83,20 +81,26 @@
* @throws Exception
*/
@Override
- public void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) throws Exception {
+ public void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) {
InetSocketAddress ipSocket = (InetSocketAddress) ctx.channel().remoteAddress();
- int port = ipSocket.getPort();
- String host = ipSocket.getHostString();
- ArdEquipRadar ardEquipRadar = ClientInitialize.tureConnectMap.get(host+":"+port);
- MessageParsing messageParsing = ClientInitialize.MessageMap.get(host + ":" + port);
+ String ipPort = ipSocket.getHostString() + ":" + ipSocket.getPort();
+ ArdEquipRadar radar = ClientInitialize.trueConnectMap.get(ipPort);
+ if (radar == null) {
+ return;
+ }
+ MessageHandler messageHandler = ClientInitialize.SucMessageHandlerMap.get(ipPort);
+ if (messageHandler == null) {
+ return;
+ }
// 处理接收到的消息
byte[] byteArray = new byte[msg.readableBytes()];
msg.getBytes(msg.readerIndex(), byteArray);
- byte[] bytes = messageParsing.receiveCompletePacket(byteArray);
+ byte[] bytes = messageHandler.receiveCompletePacket(byteArray);
if (bytes != null) {
- processData(ardEquipRadar, bytes);
+ processData(radar, bytes);
}
}
+
/**
* 通道数据处理完成
@@ -148,7 +152,7 @@
byte[] payloadCrc32 = ByteUtils.parseCrc32(payload);
byte[] footer = {0x01, 0x02, 0x00};
byte[] heart = ByteUtils.appendArrays(header, payload, payloadCrc32, footer);
-// byte[] heart = {0x01, 0x02, 0x01, 0x10, 0x00, 0x00, 0x00, (byte) 0x83, (byte) 0x88, 0x5d, 0x71, 0x01, 0x02, 0x00};
+ //byte[] heart = {0x01, 0x02, 0x01, 0x10, 0x00, 0x00, 0x00, (byte) 0x83, (byte) 0x88, 0x5d, 0x71, 0x01, 0x02, 0x00};
String hexString = DatatypeConverter.printHexBinary(heart);
// log.debug("发送心跳:" + hexString);
message.writeBytes(heart);
@@ -170,15 +174,15 @@
/**
* 解析报警数据
*/
- public void processData(ArdEquipRadar ardEquipRadarbyte, byte[] data) {
+ public void processData(ArdEquipRadar radar, byte[] data) {
try {
- String radarId = ardEquipRadarbyte.getId();
- String radarName = ardEquipRadarbyte.getName();
- Double radarLongitude = ardEquipRadarbyte.getLongitude();
- Double radarLagitude = ardEquipRadarbyte.getLatitude();
- Double radarAltitude = ardEquipRadarbyte.getAltitude();
+ String radarId = radar.getId();
+ String radarName = radar.getName();
+ Double radarLongitude = radar.getLongitude();
+ Double radarLagitude = radar.getLatitude();
+ Double radarAltitude = radar.getAltitude();
//region crc校验-目前仅用于显示校验结果
- Boolean crc32Check = MessageParsing.CRC32Check(data);
+ Boolean crc32Check = MessageHandler.CRC32Check(data);
if (!crc32Check) {
log.debug("CRC32校验不通过");
} else {
@@ -187,12 +191,12 @@
//endregion
//log.info("原始数据:" + DatatypeConverter.printHexBinary(data));
//log.info("雷达信息:" + host + "【port】" + port + "【X】" + longitude + "【Y】" + lagitude + "【Z】" + altitude);
- data = MessageParsing.transferData(data);//去掉包头和包尾、校验及转义
+ data = MessageHandler.transferData(data);//去掉包头和包尾、校验及转义
//region 负载头解析
byte[] type = Arrays.copyOfRange(data, 0, 1);//命令类型
// log.info("命令类型:" + DatatypeConverter.printHexBinary(type));
byte[] cmdId = Arrays.copyOfRange(data, 1, 2);//命令ID
- String cmdIdStr=DatatypeConverter.printHexBinary(cmdId);
+ String cmdIdStr = DatatypeConverter.printHexBinary(cmdId);
//log.info("命令ID:" + DatatypeConverter.printHexBinary(cmdId));
byte[] payloadSize = Arrays.copyOfRange(data, 2, 4);//有效负载大小
payloadSize = toLittleEndian(payloadSize);
@@ -206,7 +210,7 @@
List<ArdAlarmRadar> well = new ArrayList<>();
String alarmTime = "";
Integer targetNum = 0;
- log.debug("处理雷达"+radarName+"数据-->命令ID:"+cmdIdStr);
+ log.debug("处理雷达" + radarName + "数据-->命令ID:" + cmdIdStr);
//雷达移动防火报警
if (Arrays.equals(cmdId, new byte[]{0x01})) {
//region 告警信息反馈
@@ -315,15 +319,14 @@
double thetaRadians = Math.toRadians(fTy + 90);
// 使用正弦函数计算对边长度
Distance = Math.sin(thetaRadians) * Distance;
- if(Distance<0)
- {
+ if (Distance < 0) {
continue;//过滤距离小于0的脏数据
}
//log.debug("目标投影距离(m):" + Distance);
Double[] radarXY = {radarLongitude, radarLagitude};
Double[] alarmXY = GisUtils.CalculateCoordinates(radarXY, Distance, (double) fTx);
- log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】"+ targetId + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime + "【name】" + alarmPointName+"【Distance】"+Distance);
+ log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】" + targetId + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime + "【name】" + alarmPointName + "【Distance】" + Distance);
ArdAlarmRadar ardAlarmRadar = new ArdAlarmRadar();
ardAlarmRadar.setTargetId(targetId);
ardAlarmRadar.setName(alarmPointName);
@@ -350,7 +353,7 @@
radarAlarmData.setAlarmTime(alarmTime);
radarAlarmData.setArdAlarmRadars(radarAlarmInfos);
MqttProducer.publish(2, false, "radar", JSON.toJSONString(radarAlarmData));
- if (radarFollowInfos.size() >0) {
+ if (radarFollowInfos.size() > 0) {
radarAlarmData.setArdFollowRadars(radarFollowInfos);
//当前雷达扫描存在引导跟踪数据,只保留最后一次锁定的数据
MqttProducer.publish(2, false, "radarFollowGuide", JSON.toJSONString(radarAlarmData));
@@ -437,7 +440,7 @@
//log.info("所属告警区域名称:" + DatatypeConverter.printHexBinary(szName));
String alarmPointName = ByteUtils.bytesToStringZh(szName);
// log.info("所属告警区域名称:" + alarmPointName);
- log.debug("报警信息:"+ "【radarName】" + radarName + "【targetId】" + targetId + "【name】" + alarmPointName + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime);
+ log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】" + targetId + "【name】" + alarmPointName + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime);
ArdAlarmRadar ardAlarmRadar = new ArdAlarmRadar();
ardAlarmRadar.setTargetId(targetId);
ardAlarmRadar.setName(alarmPointName);
diff --git a/src/main/java/com/ard/utils/netty/tcp/ClientHelper.java b/src/main/java/com/ard/utils/netty/tcp/ClientHelper.java
deleted file mode 100644
index cbad3c0..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/ClientHelper.java
+++ /dev/null
@@ -1,41 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.ard.utils.util.ByteUtils;
-import io.netty.buffer.ByteBuf;
-import io.netty.buffer.Unpooled;
-import org.springframework.scheduling.annotation.EnableScheduling;
-import org.springframework.scheduling.annotation.Scheduled;
-import org.springframework.stereotype.Component;
-
-import javax.xml.bind.DatatypeConverter;
-import java.util.Map;
-
-@Component
-@EnableScheduling
-public class ClientHelper {
-
- // 使用定时器发送心跳
- @Scheduled(cron = "0/3 * * * * ?")
- public void heart_timer() {
- //System.err.println("BootNettyClientChannelCache.channelMapCache.size():" + BootNettyClientChannelCache.channelMapCache.size());
- if (BootNettyClientChannelCache.channelMapCache.size() > 0) {
- for (Map.Entry<String, BootNettyClientChannel> entry : BootNettyClientChannelCache.channelMapCache.entrySet()) {
- BootNettyClientChannel bootNettyChannel = entry.getValue();
- //System.out.println(bootNettyChannel.getCode());
- try {
- byte[] header = {0x01, 0x02, 0x01};
- byte[] payload = {0x10, 0x00, 0x00, 0x00};
- byte[] payloadCrc32 = ByteUtils.parseCrc32(payload);
- byte[] footer = {0x01, 0x02, 0x00};
- byte[] heart = ByteUtils.appendArrays(header, payload, payloadCrc32, footer);
- // log.debug("发送心跳:" + hexString);
- //message.writeBytes(heart);
- bootNettyChannel.getChannel().writeAndFlush(Unpooled.buffer().writeBytes(heart));
- } catch (Exception e) {
- continue;
- }
- }
- }
-
- }
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/ClientInitialize.java b/src/main/java/com/ard/utils/netty/tcp/ClientInitialize.java
index ddd4f9c..b30b2b5 100644
--- a/src/main/java/com/ard/utils/netty/tcp/ClientInitialize.java
+++ b/src/main/java/com/ard/utils/netty/tcp/ClientInitialize.java
@@ -1,13 +1,5 @@
package com.ard.utils.netty.tcp;
-/**
- * @Description:
- * @ClassName: init
- * @Author: 刘苏义
- * @Date: 2023年07月05日13:11
- * @Version: 1.0
- **/
-
import com.ard.alarm.radar.domain.ArdEquipRadar;
import com.ard.alarm.radar.service.IArdEquipRadarService;
import com.ard.utils.netty.config.NettyTcpConfiguration;
@@ -37,17 +29,18 @@
@Component
@Slf4j(topic = "netty")
@Order(2)
-public class ClientInitialize implements ApplicationRunner{
+public class ClientInitialize implements ApplicationRunner {
@Resource
NettyTcpConfiguration nettyTcpConfig;
@Resource
IArdEquipRadarService ardEquipRadarService;
private Bootstrap bootstrap;
- public static CopyOnWriteArraySet<ArdEquipRadar> falseConnectSet = new CopyOnWriteArraySet();
- public static ConcurrentHashMap<String, ArdEquipRadar> tureConnectMap = new ConcurrentHashMap();
- public static ConcurrentHashMap<String, Object> SuccessConnectMap = new ConcurrentHashMap();
- public static ConcurrentHashMap<String, MessageParsing> MessageMap = new ConcurrentHashMap();
+ public static CopyOnWriteArraySet<ArdEquipRadar> falseConnectSet = new CopyOnWriteArraySet();//失败连接的雷达Set
+ public static ConcurrentHashMap<String, ArdEquipRadar> trueConnectMap = new ConcurrentHashMap();//成功连接的ip端口对应的雷达
+ public static ConcurrentHashMap<String, MessageHandler> SucMessageHandlerMap = new ConcurrentHashMap();//成功连接的ip端口对应的报文解析器
+ public static ConcurrentHashMap<String, Channel> SucChannelMap = new ConcurrentHashMap();//成功连接的ip端口对应的netty通道
+
/**
* Netty初始化配置
*/
@@ -65,7 +58,7 @@
}
});
- // 异步持续监听连接失败的地址
+ //异步持续监听连接失败的地址
CompletableFuture.runAsync(new Runnable() {
@Override
public void run() {
@@ -75,8 +68,8 @@
// 循环集合内元素
falseConnectSet.forEach(new Consumer<ArdEquipRadar>() {
@Override
- public void accept(ArdEquipRadar ardEquipRadar) {
- connectServer(ardEquipRadar);
+ public void accept(ArdEquipRadar radar) {
+ connectServer(radar);
}
});
}
@@ -99,24 +92,23 @@
// 获取地址及端口
String host = ardEquipRadar.getIp();
Integer port = ardEquipRadar.getPort();
+ String ipPort = host + ":" + port;
// 异步连接tcp服务端
bootstrap.remoteAddress(host, port).connect().addListener((ChannelFuture futureListener) -> {
- if (!futureListener.isSuccess()) {
- log.debug("雷达【" + host + ":" + port + "】连接失败");
- futureListener.channel().close();
- // 连接失败信息插入Set
- falseConnectSet.add(ardEquipRadar);
- // 连接失败信息从map移除
- tureConnectMap.remove( host + ":" + port);
- SuccessConnectMap.remove(ardEquipRadar.getId());
- } else {
- log.debug("雷达【" + host + ":" + port + "】连接成功");
+ if (futureListener.isSuccess()) {
+ log.debug("雷达【" + ipPort + "】连接成功");
// 连接成功信息从Set拔除
falseConnectSet.remove(ardEquipRadar);
// 连接成功信息写入map
- tureConnectMap.put(host+":"+port, ardEquipRadar);
- MessageMap.put(host+":"+port,new MessageParsing());
- SuccessConnectMap.put(ardEquipRadar.getId(),futureListener.channel());
+ trueConnectMap.put(ipPort, ardEquipRadar);
+ SucMessageHandlerMap.put(ipPort, new MessageHandler());
+ SucChannelMap.put(ipPort, futureListener.channel());
+ } else {
+ log.debug("雷达【" + ipPort + "】连接失败");
+ futureListener.channel().close();
+ // 连接失败信息插入Set
+ falseConnectSet.add(ardEquipRadar);
+
}
});
}
@@ -134,7 +126,7 @@
for (ArdEquipRadar ardEquipRadar : ardEquipRadars) {
String host = ardEquipRadar.getIp();
Integer port = Integer.valueOf(ardEquipRadar.getPort());
- log.debug("TCP client try to connect radar【" + host + ":" + port+"】");
+ log.debug("TCP client try to connect radar【" + host + ":" + port + "】");
connectServer(ardEquipRadar);//连接每一个雷达服务
}
}
diff --git a/src/main/java/com/ard/utils/netty/tcp/DynamicClient.java b/src/main/java/com/ard/utils/netty/tcp/DynamicClient.java
deleted file mode 100644
index 82dfcdc..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/DynamicClient.java
+++ /dev/null
@@ -1,151 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.ard.alarm.radar.domain.ArdEquipRadar;
-import com.ard.alarm.radar.service.IArdEquipRadarService;
-import com.ard.utils.netty.config.NettyTcpConfiguration;
-import io.netty.bootstrap.Bootstrap;
-import io.netty.channel.*;
-import io.netty.channel.nio.NioEventLoopGroup;
-import io.netty.channel.socket.nio.NioSocketChannel;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.boot.ApplicationArguments;
-import org.springframework.boot.ApplicationRunner;
-import org.springframework.stereotype.Component;
-
-import javax.annotation.Resource;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.TimeUnit;
-
-
-/**
- * @Description: 雷达动态tcp客户端(备用)
- * @ClassName: DynamicClient
- * @Author: 刘苏义
- * @Date: 2023年11月30日9:25:48
- **/
-@Slf4j(topic = "netty")
-@Component
-public class DynamicClient{
- @Resource
- IArdEquipRadarService ardEquipRadarService;
- @Resource
- NettyTcpConfiguration nettyTcpConfig;
-
- private static List<Channel> serverChannels = new ArrayList<>();
- public static ConcurrentHashMap<Channel, ArdEquipRadar> ConnectMap = new ConcurrentHashMap();
-
- public static void main(String[] args) throws InterruptedException {
- EventLoopGroup group = new NioEventLoopGroup();
- Bootstrap bootstrap = new Bootstrap();
- bootstrap.group(group)
- .channel(NioSocketChannel.class)
- .option(ChannelOption.TCP_NODELAY, true)
- .option(ChannelOption.SO_KEEPALIVE, true)
- .handler(new DynamicClientInitializer());
-
- DynamicClient dynamicClient = new DynamicClient();
- ArdEquipRadar radar1 = new ArdEquipRadar();
- radar1.setName("511");
- radar1.setIp("112.98.126.2");
- radar1.setPort(1200);
- dynamicClient.connect(bootstrap, radar1);
-
- Thread.sleep(2000); // 等待连接建立
-
- // 模拟动态增加一个服务端
- ArdEquipRadar radar2 = new ArdEquipRadar();
- radar2.setName("140");
- radar2.setIp("112.98.126.2");
- radar2.setPort(1201);
- dynamicClient.connect(bootstrap, radar2);
-
- }
-
- public void connectToServer(ArdEquipRadar radar) {
- EventLoopGroup group = new NioEventLoopGroup();
- try {
- Bootstrap bootstrap = new Bootstrap();
- bootstrap.group(group)
- .channel(NioSocketChannel.class)
- .option(ChannelOption.TCP_NODELAY, true)
- .option(ChannelOption.SO_KEEPALIVE, true)
- .handler(new DynamicClientInitializer());
- // 连接服务端
- ChannelFuture channelFuture = bootstrap.connect(radar.getIp(), radar.getPort()).sync();
- Channel serverChannel = channelFuture.channel();
- // 将连接的服务端 Channel 添加到管理列表
- serverChannels.add(serverChannel);
- ConnectMap.put(serverChannel, radar);
- log.debug("雷达【" + radar.getIp() + ":" + radar.getPort() + "】连接成功");
- } catch (Exception e) {
- e.printStackTrace();
- }
- }
-
- public Channel connect(Bootstrap bootstrap, ArdEquipRadar radar) {
- ChannelFuture future = bootstrap.connect(radar.getIp(), radar.getPort());
- future.addListener((ChannelFutureListener) futureListener -> {
- if (futureListener.isSuccess()) {
- log.info("Connected to radar device: " + radar.getName() + "【" + radar.getIp() + ":" + radar.getPort() + "】" + " successful");
- // 在连接建立后,你可以在这里添加业务逻辑或其他处理
- serverChannels.add(future.channel());
- ConnectMap.put(future.channel(), radar);
- } else {
- log.error("Connection to radar device " + radar.getName() + "【" + radar.getIp() + ":" + radar.getPort() + "】" + " failed. Retrying...");
- // 连接失败时,定时进行重连
- futureListener.channel().eventLoop().schedule(
- () -> connect(bootstrap, radar),
- 1L, TimeUnit.SECONDS
- );
- }
- });
-
- return future.channel();
- }
-
- // 在连接建立后可以通过调用这个方法向指定的服务端发送数据
- public void sendDataToServer(Channel serverChannel, Object data) {
- serverChannel.writeAndFlush(data);
- }
-
- // 关闭指定的服务端连接
- public void closeServerConnection(Channel serverChannel) {
- serverChannel.close();
- serverChannels.remove(serverChannel);
- }
-
- // 关闭所有服务端连接
- public void closeAllServerConnections() {
- for (Channel serverChannel : serverChannels) {
- serverChannel.close();
- }
- serverChannels.clear();
- }
-
- /**
- * 初始化方法
- */
-
- public void run(ApplicationArguments args) {
- if (!nettyTcpConfig.getEnabled()) {
- return;
- }
- EventLoopGroup group = new NioEventLoopGroup();
- Bootstrap bootstrap = new Bootstrap();
- bootstrap.group(group)
- .channel(NioSocketChannel.class)
- .option(ChannelOption.TCP_NODELAY, true)
- .option(ChannelOption.SO_KEEPALIVE, true)
- .handler(new DynamicClientInitializer());
- List<ArdEquipRadar> ardEquipRadars = ardEquipRadarService.selectArdEquipRadarList(new ArdEquipRadar());
- for (ArdEquipRadar ardEquipRadar : ardEquipRadars) {
- String host = ardEquipRadar.getIp();
- Integer port = Integer.valueOf(ardEquipRadar.getPort());
- log.debug("TCP client try to connect radar【:" + host + ":" + port + "】");
- // connectServer(ardEquipRadar);//连接每一个雷达服务
- connect(bootstrap, ardEquipRadar);
- }
- }
-}
\ No newline at end of file
diff --git a/src/main/java/com/ard/utils/netty/tcp/DynamicClientHandler.java b/src/main/java/com/ard/utils/netty/tcp/DynamicClientHandler.java
deleted file mode 100644
index 1809740..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/DynamicClientHandler.java
+++ /dev/null
@@ -1,436 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import com.alibaba.fastjson2.JSON;
-import com.ard.alarm.radar.domain.ArdAlarmRadar;
-import com.ard.alarm.radar.domain.ArdEquipRadar;
-import com.ard.alarm.radar.domain.RadarAlarmData;
-import com.ard.utils.mqtt.MqttProducer;
-import com.ard.utils.util.ByteUtils;
-import com.ard.utils.util.GisUtils;
-import io.netty.buffer.ByteBuf;
-import io.netty.channel.Channel;
-import io.netty.channel.ChannelHandlerContext;
-import io.netty.channel.SimpleChannelInboundHandler;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.lang3.StringUtils;
-import javax.xml.bind.DatatypeConverter;
-import java.text.SimpleDateFormat;
-import java.util.*;
-import java.util.concurrent.ScheduledFuture;
-import java.util.concurrent.TimeUnit;
-import static com.ard.utils.util.ByteUtils.byteToBitString;
-import static com.ard.utils.util.ByteUtils.toLittleEndian;
-
-/**
- * @Description: 客户端数据处理器(备用)
- * @ClassName: DynamicClientHandler
- * @Author: 刘苏义
- * @Date: 2023年11月30日9:27:55
- **/
-@Slf4j(topic = "netty")
-class DynamicClientHandler extends SimpleChannelInboundHandler<ByteBuf> {
- /**
- * 连接建立
- *
- * @param ctx
- * @throws Exception
- */
- @Override
- public void channelActive(ChannelHandlerContext ctx) {
- context = ctx;
- startHeartbeatTask();//开始发送心跳
- }
-
- @Override
- protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) {
- // 处理接收到的数据
- Channel channel = ctx.channel();
- ArdEquipRadar ardEquipRadar = DynamicClient.ConnectMap.get(channel);
- // 处理接收到的消息
- //byte[] byteArray = new byte[msg.readableBytes()];
- //msg.getBytes(msg.readerIndex(), byteArray);
- //byte[] bytes = messageParsing.receiveCompletePacket(byteArray);
- //if (bytes != null) {
- // processData(ardEquipRadar, bytes);
- //}
- }
-
- @Override
- public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
- // 处理异常
- // ...
- log.error("处理异常"+cause.getMessage());
- }
-
- private ScheduledFuture<?> heartbeatTask;
- private ChannelHandlerContext context;
-
- /**
- * 开始心跳任务
- */
- private void startHeartbeatTask() {
- heartbeatTask = context.executor().scheduleAtFixedRate(() -> {
- // 发送心跳消息
- ByteBuf message = context.alloc().buffer();
- byte[] header = {0x01, 0x02, 0x01};
- byte[] payload = {0x10, 0x00, 0x00, 0x00};
- byte[] payloadCrc32 = ByteUtils.parseCrc32(payload);
- byte[] footer = {0x01, 0x02, 0x00};
- byte[] heart = ByteUtils.appendArrays(header, payload, payloadCrc32, footer);
-// byte[] heart = {0x01, 0x02, 0x01, 0x10, 0x00, 0x00, 0x00, (byte) 0x83, (byte) 0x88, 0x5d, 0x71, 0x01, 0x02, 0x00};
- String hexString = DatatypeConverter.printHexBinary(heart);
- // log.debug("发送心跳:" + hexString);
- message.writeBytes(heart);
- context.writeAndFlush(message);
-
- }, 0, 5, TimeUnit.SECONDS);
- }
-
- /**
- * 停止心跳任务
- */
- private void stopHeartbeatTask() {
- if (heartbeatTask != null) {
- heartbeatTask.cancel(false);
- heartbeatTask = null;
- }
- }
-
- /**
- * 解析报警数据
- */
- public void processData(ArdEquipRadar ardEquipRadarbyte, byte[] data) {
- try {
- String radarId = ardEquipRadarbyte.getId();
- String radarName = ardEquipRadarbyte.getName();
- Double radarLongitude = ardEquipRadarbyte.getLongitude();
- Double radarLagitude = ardEquipRadarbyte.getLatitude();
- Double radarAltitude = ardEquipRadarbyte.getAltitude();
- //region crc校验-目前仅用于显示校验结果
- Boolean crc32Check = MessageParsing.CRC32Check(data);
- if (!crc32Check) {
- log.debug("CRC32校验不通过");
- } else {
- //log.debug("CRC32校验通过");
- }
- //endregion
- //log.info("原始数据:" + DatatypeConverter.printHexBinary(data));
- //log.info("雷达信息:" + host + "【port】" + port + "【X】" + longitude + "【Y】" + lagitude + "【Z】" + altitude);
- data = MessageParsing.transferData(data);//去掉包头和包尾、校验及转义
- //region 负载头解析
- byte[] type = Arrays.copyOfRange(data, 0, 1);//命令类型
- // log.info("命令类型:" + DatatypeConverter.printHexBinary(type));
- byte[] cmdId = Arrays.copyOfRange(data, 1, 2);//命令ID
- String cmdIdStr = DatatypeConverter.printHexBinary(cmdId);
- //log.info("命令ID:" + DatatypeConverter.printHexBinary(cmdId));
- byte[] payloadSize = Arrays.copyOfRange(data, 2, 4);//有效负载大小
- payloadSize = toLittleEndian(payloadSize);
- //log.info("payloadSize:" + DatatypeConverter.printHexBinary(payloadSize));
- int payloadSizeToDecimal = ByteUtils.bytesToDecimal(payloadSize);
- // log.info("有效负载大小(转整型):" + payloadSizeToDecimal);
- //endregion
- List<ArdAlarmRadar> radarAlarmInfos = new ArrayList<>();
- ArdAlarmRadar radarFollowInfo = null;
- //抽油机状态雷达推送集合
- List<ArdAlarmRadar> well = new ArrayList<>();
- String alarmTime = "";
- Integer targetNum = 0;
- log.debug("Processing radar data 【" + radarName + "】数据-->命令ID:" + cmdIdStr + "二进制:" + byteToBitString(cmdId[0]));
- //雷达移动防火报警
- if (Arrays.equals(cmdId, new byte[]{0x01})) {
- //region 告警信息反馈
- byte[] dwTim = Arrays.copyOfRange(data, 4, 8);
- dwTim = toLittleEndian(dwTim);
- SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
- long l = ByteUtils.bytesToDecimal(dwTim);
- alarmTime = sdf.format(l * 1000);
- // log.info("周视图像的出现时间(转date):" + alarmTime);
-
- byte[] wTargetNum = Arrays.copyOfRange(data, 8, 10);
- wTargetNum = toLittleEndian(wTargetNum);
- targetNum = ByteUtils.bytesToDecimal(wTargetNum);
- if (targetNum == 0) {
- return;
- }
- //log.debug("目标总点数(转整型):" + targetNum);
-
- //解析NET_TARGET_UNIT(64是NET_TARGET_HEAD的字节数)
- int uintSize = (payloadSizeToDecimal - 64) / targetNum;
- // log.info("单条报警大小:" + uintSize);
-
- for (int i = 0; i < targetNum; i++) {
-
- Integer index = 68 + uintSize * i;
- byte[] dwID = Arrays.copyOfRange(data, index, index + 4);
- // log.info("目标ID:" + DatatypeConverter.printHexBinary(cmdId));
- dwID = toLittleEndian(dwID);
- int targetId = ByteUtils.bytesToDecimal(dwID);
- // log.info("目标ID号:" + targetId);
-
- byte[] iDistance = Arrays.copyOfRange(data, index + 8, index + 12);
- iDistance = toLittleEndian(iDistance);
- double Distance = ByteUtils.bytesToDecimal(iDistance);
- //log.debug("目标当前直线距离(m):" + Distance);
-
- //region 不需要的字段
-// byte[] dwGSum = Arrays.copyOfRange(data, index + 4, index + 8);
-// dwGSum = toLittleEndian(dwGSum);
-// int GSum = byteArrayToDecimal(dwGSum);
-// log.info("目标当前像素灰度和:" + GSum);
-// byte[] iTw = Arrays.copyOfRange(data, index + 12, index + 16);
-// iTw = toLittleEndian(iTw);
-// int Tw = byteArrayToDecimal(iTw);
-// log.info("目标当前的像素宽度:" + Tw);
-//
-// byte[] iTh = Arrays.copyOfRange(data, index + 16, index + 20);
-// iTh = toLittleEndian(iTh);
-// int Th = byteArrayToDecimal(iTh);
-// log.info("目标当前的像素高度:" + Th);
-//
-// byte[] wPxlArea = Arrays.copyOfRange(data, index + 20, index + 22);
-// wPxlArea = toLittleEndian(wPxlArea);
-// int PxlArea = byteArrayToDecimal(wPxlArea);
-// log.info("目标当前像素面积:" + PxlArea);
-//
-// byte[] cTrkNum = Arrays.copyOfRange(data, index + 22, index + 23);
-// cTrkNum = toLittleEndian(cTrkNum);
-// int TrkNum = byteArrayToDecimal(cTrkNum);
-// log.info("轨迹点数:" + TrkNum);
-
-// byte[] sVx = Arrays.copyOfRange(data, index + 24, index + 26);
-// sVx = toLittleEndian(sVx);
-// int Vx = byteArrayToDecimal(sVx);
-// log.info("目标当前速度矢量(像素距离)X:" + Vx);
-//
-// byte[] sVy = Arrays.copyOfRange(data, index + 26, index + 28);
-// sVy = toLittleEndian(sVy);
-// int Vy = byteArrayToDecimal(sVy);
-// log.info("目标当前速度矢量(像素距离)Y:" + Vy);
-//
-// byte[] sAreaNo = Arrays.copyOfRange(data, index + 28, index + 30);
-// sAreaNo = toLittleEndian(sAreaNo);
-// int AreaNo = byteArrayToDecimal(sAreaNo);
-// log.info("目标归属的告警区域号:" + AreaNo);
-//
-// byte[] cGrp = Arrays.copyOfRange(data, index + 30, index + 31);
-// cGrp = toLittleEndian(cGrp);
-// int Grp = byteArrayToDecimal(cGrp);
-// log.info("所属组:" + Grp);
- //endregion
- String alarmType = "";
- byte[] cStat = Arrays.copyOfRange(data, index + 23, index + 24);
- log.info("原始状态:" + byteToBitString(cStat[0]));
- // cStat = toLittleEndian(cStat);
- // 提取第4位至第6位的值
- int extractedValue = (cStat[0] >> 4) & 0b00001111;
- // 判断提取的值
- if (extractedValue == 0b0000) {
- alarmType = "运动目标检测";
- } else if (extractedValue == 0b0001) {
- alarmType = "热源检测";
- }
- // log.info("报警类型:" + alarmType);
- byte[] szName = Arrays.copyOfRange(data, index + 64, index + 96);
- String alarmPointName = ByteUtils.bytesToStringZh(szName);
- // log.info("所属告警区域名称:" + alarmPointName);
- byte[] afTx = Arrays.copyOfRange(data, index + 96, index + 100);
- afTx = toLittleEndian(afTx);
- float fTx = ByteUtils.bytesToFloat(afTx);
- // log.info("水平角度:" + fTx);
- byte[] afTy = Arrays.copyOfRange(data, index + 112, index + 116);
- afTy = toLittleEndian(afTy);
- float fTy = ByteUtils.bytesToFloat(afTy);
- //log.debug("垂直角度:" + fTy);
- // 将角度转换为弧度
- double thetaRadians = Math.toRadians(fTy + 90);
- // 使用正弦函数计算对边长度
- Distance = Math.sin(thetaRadians) * Distance;
- //log.debug("目标投影距离(m):" + Distance);
-
- Double[] radarXY = {radarLongitude, radarLagitude};
- Double[] alarmXY = GisUtils.CalculateCoordinates(radarXY, Distance, (double) fTx);
- log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】" + targetId + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime + "【name】" + alarmPointName);
- ArdAlarmRadar ardAlarmRadar = new ArdAlarmRadar();
- ardAlarmRadar.setTargetId(targetId);
- ardAlarmRadar.setName(alarmPointName);
- ardAlarmRadar.setLongitude(alarmXY[0]);
- ardAlarmRadar.setLatitude(alarmXY[1]);
- ardAlarmRadar.setAlarmType(alarmType);
- radarAlarmInfos.add(ardAlarmRadar);
-
- int bit1 = (cStat[0] >> 1) & 0x1;
- //目标的B1=1 锁定
- if (bit1 == 1) {
- radarFollowInfo = ardAlarmRadar;
- //将追踪锁定的报警对象属性复制给radarFollowInfo对象
- //BeanUtils.copyProperties(ardAlarmRadar, radarFollowInfo);
- }
- }
- //endregion
- if (StringUtils.isEmpty(alarmTime)) {
- return;
- }
- if (targetNum == 0) {
- return;
- }
- RadarAlarmData radarAlarmData = new RadarAlarmData();
- radarAlarmData.setRadarId(radarId);
- radarAlarmData.setRadarName(radarName);
- radarAlarmData.setAlarmTime(alarmTime);
- radarAlarmData.setArdAlarmRadars(radarAlarmInfos);
- MqttProducer.publish(2, false, "radar", JSON.toJSONString(radarAlarmData));
- if (radarFollowInfo != null) {
- //当前雷达扫描存在引导跟踪数据,只保留最后一次锁定的数据
- MqttProducer.publish(2, false, "radarFollowGuide", JSON.toJSONString(radarFollowInfo));
- }
- //抽油机状态MQTT队列
- radarAlarmData.setArdAlarmRadars(well);
- MqttProducer.publish(2, false, "radarWellData", JSON.toJSONString(radarAlarmData));
-
- }
- //抽油机AI状态反馈
- if (Arrays.equals(cmdId, new byte[]{0x04})) {
- //region抽油机AI状态反馈
- String hexString = DatatypeConverter.printHexBinary(data);
- //log.info(hexString);
-
- byte[] dwTim = Arrays.copyOfRange(data, 4, 8);
- dwTim = toLittleEndian(dwTim);
- SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
- long l = ByteUtils.bytesToDecimal(dwTim);
- alarmTime = sdf.format(l * 1000);
- //log.info("周视图像的出现时间(转date):" + alarmTime);
-
- byte[] wTargetNum = Arrays.copyOfRange(data, 8, 10);
- wTargetNum = toLittleEndian(wTargetNum);
- targetNum = ByteUtils.bytesToDecimal(wTargetNum);
- //log.debug("目标总点数(转整型):" + targetNum);
- if (targetNum == 0) {
- return;
- }
- //解析NET_TARGET_UNIT(64是NET_TARGET_HEAD的字节数)
- int uintSize = (payloadSizeToDecimal - 64) / targetNum;
- //log.info("单条报警大小:" + uintSize);
- for (int i = 0; i < targetNum; i++) {
- Integer index = 68 + uintSize * i;
- byte[] dwID = Arrays.copyOfRange(data, index, index + 4);
- //log.info("目标ID:" + DatatypeConverter.printHexBinary(dwID));
- dwID = toLittleEndian(dwID);
- int targetId = ByteUtils.bytesToDecimal(dwID);
- //log.info("目标ID号:" + targetId);
- //region 不需要的字段
- byte[] iTw = Arrays.copyOfRange(data, index + 4, index + 8);
- iTw = toLittleEndian(iTw);
- int Tw = ByteUtils.bytesToDecimal(iTw);
- // log.info("目标当前的像素宽度:" + Tw);
-
- byte[] iTh = Arrays.copyOfRange(data, index + 8, index + 12);
- iTh = toLittleEndian(iTh);
- int Th = ByteUtils.bytesToDecimal(iTh);
- //log.info("目标当前的像素高度:" + Th);
-
- byte[] fTx = Arrays.copyOfRange(data, index + 12, index + 16);
- fTx = toLittleEndian(fTx);
- float fTxAngle = ByteUtils.bytesToFloat(fTx);
- //log.debug("p角度:" + fTxAngle);
- byte[] fTy = Arrays.copyOfRange(data, index + 16, index + 20);
- fTy = toLittleEndian(fTy);
- float fTyAngle = ByteUtils.bytesToFloat(fTy);
- //log.debug("t角度:" + fTyAngle);
-
- byte[] sAreaNo = Arrays.copyOfRange(data, index + 20, index + 22);
- sAreaNo = toLittleEndian(sAreaNo);
- int AreaNo = ByteUtils.bytesToDecimal(sAreaNo);
- //log.debug("目标归属的告警区域号:" + AreaNo);
-
- byte[] cGrp = Arrays.copyOfRange(data, index + 22, index + 23);
- cGrp = toLittleEndian(cGrp);
- int Grp = ByteUtils.bytesToDecimal(cGrp);
- //log.info("所属组:" + Grp);
- //endregion
- String alarmType;
- //抽油机状态变量
- String wellType;
- byte[] cStat = Arrays.copyOfRange(data, index + 23, index + 24);
- cStat = toLittleEndian(cStat);
- //String binaryString = String.format("%8s", Integer.toBinaryString(cStat[0] & 0xFF)).replace(' ', '0');
- //log.info("目标当前状态:" + binaryString);
- // 提取第0位值
- // 使用位运算操作判断第0位是否为1
- boolean isB0 = (cStat[0] & 0x01) == 0x00;
- // 判断提取的值
- if (isB0) {
- alarmType = "雷达抽油机停机";
- byte[] szName = Arrays.copyOfRange(data, index + 32, index + 64);
- //log.info("所属告警区域名称:" + DatatypeConverter.printHexBinary(szName));
- String alarmPointName = ByteUtils.bytesToStringZh(szName);
- // log.info("所属告警区域名称:" + alarmPointName);
- //log.debug("报警信息:"+ "【radarName】" + radarName + "【targetId】" + targetId + "【name】" + alarmPointName + "【alarmType】" + alarmType + "【alarmTime】" + alarmTime);
- ArdAlarmRadar ardAlarmRadar = new ArdAlarmRadar();
- ardAlarmRadar.setTargetId(targetId);
- ardAlarmRadar.setName(alarmPointName);
- ardAlarmRadar.setAlarmType(alarmType);
- radarAlarmInfos.add(ardAlarmRadar);
- wellType = "停机";
- } else {
- wellType = "运行";
- }
- //抽油机状态集合中装入数据
- byte[] szName = Arrays.copyOfRange(data, index + 32, index + 64);
- String alarmPointName = ByteUtils.bytesToStringZh(szName);
- log.debug("报警信息:" + "【radarName】" + radarName + "【targetId】" + targetId + "【alarmType】抽油机状态报警" + "【alarmTime】" + alarmTime + "【name】" + alarmPointName + "【alarmState】" + wellType);
- ArdAlarmRadar wellAlarm = new ArdAlarmRadar();
- wellAlarm.setTargetId(targetId);
- wellAlarm.setName(alarmPointName);
- wellAlarm.setAlarmType(wellType);
- well.add(wellAlarm);
- }
- //endregion
- if (StringUtils.isEmpty(alarmTime)) {
- return;
- }
- if (targetNum == 0) {
- return;
- }
- RadarAlarmData radarAlarmData = new RadarAlarmData();
- radarAlarmData.setRadarId(radarId);
- radarAlarmData.setRadarName(radarName);
- radarAlarmData.setAlarmTime(alarmTime);
- radarAlarmData.setArdAlarmRadars(radarAlarmInfos);
- MqttProducer.publish(2, false, "radar", JSON.toJSONString(radarAlarmData));
- //抽油机状态MQTT队列
- radarAlarmData.setArdAlarmRadars(well);
- MqttProducer.publish(2, false, "radarWellData", JSON.toJSONString(radarAlarmData));
- }
- //强制引导
- if (Arrays.equals(cmdId, new byte[]{0x02})) {
- //region 告警前端发送的强制引导信息
- byte[] iDistance = Arrays.copyOfRange(data, 4, 8);
- iDistance = toLittleEndian(iDistance);
- long distance = ByteUtils.bytesToDecimal(iDistance);
- log.info("目标当前距离(m):" + distance);
- byte[] fTx = Arrays.copyOfRange(data, 8, 12);
- fTx = toLittleEndian(fTx);
- float tx = ByteUtils.bytesToFloat(fTx);
- log.debug("方位:" + tx);
- byte[] fTy = Arrays.copyOfRange(data, 12, 16);
- fTy = toLittleEndian(fTy);
- float ty = ByteUtils.bytesToFloat(fTy);
- if (ty < 0) {
- ty += 360;
- }
- log.debug("俯仰:" + ty);
- Map<String, Object> forceGuideMap = new HashMap<>();
- forceGuideMap.put("distance", distance);
- forceGuideMap.put("p", tx);
- forceGuideMap.put("t", ty);
- forceGuideMap.put("radarId", radarId);
- log.debug("强制引导信息" + forceGuideMap);
- //endregion
- MqttProducer.publish(2, false, "radarForceGuide", JSON.toJSONString(forceGuideMap));
- }
- } catch (Exception ex) {
- log.error("雷达报文解析异常:" + ex.getMessage());
- }
- }
-}
diff --git a/src/main/java/com/ard/utils/netty/tcp/DynamicClientInitializer.java b/src/main/java/com/ard/utils/netty/tcp/DynamicClientInitializer.java
deleted file mode 100644
index 909c5d9..0000000
--- a/src/main/java/com/ard/utils/netty/tcp/DynamicClientInitializer.java
+++ /dev/null
@@ -1,25 +0,0 @@
-package com.ard.utils.netty.tcp;
-
-import io.netty.channel.ChannelInitializer;
-import io.netty.channel.ChannelPipeline;
-import io.netty.channel.socket.SocketChannel;
-
-/**
- * @Description: 初始化客户端的通道(备用)
- * @ClassName: DynamicClientInitializer
- * @Author: 刘苏义
- * @Date: 2023年11月30日9:27:03
- **/
-public class DynamicClientInitializer extends ChannelInitializer<SocketChannel> {
- @Override
- protected void initChannel(SocketChannel ch){
- try {
- ChannelPipeline pipeline = ch.pipeline();
- // 添加你需要的处理器
- pipeline.addLast(new DynamicClientHandler());
- }catch (Exception e)
- {
- e.printStackTrace();
- }
- }
-}
diff --git a/src/main/java/com/ard/utils/netty/tcp/MessageParsing.java b/src/main/java/com/ard/utils/netty/tcp/MessageHandler.java
similarity index 98%
rename from src/main/java/com/ard/utils/netty/tcp/MessageParsing.java
rename to src/main/java/com/ard/utils/netty/tcp/MessageHandler.java
index 4aac19d..e6286fb 100644
--- a/src/main/java/com/ard/utils/netty/tcp/MessageParsing.java
+++ b/src/main/java/com/ard/utils/netty/tcp/MessageHandler.java
@@ -14,7 +14,7 @@
* @Date: 2023年07月03日15:30
* @Version: 1.0
**/
-public class MessageParsing {
+public class MessageHandler {
// 创建缓冲区列表
private List<Byte> buffer = new ArrayList<>();
--
Gitblit v1.9.3