|
|
@@ -0,0 +1,420 @@
|
|
|
+/**
|
|
|
+ * Copyright 2019 DH
|
|
|
+ * All right reserved.
|
|
|
+ * 项目名称: 大恒泰山系统
|
|
|
+ * 创建日期:2021/12/16
|
|
|
+ */
|
|
|
+package org.ts.ddcs.dataStream.processor.sw;
|
|
|
+
|
|
|
+
|
|
|
+import com.alibaba.fastjson.JSONObject;
|
|
|
+import lombok.extern.slf4j.Slf4j;
|
|
|
+import org.ts.ddcs.common.constant.SwDatagramConstant;
|
|
|
+import org.ts.ddcs.common.dict.BaseSettingFlag;
|
|
|
+import org.ts.ddcs.enums.SwElementCodeEnum;
|
|
|
+import org.ts.ddcs.common.metadata.DatagramMetadataEnum;
|
|
|
+import org.ts.ddcs.enums.SwDatagramTypeCodeEnum;
|
|
|
+import org.ts.ddcs.common.util.BytesHelp;
|
|
|
+import org.ts.ddcs.common.util.Sw2014DatagramTemplateHelp;
|
|
|
+import org.ts.ddcs.common.util.ValueLen;
|
|
|
+import org.ts.ddcs.dataStream.*;
|
|
|
+import org.ts.ddcs.datatemplate.TestState;
|
|
|
+import org.ts.ddcs.itemdatatemplate.DataItemTemplate;
|
|
|
+import org.ts.ddcs.metadata.DataStreamCacheMetadata;
|
|
|
+
|
|
|
+import java.io.UnsupportedEncodingException;
|
|
|
+import java.util.*;
|
|
|
+
|
|
|
+/***
|
|
|
+ * Date:2021/12/16
|
|
|
+ * Title:数据处理器模块
|
|
|
+ * Description:读取基本参数回报 0x41
|
|
|
+ * @author dylan
|
|
|
+ * @version 1.0
|
|
|
+ * Remark:认为有必要的其他信息
|
|
|
+ */
|
|
|
+@Slf4j
|
|
|
+public class DataStreamSwReadSettingDatagramProcessor extends DataStreamProcessor {
|
|
|
+
|
|
|
+ private DataStreamCache dataStreamCache;
|
|
|
+ private Map<String, List<Map<String,Object>>> keyStore = new HashMap<>();
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void init(DataStreamLineManager context) {
|
|
|
+ this.context = context;
|
|
|
+ this.dataStreamCache = context.getCache();
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void process(DataStreamAdapter dataStreamAdapter) throws UnsupportedEncodingException {
|
|
|
+ this.dataStreamCache.setState(TestState.MATHCE);
|
|
|
+ this.dataStreamCache.setMatchProcessor(this);
|
|
|
+ String datagramCode = this.dataStreamCache.getDatagramCode();
|
|
|
+ if (SwDatagramTypeCodeEnum.DATAGRAM_CODE_41.getName().equals(datagramCode)) {
|
|
|
+ this.keyStore.clear();
|
|
|
+ String agreement = dataStreamAdapter.getAgreement();
|
|
|
+ Map<String, DataItemTemplate> itemTemplateMap = this.context.getItemTemplate(agreement);
|
|
|
+ //解析报文
|
|
|
+ byte[] dataArea = this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_DATA_AREA.getName());
|
|
|
+ int count = 0;
|
|
|
+ //流水号
|
|
|
+ byte[] serialNo = BytesHelp.subBytes(dataArea, count, 2);
|
|
|
+ dataStreamCache.putValue(DataStreamCacheMetadata.METADATA_SERIAL_NO.getName(), serialNo);
|
|
|
+ count += 2;
|
|
|
+ //String serialNo = Sw2014DatagramTemplateHelp.toSerialNo(dataArea, count);
|
|
|
+ // dataStreamCache.putValue(DatagramConstant.serialNo, serialNo);
|
|
|
+ // 发报时间
|
|
|
+ byte[] sendPacketTime = BytesHelp.subBytes(dataArea, count, 6);
|
|
|
+ dataStreamCache.putValue(DataStreamCacheMetadata.METADATA_UP_TIME_BYTE.getName(), sendPacketTime);
|
|
|
+ count += 6;
|
|
|
+ //String sendPacketTime = Sw2014DatagramTemplateHelp.toSendPacketTime(dataArea, count);
|
|
|
+ // dataStreamCache.putValue(DatagramConstant.sendPacketTime, sendPacketTime);
|
|
|
+ //遥测站地址
|
|
|
+ DataItemTemplate itemTemplateSt = itemTemplateMap.get(SwElementCodeEnum.CODE_ST.getCode());
|
|
|
+ // itemTemplateSt.clear();
|
|
|
+ int readSize = itemTemplateSt.stream(dataArea, count);
|
|
|
+ if (itemTemplateSt.isMacth()) {
|
|
|
+ itemTemplateSt.saveValue(this.keyStore);
|
|
|
+ }
|
|
|
+ count += readSize;
|
|
|
+
|
|
|
+ while (true) {
|
|
|
+ if (count >= dataArea.length) {
|
|
|
+ break;
|
|
|
+ }
|
|
|
+ byte[] f = BytesHelp.subBytes(dataArea, count, 1);
|
|
|
+ count += 1;
|
|
|
+ String flagCode = BytesHelp.byte2HexStr(f);
|
|
|
+ ValueLen valueLen = Sw2014DatagramTemplateHelp.getValueLen(dataArea, count);
|
|
|
+ int dataTotalLen = valueLen.getByteLen();
|
|
|
+ int decimalLen = valueLen.getDecimalLen();
|
|
|
+ count += 1;
|
|
|
+
|
|
|
+ if (flagCode.equals(BaseSettingFlag.CENTER_ADD.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.RTU_CODE.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("rtucode", BytesHelp.byte2HexStr(value));
|
|
|
+ data.put("code", BaseSettingFlag.RTU_CODE.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.RTU_CODE.getId(), list);
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.PASSWORD.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP1.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ if (value[0] != 0x00) {
|
|
|
+ data.put("ip1", BytesHelp.toIp(BytesHelp.subBytes(value, 1, 6)));
|
|
|
+ data.put("ip1port", String.format("%d", Integer.parseInt(BytesHelp.byte2HexStr(BytesHelp.subBytes(value, 7, 3)))));
|
|
|
+ data.put("code", BaseSettingFlag.IP1.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.IP1.getId(), list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP1_SUB.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP2.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ if (value[0] != 0x00) {
|
|
|
+ data.put("ip2", BytesHelp.toIp(BytesHelp.subBytes(value, 1, 6)));
|
|
|
+ data.put("ip2port", String.format("%d", Integer.parseInt(BytesHelp.byte2HexStr(BytesHelp.subBytes(value, 7, 3)))));
|
|
|
+ data.put("code", BaseSettingFlag.IP2.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.IP2.getId(), list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP2_SUB.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP3.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ if (value[0] != 0x00) {
|
|
|
+ data.put("ip3", BytesHelp.toIp(BytesHelp.subBytes(value, 1, 6)));
|
|
|
+ data.put("ip3port", String.format("%d", Integer.parseInt(BytesHelp.byte2HexStr(BytesHelp.subBytes(value, 7, 3)))));
|
|
|
+ data.put("code", BaseSettingFlag.IP3.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.IP3.getId(), list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP3_SUB.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP4.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ if (value[0] != 0x00) {
|
|
|
+ data.put("ip4", BytesHelp.toIp(BytesHelp.subBytes(value, 1, 6)));
|
|
|
+ data.put("ip4port", String.format("%d", Integer.parseInt(BytesHelp.byte2HexStr(BytesHelp.subBytes(value, 7, 3)))));
|
|
|
+ data.put("code", BaseSettingFlag.IP4.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.IP4.getId(), list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.IP4_SUB.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals("34")) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, 39);
|
|
|
+ String ipText = BytesHelp.bytes2String(value, 0, value.length);
|
|
|
+ count += 39;
|
|
|
+ byte[] port = BytesHelp.subBytes(dataArea, count, 2);
|
|
|
+ String portText = "" + BytesHelp.bytes2short(port);
|
|
|
+ count += 2;
|
|
|
+ if (null != portText && !portText.equals("0")) {
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("ip1", ipText);
|
|
|
+ data.put("ip1port", portText);
|
|
|
+ data.put("code", "34");
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put("34", list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals("35")) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, 39);
|
|
|
+ String ipText = BytesHelp.bytes2String(value, 0, value.length);
|
|
|
+ count += 39;
|
|
|
+ byte[] port = BytesHelp.subBytes(dataArea, count, 2);
|
|
|
+ String portText = "" + BytesHelp.bytes2short(port);
|
|
|
+ count += 2;
|
|
|
+ if (null != portText && !portText.equals("0")) {
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("ip2", ipText);
|
|
|
+ data.put("ip2port", portText);
|
|
|
+ data.put("code", "35");
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put("35", list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals("36")) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, 39);
|
|
|
+ String ipText = BytesHelp.bytes2String(value, 0, value.length);
|
|
|
+ count += 39;
|
|
|
+ byte[] port = BytesHelp.subBytes(dataArea, count, 2);
|
|
|
+ String portText = "" + BytesHelp.bytes2short(port);
|
|
|
+ count += 2;
|
|
|
+ if (null != portText && !portText.equals("0")) {
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("ip3", ipText);
|
|
|
+ data.put("ip3port", portText);
|
|
|
+ data.put("code", "36");
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put("36", list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals("37")) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, 39);
|
|
|
+ String ipText = BytesHelp.bytes2String(value, 0, value.length);
|
|
|
+ count += 39;
|
|
|
+ byte[] port = BytesHelp.subBytes(dataArea, count, 2);
|
|
|
+ String portText = "" + BytesHelp.bytes2short(port);
|
|
|
+ count += 2;
|
|
|
+ if (null != portText && !portText.equals("0")) {
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("ip4", ipText);
|
|
|
+ data.put("ip4port", portText);
|
|
|
+ data.put("code", "37");
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put("37", list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.WORK_TYPE.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.COLLECT_SETTING.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ if (value.length == 8) {
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ for (int i = 0; i < 8; i++) {
|
|
|
+ String binStr = BytesHelp.BytesToBinStr(BytesHelp.subBytes(value, i, 1));
|
|
|
+ binStr = new StringBuilder(binStr).reverse().toString();
|
|
|
+ for (int j = 0; j < 8; j++) {
|
|
|
+ String key = "G_" + String.format("%02d", i) + "_V_" + String.format("%02d", j);
|
|
|
+ data.put(key, binStr.substring(j, j + 1));
|
|
|
+ }
|
|
|
+ }
|
|
|
+ data.put("code", BaseSettingFlag.COLLECT_SETTING.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.COLLECT_SETTING.getId(), list);
|
|
|
+ }
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.SERVER_ADD.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.DEVICE_ID.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.WORK_MODEL.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ if (value[0] == 0x00) {
|
|
|
+ data.put("workModel", "0");
|
|
|
+ } else {
|
|
|
+ data.put("workModel", "1");
|
|
|
+ }
|
|
|
+ data.put("code", BaseSettingFlag.WORK_MODEL.getId());
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.WORK_MODEL.getId(), list);
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.RTU_KIND.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("code", BaseSettingFlag.RTU_KIND.getId());
|
|
|
+ if (value[0] == 0x50) {
|
|
|
+ data.put("rtuKind", "P");
|
|
|
+ } else if (value[0] == 0x48) {
|
|
|
+ data.put("rtuKind", "H");
|
|
|
+ } else if (value[0] == 0x4B) {
|
|
|
+ data.put("rtuKind", "K");
|
|
|
+ } else if (value[0] == 0x5A) {
|
|
|
+ data.put("rtuKind", "Z");
|
|
|
+ } else if (value[0] == 0x44) {
|
|
|
+ data.put("rtuKind", "D");
|
|
|
+ } else if (value[0] == 0x54) {
|
|
|
+ data.put("rtuKind", "T");
|
|
|
+ } else if (value[0] == 0x4D) {
|
|
|
+ data.put("rtuKind", "M");
|
|
|
+ } else if (value[0] == 0x47) {
|
|
|
+ data.put("rtuKind", "G");
|
|
|
+ } else if (value[0] == 0x51) {
|
|
|
+ data.put("rtuKind", "Q");
|
|
|
+ } else if (value[0] == 0x49) {
|
|
|
+ data.put("rtuKind", "I");
|
|
|
+ } else if (value[0] == 0x4F) {
|
|
|
+ data.put("rtuKind", "O");
|
|
|
+ }
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.RTU_KIND.getId(), list);
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.RAIN_PERCEN.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("code", BaseSettingFlag.RAIN_PERCEN.getId());
|
|
|
+ data.put("rainPercenValue", BytesHelp.byte2HexStr(value).substring(1, 2));
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.RAIN_PERCEN.getId(), list);
|
|
|
+ } else if (flagCode.equals(BaseSettingFlag.BLUETOOTH.getId())) {
|
|
|
+ byte[] value = BytesHelp.subBytes(dataArea, count, dataTotalLen);
|
|
|
+ count += dataTotalLen;
|
|
|
+ Map<String,Object> data = new HashMap<>();
|
|
|
+ data.put("code", BaseSettingFlag.BLUETOOTH.getId());
|
|
|
+ data.put("bluetoothValue", BytesHelp.byte2HexStr(value).substring(1, 2));
|
|
|
+ List<Map<String,Object>> list = new ArrayList<>();
|
|
|
+ list.add(data);
|
|
|
+ this.keyStore.put(BaseSettingFlag.BLUETOOTH.getId(), list);
|
|
|
+ } else {
|
|
|
+ count += dataTotalLen;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ this.next(false);
|
|
|
+ } else {
|
|
|
+ this.next(true);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void close(DataStreamAdapter dataStreamAdapter) {
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void batchComplete(DataStreamAdapter dataStreamAdapter) {
|
|
|
+ log.info("接收到41读取基本参数,回复报文");
|
|
|
+ //帧起始符
|
|
|
+ byte[] respPacket = new byte[2];
|
|
|
+ respPacket[0] = 0x7e;
|
|
|
+ respPacket[1] = 0x7e;
|
|
|
+ //遥测站地址
|
|
|
+ byte[] rtuCodeByte = this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_RTU_CODE.getName());
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, rtuCodeByte);
|
|
|
+ //中心站地址
|
|
|
+ byte[] centerStationByte = this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_CENTER_STATION_CODE.getName());
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, centerStationByte);
|
|
|
+ //密码
|
|
|
+ byte[] rtuPwByte = this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_SW_PW.getName());
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, rtuPwByte);
|
|
|
+ //功能码
|
|
|
+ byte[] funcCodeByte = this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_SW_AFN.getName());
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, funcCodeByte);
|
|
|
+ //报文下行标识及长度
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, BytesHelp.hexStr2Bytes("8008"));
|
|
|
+ //报文起始符
|
|
|
+ byte[] dataAreaStartFlagByte =this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_DATA_AREA_START_FLAG.getName());
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, dataAreaStartFlagByte);
|
|
|
+ //流水号
|
|
|
+ byte[] serialNoByte = this.dataStreamCache.getBytesValue(DataStreamCacheMetadata.METADATA_SERIAL_NO.getName());
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, serialNoByte);
|
|
|
+ //发报时间
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, Sw2014DatagramTemplateHelp.getFBTm());
|
|
|
+ //RTU2014年报文正文结束符
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, SwDatagramConstant.PACKET_CONTROL_FLAG_ESC);
|
|
|
+ //CRC16
|
|
|
+ byte[] crc16 = Sw2014DatagramTemplateHelp.crcFun2014(respPacket);
|
|
|
+ respPacket = BytesHelp.arrayApend(respPacket, crc16);
|
|
|
+ dataStreamAdapter.output(respPacket);
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public List<Map<String, Object>> sink(DataStreamAdapter dataStreamAdapter) {
|
|
|
+ try {
|
|
|
+ List<Map<String, Object>> elementDataList = new LinkedList<>();
|
|
|
+ Set<String> keys = this.keyStore.keySet();
|
|
|
+ for (String key : keys) {
|
|
|
+ Map<String, Object> data = new HashMap<>();
|
|
|
+// data.put(JsonPackMetadataEnum.PACK_METADATA_DATAGRAM.getName(), BytesHelp.byte2HexStr(this.dataStreamCache.getDatagramBuff()));
|
|
|
+// data.put(JsonPackMetadataEnum.PACK_METADATA_FROM_TIME.getName(), this.dataStreamCache.getPickPacketTime());
|
|
|
+ // data.put(SinkJsonPackMetadataEnum.PACK_METADATA_FROM_TIME.getName(), this.dataStreamCache.getValue(DatagramConstant.pickPacketTime));
|
|
|
+// data.put(JsonPackMetadataEnum.PACK_METADATA_RTU_CODE.getName(), this.dataStreamCache.getRtuCode());
|
|
|
+//
|
|
|
+// String collectTime = this.dataStreamCache.getStringValue(DatagramConstant.collectTime);
|
|
|
+// if (null != collectTime) {
|
|
|
+// data.put(JsonPackMetadataEnum.PACK_METADATA_COLLECT_TIME.getName(), collectTime);
|
|
|
+// } else {
|
|
|
+// data.put(JsonPackMetadataEnum.PACK_METADATA_COLLECT_TIME.getName(), DatagramHelp.getCollectTime(this.keyStore));
|
|
|
+// }
|
|
|
+ // data.put(SinkJsonPackMetadataEnum.PACK_METADATA_COLLECT_TIME.getName(), this.dataStreamCache.getStringValue(DatagramConstant.collectTime));
|
|
|
+ // data.put(JsonPackMetadataEnum.PACK_METADATA_UP_TIME.getName(), this.dataStreamCache.getStringValue(DatagramConstant.sendPacketTime));
|
|
|
+ data.put(DatagramMetadataEnum.CACHE_METADATA_ELEMENT_CODE.getName(), key);
|
|
|
+ List<Map<String,Object>> elements = this.keyStore.get(key);
|
|
|
+ for (Map<String,Object> jsonObject : elements) {
|
|
|
+ log.info("element {}", JSONObject.toJSONString(jsonObject));
|
|
|
+ }
|
|
|
+ if (null != elements) {
|
|
|
+ data.put(DatagramMetadataEnum.CACHE_METADATA_DATA.getName(), elements);
|
|
|
+ }
|
|
|
+ elementDataList.add(data);
|
|
|
+ }
|
|
|
+ return elementDataList;
|
|
|
+ } catch (Exception e) {
|
|
|
+ return null;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+}
|