From 66b091963fb0f3356f27ec094c013369bf91db89 Mon Sep 17 00:00:00 2001 From: jiazx0107@163.com <jiazx0107@163.com> Date: 星期日, 24 十二月 2023 14:02:19 +0800 Subject: [PATCH] 游仙协议解析-3 --- src/main/java/com/fzzy/gateway/hx2023/service/DeviceReportServiceImpl.java | 134 ++++++++++++++++++++++++++++++++++++++------ 1 files changed, 115 insertions(+), 19 deletions(-) diff --git a/src/main/java/com/fzzy/gateway/hx2023/service/DeviceReportServiceImpl.java b/src/main/java/com/fzzy/gateway/hx2023/service/DeviceReportServiceImpl.java index 7f76daa..2d572be 100644 --- a/src/main/java/com/fzzy/gateway/hx2023/service/DeviceReportServiceImpl.java +++ b/src/main/java/com/fzzy/gateway/hx2023/service/DeviceReportServiceImpl.java @@ -1,56 +1,152 @@ package com.fzzy.gateway.hx2023.service; -import com.fzzy.api.data.GatewayProtocol; +import com.alibaba.fastjson2.JSONObject; import com.fzzy.api.data.PushProtocol; -import com.fzzy.gateway.api.DeviceReportService; +import com.fzzy.data.ConfigData; +import com.fzzy.gateway.api.GatewayDeviceReportService; +import com.fzzy.gateway.data.BaseReqData; +import com.fzzy.gateway.data.BaseResp; import com.fzzy.gateway.entity.GatewayDevice; +import com.fzzy.gateway.hx2023.ScConstant; +import com.fzzy.gateway.hx2023.data.LprData; import com.fzzy.gateway.hx2023.data.WebSocketPacket; import com.fzzy.gateway.hx2023.data.WebSocketPacketHeader; -import com.fzzy.gateway.hx2023.websocket.WebSocketDeviceReport; +import com.fzzy.gateway.hx2023.data.WeightInfo; +import com.fzzy.gateway.hx2023.kafka.KafkaDeviceReportService; +import com.fzzy.mqtt.MqttGatewayService; +import jdk.nashorn.internal.runtime.regexp.joni.Config; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang.StringUtils; import org.springframework.stereotype.Component; import javax.annotation.Resource; @Slf4j @Component -public class DeviceReportServiceImpl implements DeviceReportService { - +public class DeviceReportServiceImpl implements GatewayDeviceReportService { @Resource - private WebSocketDeviceReport webSocketDeviceReport; + private KafkaDeviceReportService kafkaDeviceReportService; + @Resource + private MqttGatewayService publishService; + @Resource + private ConfigData configData; @Override - public String getProvinceProtocol() { + public String getProtocol() { return PushProtocol.GATEWAY_SC_2023.getCode(); } @Override - public String report2GatewayBySn(double weigh, GatewayDevice device) { + public BaseResp reportGrainData(BaseReqData reqData) { - if (null == device) { - log.error("-----------娌℃湁鑾峰彇鍒拌澶囬厤缃俊鎭�-----"); - return "ERROR:娌℃湁鑾峰彇鍒拌澶囬厤缃俊鎭�"; + String topic = ScConstant.TOPIC_REPORT; + topic = topic.replace("${productId}", reqData.getProductId()).replace("${deviceId}", reqData.getDeviceId()); + + //濡傛灉鏄祴璇曟ā寮忎笉鎵ц鎺ㄩ�� + if(configData.getActive().indexOf("dev")>=0){ + + log.info("----------------------------鎺ㄩ�丮QTT绮儏淇℃伅锛屾敞锛氳皟璇曟ā寮忎笉鎺ㄩ��---------------------------"); + log.info("-----TOPIC-----{}", topic); + log.info("-----Message-----{}", reqData.getData()); + + return new BaseResp(); } - //浣跨敤WEBSOCKET杩斿洖 - if (GatewayProtocol.GATE_WEBSOCKET.equals(device.getPushProtocol())) { + publishService.publishMqttWithTopic(reqData.getData(), topic); + + log.info("----------------------------鎺ㄩ�丮QTT绮儏淇℃伅---------------------------"); + log.info("-----TOPIC-----{}", topic); + log.info("-----Message-----{}", reqData.getData()); + + return new BaseResp(); + } + + @Override + public BaseResp reportWeightData(BaseReqData reqData) { + + String topic = ScConstant.TOPIC_MESSAGE_REPORT; + + topic = topic.replace("${productId}", reqData.getProductId()).replace("${deviceId}", reqData.getDeviceId()); + + if (null == reqData.getData()) { + GatewayDevice device = reqData.getDevice(); WebSocketPacket packet = new WebSocketPacket(); WebSocketPacketHeader header = new WebSocketPacketHeader(); header.setDeviceName(device.getDeviceName()); + header.setProductId(device.getProductId()); packet.setHeaders(header); - packet.setMessageType(""); + packet.setMessageType(ScConstant.MESSAGE_TYPE_REPORT_PROPERTY); packet.setDeviceId(device.getDeviceId()); - packet.setProperties(null); + + //璁剧疆淇℃伅涓讳綋 + WeightInfo weightInfo = new WeightInfo(); + weightInfo.setGrossWeight(reqData.getWeight()); + weightInfo.setNetWeight(reqData.getWeight()); + weightInfo.setNetWeight(reqData.getWeight()); + weightInfo.setWeightUnit("KG"); + JSONObject jsonObject = new JSONObject(); + jsonObject.put("weightInfo", JSONObject.toJSONString(weightInfo)); + + packet.setProperties(jsonObject); + packet.setTimestamp(System.currentTimeMillis()); - - webSocketDeviceReport.sendByPacket(packet); - + reqData.setData(JSONObject.toJSONString(packet)); } - return null; + publishService.publishMqttWithTopic(reqData.getData(), topic); + + log.info("----------------------------鎺ㄩ�丮QTT鍦扮淇℃伅---------------------------"); + log.info("-----TOPIC-----{}", topic); + log.info("-----Message-----{}", reqData.getData()); + return new BaseResp(); + } + + @Override + public BaseResp reportLprData(BaseReqData reqData) { + String topic = ScConstant.TOPIC_MESSAGE_REPORT; + topic = topic.replace("${productId}", reqData.getProductId()).replace("${deviceId}", reqData.getDeviceId()); + + GatewayDevice device = reqData.getDevice(); + + if (StringUtils.isEmpty(reqData.getData())) { + WebSocketPacket packet = new WebSocketPacket(); + WebSocketPacketHeader header = new WebSocketPacketHeader(); + header.setDeviceName(reqData.getDeviceName()); + header.setProductId(reqData.getProductId()); + + packet.setHeaders(header); + packet.setMessageType(ScConstant.MESSAGE_TYPE_REPORT_PROPERTY); + packet.setDeviceId(reqData.getDeviceId()); + packet.setMessageId(System.currentTimeMillis() + ""); + //璁剧疆淇℃伅涓讳綋 + LprData lpr = new LprData(); + lpr.setDeviceId(reqData.getDeviceId()); + lpr.setCarNumber(reqData.getCarNumber()); + JSONObject jsonObject = new JSONObject(); + jsonObject.put("carNumber", reqData.getCarNumber()); + jsonObject.put("position", device.getPosition()); + packet.setProperties(jsonObject); + packet.setTimestamp(System.currentTimeMillis()); + + reqData.setData(JSONObject.toJSONString(packet)); + } + + publishService.publishMqttWithTopic(reqData.getData(), topic); + + log.info("----------------------------鎺ㄩ�丮QTT杞︾墝璇嗗埆淇℃伅---------------------------"); + log.info("-----TOPIC-----{}", topic); + log.info("-----Message-----{}", reqData.getData()); + return new BaseResp(); + } + + @Override + public BaseResp reportGrainDataByKafka(BaseReqData reqData) { + String topic = ScConstant.TOPIC_MESSAGE_REPORT; + kafkaDeviceReportService.publishWithTopic(reqData.getData(), topic); + return new BaseResp(); } } -- Gitblit v1.9.3