vince
2024-06-14 a230c790107e4be5743473849bbc9aeb93583ddd
src/main/java/com/fzzy/gateway/hx2023/service/DeviceReportServiceImpl.java
@@ -1,8 +1,9 @@
package com.fzzy.gateway.hx2023.service;
import com.alibaba.fastjson2.JSONObject;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.fzzy.api.data.DepotType;
import com.fzzy.api.data.PushProtocol;
import com.fzzy.api.utils.NumberUtil;
import com.fzzy.async.fzzy40.entity.Fz40Grain;
import com.fzzy.data.ConfigData;
import com.fzzy.gateway.GatewayUtils;
@@ -15,7 +16,6 @@
import com.fzzy.gateway.hx2023.data.*;
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.apache.commons.lang.time.DateFormatUtils;
@@ -23,6 +23,7 @@
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
@Slf4j
@@ -65,6 +66,21 @@
        return new BaseResp();
    }
    @Override
    public BaseResp reportGrainDataByKafka(BaseReqData reqData) {
        String topic = ScConstant.TOPIC_ZLJ_GRAIN_TEMPERATURE;
        //如果是测试模式不执行推送
        if (configData.getActive().indexOf("dev") >= 0) {
            log.info("----------------------------推送KAFKA粮情信息,注:调试模式不推送---------------------------");
            log.info("-----TOPIC-----{}", topic);
            log.info("-----Message-----{}", reqData.getData());
            return new BaseResp();
        }
        kafkaDeviceReportService.publishWithTopic(reqData.getData(), topic);
        return new BaseResp();
    }
    @Override
    public BaseResp reportWeightData(BaseReqData reqData) {
@@ -92,7 +108,7 @@
            weightInfo.setNetWeight(reqData.getWeight());
            weightInfo.setWeightUnit("KG");
            JSONObject jsonObject = new JSONObject();
            jsonObject.put("weightInfo", JSONObject.toJSONString(weightInfo));
            jsonObject.put("weightInfo", JSON.toJSONString(weightInfo));
            packet.setProperties(jsonObject);
@@ -170,6 +186,21 @@
        BaseResp resp = new BaseResp();
        GrainCableData cableData = GatewayUtils.getCableData(device);
        if (null == device.getDepotType()) device.setDepotType(DepotType.TYPE_01.getCode());
        //表示筒仓
        if (DepotType.TYPE_02.getCode().equals(device.getDepotType()) || DepotType.TYPE_04.getCode().equals(device.getDepotType())) {
            return grainData2GatewayApiInfo2(grainData, device, cableData);
        }
        //表示为筒仓包括油罐仓
        if (DepotType.TYPE_03.getCode().equals(device.getDepotType())) {
            return grainData2GatewayApiInfo3(grainData, device, cableData);
        }
        KafaGrainData result = new KafaGrainData();
        result.setMessageId(ScConstant.getMessageId());
@@ -179,13 +210,6 @@
        result.setMinTemperature(grainData.getTempMin() + "");
        result.setMaxTemperature(grainData.getTempMax() + "");
        result.setCollectTime(DateFormatUtils.format(grainData.getReceiveDate(), "yyyy-MM-dd HH:mm:ss"));
        GrainCableData cableData = GatewayUtils.getCableData(device);
        if(cableData.isCir()){
            return grainData2GatewayApiInfo2(grainData, device,cableData);
        }
        //层行列
        int cableZ = cableData.getCableZ();
@@ -254,7 +278,7 @@
        temperatureAndhumidity.add(grainTH);
        trhInfo.put("temperatureAndhumidity", temperatureAndhumidity);
        //trhInfo.put("temperatureAndhumidity",grainTH);
        JSONObject params = new JSONObject();
        params.put("TRHInfo", trhInfo);
@@ -262,114 +286,434 @@
        result.setParams(params);
        //resp.setObj(result);
        resp.setData(JSONObject.toJSONString(result));
        return resp;
    }
    /**
     * 筒仓
     *
     * @param grainData
     * @param device
     * @param cableData
     * @return
     */
    private BaseResp grainData2GatewayApiInfo2(Fz40Grain grainData, GatewayDevice device, GrainCableData cableData) {
        BaseResp resp = new BaseResp();
        KafaGrainData result = new KafaGrainData();
        result.setMessageId(ScConstant.getMessageId());
        result.setMessgeId(result.getMessageId());
        result.setDeviceID(device.getDeviceId());
        result.setAvgTemperature(grainData.getTempAve() + "");
        result.setMinTemperature(grainData.getTempMin() + "");
        result.setMaxTemperature(grainData.getTempMax() + "");
        result.setCollectTime(DateFormatUtils.format(grainData.getReceiveDate(), "yyyy-MM-dd HH:mm:ss"));
        //层行列
        int cableZ = cableData.getCableZ();
        //温度集合
        String[] attr = grainData.getPoints().split(",");
        //根号
        int cableNum = 1, position = 0;
        String curTemp;
        List<KafkaGrainDataDetail1> temperature = new ArrayList<>();
        JSONObject totalCircle = new JSONObject();
        totalCircle.put("totalCircle", cableData.getTotalCircle() + "");
        totalCircle.put("smallCircle", cableData.getSmallCircle());
        String totalCircleStr = totalCircle.toJSONString();
        int x = 0, y = 0, z = 0, circle = 1;
        for (int i = 0; i < attr.length; i++) {
            position = i;
            z = i % cableZ + 1;
            //根号
            cableNum = (i / cableZ) + 1;
            curTemp = attr[i];
            circle = this.getCirCle(position, cableNum, cableData);
            y = 1;
            //如果是异常值,执行调整数据 TODO
            if (Double.valueOf(curTemp) < -99.9) {
                curTemp = grainData.getTempAve() + "";
            } else {
                //判断最大
                if (curTemp.equals(result.getMaxTemperature())) {
                    result.setMaxX(cableNum + "");
                    result.setMaxY(y + "");
                    result.setMaxZ(position + "");
                }
                //判断最小
                if (curTemp.equals(result.getMinTemperature())) {
                    result.setMinX(cableNum + "");
                    result.setMinY(y + "");
                    result.setMinZ(position + "");
                }
            }
            temperature.add(new KafkaGrainDataDetail1(cableNum + "", z + "", curTemp, position + "", null, null, circle + "", totalCircleStr));
        }
        //粮温信息
        JSONObject trhInfo = new JSONObject();
        trhInfo.put("temperature", temperature);
        //仓温度信息
        KafkaGrainTH grainTH = new KafkaGrainTH();
        grainTH.setHumidity(grainData.getHumidityIn() + "");
        grainTH.setTemperature(grainData.getTempIn() + "");
        grainTH.setAirHumidity(grainData.getHumidityOut() + "");
        grainTH.setAirTemperature(grainData.getTempOut() + "");
        List<KafkaGrainTH> temperatureAndhumidity = new ArrayList<>();
        temperatureAndhumidity.add(grainTH);
        trhInfo.put("temperatureAndhumidity", temperatureAndhumidity);
        JSONObject params = new JSONObject();
        params.put("TRHInfo", trhInfo);
        result.setParams(params);
        resp.setData(JSONObject.toJSONString(result));
        return resp;
    }
    /**
     * 获取当点所在圈
     *
     * @param cableNum
     * @param cableData
     * @return
     */
    private int getCirCle(int position, int cableNum, GrainCableData cableData) {
        int num1 = 1, num2 = 2;
        String[] attCable = cableData.getCableRule().split("-");
        if (cableData.getTotalCircle() == 1) return 1;
        if (cableData.getTotalCircle() == 2) {
            num1 = Integer.valueOf(attCable[0]);
            if (cableNum <= num1) return 1;
            return 2;
        }
        if (cableData.getTotalCircle() == 3) {
            num1 = Integer.valueOf(attCable[0]);
            num2 = num1 + Integer.valueOf(attCable[1]);
            if (cableNum <= num1) return 1;
            if (cableNum <= num2) return 2;
            return 3;
        }
        if (cableData.getTotalCircle() == 4) {
            num1 = Integer.valueOf(attCable[0]);
            num2 = num1 + Integer.valueOf(attCable[1]);
            if (cableNum <= num1) return 1;
            if (cableNum <= num2) return 2;
            num2 = num1 + Integer.valueOf(attCable[1]) + Integer.valueOf(attCable[2]);
            if (cableNum <= num2) return 3;
            return 4;
        }
        return 1;
    }
    /**
     * 油罐仓的处理
     * <p>
     * 2024年1月20日 暂时用的平房仓报文
     *
     * @param grainData
     * @param device
     * @param cableData
     * @return
     */
    private BaseResp grainData2GatewayApiInfo3(Fz40Grain grainData, GatewayDevice device, GrainCableData cableData) {
        BaseResp resp = new BaseResp();
        KafaGrainData result = new KafaGrainData();
        result.setMessageId(ScConstant.getMessageId());
        result.setMessgeId(result.getMessageId());
        result.setDeviceID(device.getDeviceId());
        result.setAvgTemperature(grainData.getTempAve() + "");
        result.setMinTemperature(grainData.getTempMin() + "");
        result.setMaxTemperature(grainData.getTempMax() + "");
        result.setCollectTime(DateFormatUtils.format(grainData.getReceiveDate(), "yyyy-MM-dd HH:mm:ss"));
        //层行列
        int cableZ = cableData.getCableZ();
        int cableY = cableData.getCableY();
        int cableX = cableData.getCableX();
        //温度集合
        String[] attr = grainData.getPoints().split(",");
        //根号
        int cableNum = 1, position = 0;
        String curTemp;
        List<KafkaGrainDataDetail1> temperature = new ArrayList<>();
        int x = 0, y = 0, z = 0;
        for (int i = 0; i < attr.length; i++) {
            position = i;
            z = i % cableZ + 1;
            x = i / (cableZ * cableY);
            y = x * (cableZ * cableY);
            y = (i - y) / cableZ;
            // 倒转X轴
            x = cableX - 1 - x;
            //根号
            cableNum = (i / cableZ) + 1;
            curTemp = attr[i];
            //如果是异常值,执行调整数据 TODO
            if (Double.valueOf(curTemp) < -99.9) {
                curTemp = grainData.getTempAve() + "";
            } else {
                //判断最大
                if (curTemp.equals(result.getMaxTemperature())) {
                    result.setMaxX(x + "");
                    result.setMaxY(y + "");
                    result.setMaxZ(position + "");
                }
                //判断最小
                if (curTemp.equals(result.getMinTemperature())) {
                    result.setMinX(x + "");
                    result.setMinY(y + "");
                    result.setMinZ(position + "");
                }
            }
            temperature.add(new KafkaGrainDataDetail1(cableNum + "", z + "", curTemp, position + "", x + "", y + ""));
        }
        //粮温信息
        JSONObject trhInfo = new JSONObject();
        // TRHInfo trhInfo = new TRHInfo();
        trhInfo.put("temperature", temperature);
        //仓温度信息
        KafkaGrainTH grainTH = new KafkaGrainTH();
        grainTH.setHumidity(grainData.getHumidityIn() + "");
        grainTH.setTemperature(grainData.getTempIn() + "");
        grainTH.setAirHumidity(grainData.getHumidityOut() + "");
        grainTH.setAirTemperature(grainData.getTempOut() + "");
        List<KafkaGrainTH> temperatureAndhumidity = new ArrayList<>();
        temperatureAndhumidity.add(grainTH);
        trhInfo.put("temperatureAndhumidity", temperatureAndhumidity);
        JSONObject params = new JSONObject();
        params.put("TRHInfo", trhInfo);
        result.setParams(params);
        resp.setData(JSONObject.toJSONString(result));
        return resp;
    }
    private BaseResp grainData2GatewayApiInfo2(Fz40Grain grainData, GatewayDevice device,GrainCableData cableData) {
    //  ----------------------------------------------------
    @Override
    public BaseResp grainData2GatewayApiInfoKafka(GrainData grainData, GatewayDevice device) {
        BaseResp resp = new BaseResp();
//        int cableZ = cableData.getCableZ();
//        int cableY = cableData.getCableY();
//        int cableX = cableData.getCableX();
//
//        int sumNum = cableData.getSumNum();
//
//        //数据封装
//        GrainData grain = new GrainData();
//        grain.setMessageId(ScConstant.getMessageId());
//        grain.setDeviceId(device.getDeviceId());
//        grain.setTimestamp(System.currentTimeMillis() + "");
//
//        ClientHeaders headers = new ClientHeaders();
//        headers.setDeviceName(device.getDeviceName());
//        headers.setProductId(device.getProductId());
//        headers.setOrgId(device.getOrgId());
//        headers.setMsgId(ScConstant.getMessageId());
//        grain.setHeaders(headers);
//
//
//        GrainOutPut outPut = new GrainOutPut();
//
//
//        double max = com.fzzy.protocol.bhzn.v0.cmd.ReMessageBuilder.MAX_TEMP, min = com.fzzy.protocol.bhzn.v0.cmd.ReMessageBuilder.MIN_TEMP, sumT = 0.0;
//
//        List<GrainTemp> temperature = new ArrayList<>();
//        //根号
//        int cableNum = 1, position = 0;
//
//        double curTemp;
//        int x = 0, y = 0, z = 0;
//        for (int i = 0; i < sumNum; i++) {
//            curTemp = temps.get(i);
//            position = i;
//
//            z = i % cableZ + 1;
//            x = i / (cableZ * cableY);
//            y = x * (cableZ * cableY);
//            y = (i - y) / cableZ;
//            //根号
//            cableNum = (i / cableZ) + 1;
//
//            temperature.add(new GrainTemp(cableNum + "", z + "", curTemp + "", position + ""));
//
//            //求最大最小值
//            if (curTemp < -900) {
//                sumNum--;
//            } else {
//                sumT += curTemp;
//                if (curTemp > max) {
//                    max = curTemp;
//                }
//                if (curTemp < min) {
//                    min = curTemp;
//                }
//            }
        GrainCableData cableData = GatewayUtils.getCableData(device);
        if (null == device.getDepotType()) device.setDepotType(DepotType.TYPE_01.getCode());
//        //表示筒仓
//        if (DepotType.TYPE_02.getCode().equals(device.getDepotType()) || DepotType.TYPE_04.getCode().equals(device.getDepotType())) {
//            return grainData2GatewayApiInfo2(grainData, device, cableData);
//        }
//
//        if (sumNum == 0) {
//            sumNum = 1;
//            log.warn("---当前粮情采集异常--");
//        //表示为筒仓包括油罐仓
//        if (DepotType.TYPE_03.getCode().equals(device.getDepotType())) {
//            return grainData2GatewayApiInfo3(grainData, device, cableData);
//        }
//        //过滤比较用的最大最小值
//        if (max == com.fzzy.protocol.bhzn.v0.cmd.ReMessageBuilder.MAX_TEMP) {
//            max = 0.0;
//        }
//        if (min == com.fzzy.protocol.bhzn.v0.cmd.ReMessageBuilder.MIN_TEMP) {
//            min = 0.0;
//        }
//
//        outPut.setTemperature(temperature);
//        outPut.setAvgTemperature(NumberUtil.keepPrecision((sumT / sumNum), 1) + "");
//        outPut.setMinTemperature(min + "");
//        outPut.setMaxTemperature(min + "");
//
//
//        com.alibaba.fastjson.JSONObject properties = new com.alibaba.fastjson.JSONObject();
//        properties.put("data", com.alibaba.fastjson.JSONObject.toJSONString(outPut));
//        properties.put("timestamp", grain.getTimestamp());
//
//        String height = this.getCacheHeight(device);
//        if (org.apache.commons.lang3.StringUtils.isEmpty(height)) height = "0.0";
//        properties.put("liquidHeight", height);
//
//        grain.setProperties(properties.toJSONString());
//
//        //封装好的数据
//        log.info("---浅圆仓封装完成----开始执行推送");
        GrainOutPut output = JSONObject.parseObject(grainData.getOutput(),GrainOutPut.class);
        KafaGrainData result = new KafaGrainData();
        result.setMessageId(ScConstant.getMessageId());
        result.setMessgeId(result.getMessageId());
        result.setDeviceID(device.getDeviceId());
        result.setAvgTemperature(output.getAvgTemperature());
        result.setMinTemperature(output.getMinTemperature());
        result.setMaxTemperature(output.getMaxTemperature() );
        result.setCollectTime(DateFormatUtils.format(new Date(), "yyyy-MM-dd HH:mm:ss"));
        //层行列
        int cableZ = cableData.getCableZ();
        int cableY = cableData.getCableY();
        int cableX = cableData.getCableX();
        //温度集合
        List<GrainTemp> attr = output.getTemperature();
        //根号
        int cableNum = 1, position = 0;
        String curTemp;
        List<KafkaGrainDataDetail1> temperature = new ArrayList<>();
        int x = 0, y = 0, z = 0;
        for (int i = 0; i < attr.size(); i++) {
            position = i;
            z = i % cableZ + 1;
            x = i / (cableZ * cableY);
            y = x * (cableZ * cableY);
            y = (i - y) / cableZ;
            // 倒转X轴
            x = cableX - 1 - x;
            //根号
            cableNum = (i / cableZ) + 1;
            curTemp = attr.get(i).getTemperature();
            //如果是异常值,执行调整数据 TODO
            if (Double.valueOf(curTemp) < -99.9) {
                curTemp = output.getAvgTemperature();
            } else {
                //判断最大
                if (curTemp.equals(result.getMaxTemperature())) {
                    result.setMaxX(x + "");
                    result.setMaxY(y + "");
                    result.setMaxZ(position + "");
                }
                //判断最小
                if (curTemp.equals(result.getMinTemperature())) {
                    result.setMinX(x + "");
                    result.setMinY(y + "");
                    result.setMinZ(position + "");
                }
            }
            temperature.add(new KafkaGrainDataDetail1(cableNum + "", z + "", curTemp, position + "", x + "", y + ""));
        }
        //粮温信息
        JSONObject trhInfo = new JSONObject();
        // TRHInfo trhInfo = new TRHInfo();
        trhInfo.put("temperature", temperature);
        //仓温度信息
        KafkaGrainTH grainTH = new KafkaGrainTH();
        List<GrainTH> ths= output.getTemperatureAndhumidity();
        grainTH.setHumidity(ths.get(0).getHumidity());
        grainTH.setTemperature(ths.get(0).getTemperature() );
        GrainWeather weather = JSON.parseObject(grainData.getWeatherStation(),GrainWeather.class);
        grainTH.setAirHumidity(weather.getHumidity());
        grainTH.setAirTemperature(weather.getTemperature() );
        List<KafkaGrainTH> temperatureAndhumidity = new ArrayList<>();
        temperatureAndhumidity.add(grainTH);
        trhInfo.put("temperatureAndhumidity", temperatureAndhumidity);
        JSONObject params = new JSONObject();
        params.put("TRHInfo", trhInfo);
        resp.setCode(BaseResp.CODE_500);
        resp.setMsg("筒仓解析暂未实现");
        result.setParams(params);
        resp.setData(JSONObject.toJSONString(result));
        return resp;
    }
    /**
     * 获取当点所在圈
     *
     * @param cableNum
     * @param cableData
     * @return
     */
    private int getCirCleKafka(int position, int cableNum, GrainCableData cableData) {
        int num1 = 1, num2 = 2;
        String[] attCable = cableData.getCableRule().split("-");
        if (cableData.getTotalCircle() == 1) return 1;
        if (cableData.getTotalCircle() == 2) {
            num1 = Integer.valueOf(attCable[0]);
            if (cableNum <= num1) return 1;
            return 2;
        }
        if (cableData.getTotalCircle() == 3) {
            num1 = Integer.valueOf(attCable[0]);
            num2 = num1 + Integer.valueOf(attCable[1]);
            if (cableNum <= num1) return 1;
            if (cableNum <= num2) return 2;
            return 3;
        }
        if (cableData.getTotalCircle() == 4) {
            num1 = Integer.valueOf(attCable[0]);
            num2 = num1 + Integer.valueOf(attCable[1]);
            if (cableNum <= num1) return 1;
            if (cableNum <= num2) return 2;
            num2 = num1 + Integer.valueOf(attCable[1]) + Integer.valueOf(attCable[2]);
            if (cableNum <= num2) return 3;
            return 4;
        }
        return 1;
    }
}