package cn.iocoder.yudao.module.mes.service.dv.telemetry; import cn.hutool.core.collection.CollUtil; import cn.hutool.core.date.DateUtil; import cn.hutool.core.util.StrUtil; import cn.iocoder.yudao.framework.common.pojo.PageResult; import cn.iocoder.yudao.framework.common.util.http.HttpUtils; import cn.iocoder.yudao.framework.common.util.json.JsonUtils; import cn.iocoder.yudao.module.mes.controller.admin.dv.telemetry.vo.MesDvTelemetryDeviceRespVO; import cn.iocoder.yudao.module.mes.controller.admin.dv.telemetry.vo.MesDvTelemetryLatestReqVO; import cn.iocoder.yudao.module.mes.controller.admin.dv.telemetry.vo.MesDvTelemetryPageReqVO; import cn.iocoder.yudao.module.mes.controller.admin.dv.telemetry.vo.MesDvTelemetryPullRespVO; import cn.iocoder.yudao.module.mes.dal.dataobject.dv.telemetry.MesDvTelemetryDO; import cn.iocoder.yudao.module.mes.dal.mysql.dv.telemetry.MesDvTelemetryMapper; import cn.iocoder.yudao.module.mes.framework.config.MesTelemetryProperties; import jakarta.annotation.Resource; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import java.time.LocalDateTime; import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import static cn.iocoder.yudao.framework.common.exception.util.ServiceExceptionUtil.exception; import static cn.iocoder.yudao.module.mes.enums.ErrorCodeConstants.DV_TELEMETRY_PULL_ERROR; /** * MES 设备数采遥测数据 Service 实现类 * * @author 超级管理员 */ @Service @Slf4j public class MesDvTelemetryServiceImpl implements MesDvTelemetryService { @Resource private MesDvTelemetryMapper telemetryMapper; @Resource private MesTelemetryProperties telemetryProperties; @Override public PageResult getTelemetryPage(MesDvTelemetryPageReqVO pageReqVO) { return telemetryMapper.selectPage(pageReqVO); } @Override public List getLatestTelemetry(MesDvTelemetryLatestReqVO reqVO) { // 取最近一批数据,按设备+信号去重保留最新,用于实时监控视图 List list = telemetryMapper.selectLatest(reqVO, 5000); Map latestMap = new LinkedHashMap<>(); for (MesDvTelemetryDO row : list) { String key = StrUtil.blankToDefault(row.getTbDeviceId(), "") + ":" + StrUtil.blankToDefault(row.getParamKeyName(), ""); latestMap.putIfAbsent(key, row); } return new ArrayList<>(latestMap.values()); } @Override public List getDeviceList() { List result = new ArrayList<>(); for (MesDvTelemetryDO row : telemetryMapper.selectDeviceList()) { MesDvTelemetryDeviceRespVO vo = new MesDvTelemetryDeviceRespVO(); vo.setTbDeviceId(row.getTbDeviceId()); vo.setDeviceName(row.getDeviceName()); result.add(vo); } return result; } @Override public MesDvTelemetryPullRespVO pullTelemetry() { // 1. 校验配置 if (StrUtil.isBlank(telemetryProperties.getBaseUrl())) { throw exception(DV_TELEMETRY_PULL_ERROR, "未配置数采服务地址 yudao.mes.telemetry.base-url"); } String url = telemetryProperties.getBaseUrl() + telemetryProperties.getPath(); // 2. 调用外部数采接口 String responseText; try { responseText = HttpUtils.get(url, null); } catch (Exception e) { log.error("[pullTelemetry][调用数采接口失败,url({})]", url, e); throw exception(DV_TELEMETRY_PULL_ERROR, e.getMessage()); } if (StrUtil.isBlank(responseText)) { throw exception(DV_TELEMETRY_PULL_ERROR, "数采接口返回为空"); } // 3. 解析响应 Map response = JsonUtils.parseObject(responseText, Map.class); if (response == null || !"200".equals(String.valueOf(response.get("code")))) { log.error("[pullTelemetry][数采接口返回异常,url({}),响应({})]", url, responseText); throw exception(DV_TELEMETRY_PULL_ERROR, "数采接口返回异常"); } List dataList = (List) response.get("data"); // 4. 无数据时返回空结果,不视为错误 MesDvTelemetryPullRespVO result = new MesDvTelemetryPullRespVO(); LocalDateTime pullTime = LocalDateTime.now(); result.setPullTime(pullTime); if (CollUtil.isEmpty(dataList)) { result.setRecordCount(0); result.setDeviceCount(0); return result; } // 5. 转换并批量入库 List saveList = new ArrayList<>(dataList.size()); for (Object rowObj : dataList) { if (!(rowObj instanceof Map)) { continue; } MesDvTelemetryDO data = convert((Map) rowObj); data.setPullTime(pullTime); saveList.add(data); } telemetryMapper.insertBatch(saveList); // 6. 返回统计 result.setRecordCount(saveList.size()); result.setDeviceCount((int) saveList.stream() .map(MesDvTelemetryDO::getTbDeviceId).distinct().count()); return result; } private MesDvTelemetryDO convert(Map row) { MesDvTelemetryDO data = new MesDvTelemetryDO(); data.setTbDeviceId(StrUtil.toStringOrNull(row.get("tbDeviceId"))); data.setDeviceName(StrUtil.toStringOrNull(row.get("deviceName"))); data.setParamName(StrUtil.toStringOrNull(row.get("paramName"))); data.setParamKeyName(StrUtil.toStringOrNull(row.get("paramKeyName"))); data.setStandardValue(StrUtil.toStringOrNull(row.get("standardValue"))); data.setTimelyValue(StrUtil.toStringOrNull(row.get("timelyValue"))); data.setAvgValue(StrUtil.toStringOrNull(row.get("avgValue"))); data.setMaxValue(StrUtil.toStringOrNull(row.get("maxValue"))); data.setMinValue(StrUtil.toStringOrNull(row.get("minValue"))); data.setWhetherAnomaly(parseBoolean(row.get("whetherAnomaly"))); data.setTelemetryDataTime(parseTime(row.get("telemetryDataTime"))); data.setBillNo(StrUtil.toStringOrNull(row.get("billNo"))); data.setShiftName(StrUtil.toStringOrNull(row.get("shiftName"))); return data; } private Boolean parseBoolean(Object value) { if (value == null) { return null; } String str = String.valueOf(value).trim(); if ("1".equals(str) || "true".equalsIgnoreCase(str) || "Y".equalsIgnoreCase(str)) { return Boolean.TRUE; } if ("0".equals(str) || "false".equalsIgnoreCase(str) || "N".equalsIgnoreCase(str)) { return Boolean.FALSE; } return null; } private LocalDateTime parseTime(Object value) { if (value == null || StrUtil.isBlank(String.valueOf(value))) { return null; } try { return DateUtil.parse(String.valueOf(value)).toLocalDateTime(); } catch (Exception e) { log.warn("[parseTime][遥测时间解析失败,value({})]", value); return null; } } }