Pārlūkot izejas kodu

Merge branch 'vpp-zyj' into feature/service-vpp-20260701

james 1 dienu atpakaļ
vecāks
revīzija
b7e00022f6

+ 9 - 0
service-tsdb/service-tsdb-api/src/main/java/com/usky/demo/RemoteTsdbProxyService.java

@@ -6,7 +6,9 @@ import com.usky.demo.domain.*;
 import com.usky.demo.factory.RemoteTsdbProxyFallbackFactory;
 import com.usky.demo.feign.RemoteTsdbProxyFeignConfiguration;
 import org.springframework.cloud.openfeign.FeignClient;
+import org.springframework.http.MediaType;
 import org.springframework.web.bind.annotation.*;
+import org.springframework.web.multipart.MultipartFile;
 
 import java.util.List;
 import java.util.Map;
@@ -105,4 +107,11 @@ public interface RemoteTsdbProxyService {
      */
     @PostMapping("/energyItemTrend")
     EnergyItemTrendResultVO sumEnergyItemTrend(@RequestBody EnergyItemTrendQueryVO requestVO);
+
+    /**
+     * 解析 CSV 并导入设备时序数据:INSERT INTO _{deviceUuid} (ts, p) VALUES ...
+     */
+    @PostMapping(value = "/importDeviceCsv", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
+    ApiResult<DeviceCsvImportResultVO> importDeviceCsv(@RequestParam("deviceUuid") String deviceUuid,
+                                                       @RequestPart("file") MultipartFile file);
 }

+ 23 - 0
service-tsdb/service-tsdb-api/src/main/java/com/usky/demo/domain/DeviceCsvImportResultVO.java

@@ -0,0 +1,23 @@
+package com.usky.demo.domain;
+
+import lombok.Data;
+
+import java.io.Serializable;
+
+/**
+ * TDengine CSV 导入结果
+ */
+@Data
+public class DeviceCsvImportResultVO implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    /** 设备 UUID(入参) */
+    private String deviceUuid;
+
+    /** 实际写入的子表名(通常为 _{deviceUuid}) */
+    private String tableName;
+
+    /** 成功导入的数据行数 */
+    private Integer importedRows;
+}

+ 6 - 0
service-tsdb/service-tsdb-api/src/main/java/com/usky/demo/factory/RemoteTsdbProxyFallbackFactory.java

@@ -104,6 +104,12 @@ public class RemoteTsdbProxyFallbackFactory implements FallbackFactory<RemoteTsd
             {
                 throw new BusinessException("能耗分项趋势查询:" + throwable.getMessage());
             }
+
+            @Override
+            public ApiResult<DeviceCsvImportResultVO> importDeviceCsv(String deviceUuid, org.springframework.web.multipart.MultipartFile file)
+            {
+                throw new BusinessException("设备时序 CSV 导入:" + throwable.getMessage());
+            }
         };
     }
 }

+ 10 - 0
service-tsdb/service-tsdb-biz/src/main/java/com/usky/demo/controller/api/DataTsdbProxyControllerApi.java

@@ -16,6 +16,7 @@ import com.usky.system.RemoteUserService;
 import com.usky.system.domain.SysUserVO;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.beans.factory.annotation.Value;
+import org.springframework.http.MediaType;
 import org.springframework.web.bind.annotation.*;
 import org.springframework.web.multipart.MultipartFile;
 
@@ -175,4 +176,13 @@ public class DataTsdbProxyControllerApi implements RemoteTsdbProxyService {
         return queryTdengineDataService.sumEnergyItemTrend(requestVO);
     }
 
+    @Override
+    public ApiResult<DeviceCsvImportResultVO> importDeviceCsv(@RequestParam("deviceUuid") String deviceUuid,
+                                                              @RequestPart("file") MultipartFile file) {
+        if (!"taos".equals(sourcetype)) {
+            throw new BusinessException("当前数据源不支持 TDengine CSV 导入");
+        }
+        return ApiResult.success(queryTdengineDataService.importDeviceCsv(deviceUuid, file));
+    }
+
 }

+ 6 - 0
service-tsdb/service-tsdb-biz/src/main/java/com/usky/demo/service/QueryTdengineDataService.java

@@ -6,6 +6,7 @@ import com.usky.demo.service.vo.SuperTableVO;
 import org.apache.ibatis.annotations.Param;
 import org.springframework.web.bind.annotation.RequestBody;
 import org.springframework.web.bind.annotation.RequestParam;
+import org.springframework.web.multipart.MultipartFile;
 
 import java.util.List;
 import java.util.Map;
@@ -48,4 +49,9 @@ public interface QueryTdengineDataService extends CrudService<QueryTdengineData>
      * 能耗分项趋势:按时间粒度 INTERVAL 聚合,返回 time/value 列表
      */
     EnergyItemTrendResultVO sumEnergyItemTrend(EnergyItemTrendQueryVO requestVO);
+
+    /**
+     * 从 CSV 解析并导入设备子表时序:INSERT INTO _{deviceUuid} (ts, p) VALUES ...
+     */
+    DeviceCsvImportResultVO importDeviceCsv(String deviceUuid, MultipartFile file);
 }

+ 187 - 0
service-tsdb/service-tsdb-biz/src/main/java/com/usky/demo/service/impl/QueryTdengineDataServiceImpl.java

@@ -23,12 +23,21 @@ import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.cache.annotation.Cacheable;
 import org.springframework.jdbc.core.JdbcTemplate;
 import org.springframework.stereotype.Service;
+import org.springframework.web.multipart.MultipartFile;
 
 import javax.sql.DataSource;
+import java.io.BufferedReader;
+import java.io.IOException;
+import java.io.InputStreamReader;
 import java.math.BigDecimal;
+import java.nio.charset.StandardCharsets;
 import java.sql.*;
 import java.text.ParseException;
 import java.text.SimpleDateFormat;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.format.DateTimeFormatter;
+import java.time.format.DateTimeParseException;
 import java.util.*;
 import java.util.regex.Pattern;
 import java.util.stream.Collectors;
@@ -47,6 +56,16 @@ public class QueryTdengineDataServiceImpl extends AbstractCrudService<QueryTdeng
 
     private static final Pattern SQL_IDENTIFIER_PATTERN = Pattern.compile("^[a-zA-Z_][a-zA-Z0-9_]*$");
 
+    private static final Pattern DEVICE_UUID_PATTERN = Pattern.compile("^[a-zA-Z0-9_-]+$");
+
+    private static final int IMPORT_BATCH_SIZE = 500;
+
+    private static final DateTimeFormatter TS_FORMAT = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+
+    private static final List<String> TS_HEADER_ALIASES = Arrays.asList("ts", "timestamp", "time");
+
+    private static final List<String> P_HEADER_ALIASES = Arrays.asList("p", "power", "active_power", "value");
+
     @Autowired
     private TsdbUtils tsdbUtils;
     @Autowired
@@ -460,4 +479,172 @@ public class QueryTdengineDataServiceImpl extends AbstractCrudService<QueryTdeng
             throw new BusinessException(paramName + " 格式非法,仅允许字母、数字与下划线,且不能以数字开头");
         }
     }
+
+    @Override
+    public DeviceCsvImportResultVO importDeviceCsv(String deviceUuid, MultipartFile file) {
+        if (StringUtils.isBlank(deviceUuid)) {
+            throw new BusinessException("设备UUID不能为空");
+        }
+        if (!DEVICE_UUID_PATTERN.matcher(deviceUuid.trim()).matches()) {
+            throw new BusinessException("设备UUID格式无效");
+        }
+        if (file == null || file.isEmpty()) {
+            throw new BusinessException("导入文件不能为空");
+        }
+        String originalFilename = file.getOriginalFilename();
+        if (StringUtils.isBlank(originalFilename) || !originalFilename.toLowerCase(Locale.ROOT).endsWith(".csv")) {
+            throw new BusinessException("仅支持 CSV 文件");
+        }
+
+        String normalizedUuid = deviceUuid.trim();
+        String tableName = buildDeviceSubTableName(normalizedUuid);
+        try {
+            List<TsdbCsvRow> rows = parseCsvRows(file);
+            batchInsertRows(tableName, rows);
+
+            DeviceCsvImportResultVO result = new DeviceCsvImportResultVO();
+            result.setDeviceUuid(normalizedUuid);
+            result.setTableName(tableName);
+            result.setImportedRows(rows.size());
+            return result;
+        } catch (IOException ex) {
+            throw new BusinessException("解析 CSV 文件失败: " + ex.getMessage());
+        } catch (SQLException ex) {
+            throw new BusinessException("时序数据导入失败: " + ex.getMessage());
+        }
+    }
+
+    private List<TsdbCsvRow> parseCsvRows(MultipartFile file) throws IOException {
+        List<TsdbCsvRow> rows = new ArrayList<>();
+        try (BufferedReader reader = new BufferedReader(
+                new InputStreamReader(file.getInputStream(), StandardCharsets.UTF_8))) {
+            String line;
+            int tsIndex = 0;
+            int pIndex = 1;
+            boolean headerResolved = false;
+            int lineNo = 0;
+            while ((line = reader.readLine()) != null) {
+                lineNo++;
+                line = line.trim();
+                if (line.isEmpty()) {
+                    continue;
+                }
+                String[] parts = splitCsvLine(line);
+                if (!headerResolved) {
+                    if (isHeaderRow(parts)) {
+                        tsIndex = findColumnIndex(parts, TS_HEADER_ALIASES, "ts");
+                        pIndex = findColumnIndex(parts, P_HEADER_ALIASES, "p");
+                        headerResolved = true;
+                        continue;
+                    }
+                    headerResolved = true;
+                }
+                if (parts.length <= Math.max(tsIndex, pIndex)) {
+                    throw new BusinessException("CSV 第 " + lineNo + " 行列数不足");
+                }
+                rows.add(new TsdbCsvRow(
+                        parseTimestamp(parts[tsIndex], lineNo),
+                        parseMetricValue(parts[pIndex], lineNo)));
+            }
+        }
+        if (rows.isEmpty()) {
+            throw new BusinessException("CSV 无有效数据行");
+        }
+        return rows;
+    }
+
+    private void batchInsertRows(String tableName, List<TsdbCsvRow> rows) throws SQLException {
+        try (Connection connection = dataSource.getConnection();
+             Statement statement = connection.createStatement()) {
+            for (int start = 0; start < rows.size(); start += IMPORT_BATCH_SIZE) {
+                int end = Math.min(start + IMPORT_BATCH_SIZE, rows.size());
+                StringBuilder sql = new StringBuilder("INSERT INTO ")
+                        .append(tableName)
+                        .append(" (ts, p) VALUES ");
+                for (int i = start; i < end; i++) {
+                    TsdbCsvRow row = rows.get(i);
+                    if (i > start) {
+                        sql.append(' ');
+                    }
+                    sql.append('(').append(row.tsMillis).append(',').append(row.pValue).append(')');
+                }
+                statement.execute(sql.toString());
+            }
+        }
+    }
+
+    private static String[] splitCsvLine(String line) {
+        return Arrays.stream(line.split(",", -1))
+                .map(String::trim)
+                .toArray(String[]::new);
+    }
+
+    private static boolean isHeaderRow(String[] parts) {
+        for (String part : parts) {
+            String normalized = part.trim().toLowerCase(Locale.ROOT);
+            if (TS_HEADER_ALIASES.contains(normalized) || P_HEADER_ALIASES.contains(normalized)) {
+                return true;
+            }
+        }
+        return false;
+    }
+
+    private static int findColumnIndex(String[] headers, List<String> aliases, String columnLabel) {
+        for (int i = 0; i < headers.length; i++) {
+            if (aliases.contains(headers[i].trim().toLowerCase(Locale.ROOT))) {
+                return i;
+            }
+        }
+        throw new BusinessException("CSV 缺少 " + columnLabel + " 列");
+    }
+
+    private static long parseTimestamp(String rawValue, int lineNo) {
+        String value = rawValue == null ? "" : rawValue.trim();
+        if (value.isEmpty()) {
+            throw new BusinessException("CSV 第 " + lineNo + " 行 ts 为空");
+        }
+        if (value.matches("\\d+")) {
+            long epoch = Long.parseLong(value);
+            return value.length() <= 10 ? epoch * 1000L : epoch;
+        }
+        try {
+            return LocalDateTime.parse(value, TS_FORMAT)
+                    .atZone(ZoneId.systemDefault())
+                    .toInstant()
+                    .toEpochMilli();
+        } catch (DateTimeParseException ignored) {
+            // fall through
+        }
+        try {
+            return new SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss").parse(value).getTime();
+        } catch (ParseException ex) {
+            throw new BusinessException("CSV 第 " + lineNo + " 行 ts 格式无效: " + value);
+        }
+    }
+
+    private static double parseMetricValue(String rawValue, int lineNo) {
+        String value = rawValue == null ? "" : rawValue.trim();
+        if (value.isEmpty()) {
+            throw new BusinessException("CSV 第 " + lineNo + " 行 p 为空");
+        }
+        try {
+            return new BigDecimal(value).doubleValue();
+        } catch (NumberFormatException ex) {
+            throw new BusinessException("CSV 第 " + lineNo + " 行 p 不是有效数值: " + value);
+        }
+    }
+
+    private static String buildDeviceSubTableName(String deviceUuid) {
+        return "_" + deviceUuid;
+    }
+
+    private static final class TsdbCsvRow {
+        private final long tsMillis;
+        private final double pValue;
+
+        private TsdbCsvRow(long tsMillis, double pValue) {
+            this.tsMillis = tsMillis;
+            this.pValue = pValue;
+        }
+    }
 }

+ 16 - 0
service-vpp/service-vpp-biz/src/main/java/com/usky/vpp/controller/web/VppDeviceController.java

@@ -5,8 +5,11 @@ import com.usky.common.core.bean.CommonPage;
 import com.usky.vpp.domain.VppDevice;
 import com.usky.vpp.service.VppDeviceService;
 import com.usky.vpp.service.vo.DeviceRequest;
+import com.usky.vpp.service.vo.DeviceTsdbImportResultVO;
 import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.http.MediaType;
 import org.springframework.web.bind.annotation.*;
+import org.springframework.web.multipart.MultipartFile;
 
 import java.util.Map;
 
@@ -47,4 +50,17 @@ public class VppDeviceController {
         vppDeviceService.deleteDevice(id);
         return ApiResult.success();
     }
+
+    /**
+     * 解析 CSV 并导入设备时序数据到 TDengine。
+     * <p>CSV 需包含 ts(时间戳)与 p(数值)列,支持表头;时间支持毫秒/秒时间戳或 yyyy-MM-dd HH:mm:ss。</p>
+     *
+     * @param deviceUuid 设备 UUID,对应子表 _{deviceUuid}
+     * @param file       CSV 文件
+     */
+    @PostMapping(value = "/import-tsdb", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
+    public ApiResult<DeviceTsdbImportResultVO> importTsdbData(@RequestParam("deviceUuid") String deviceUuid,
+                                                             @RequestPart("file") MultipartFile file) {
+        return ApiResult.success(vppDeviceService.importTsdbData(deviceUuid, file));
+    }
 }

+ 4 - 0
service-vpp/service-vpp-biz/src/main/java/com/usky/vpp/service/VppDeviceService.java

@@ -3,6 +3,8 @@ package com.usky.vpp.service;
 import com.usky.common.core.bean.CommonPage;
 import com.usky.vpp.domain.VppDevice;
 import com.usky.vpp.service.vo.DeviceRequest;
+import com.usky.vpp.service.vo.DeviceTsdbImportResultVO;
+import org.springframework.web.multipart.MultipartFile;
 
 import java.util.Map;
 
@@ -24,4 +26,6 @@ public interface VppDeviceService {
     void updateDevice(Long id, DeviceRequest request);
 
     void deleteDevice(Long id);
+
+    DeviceTsdbImportResultVO importTsdbData(String deviceUuid, MultipartFile file);
 }

+ 26 - 0
service-vpp/service-vpp-biz/src/main/java/com/usky/vpp/service/impl/VppDeviceServiceImpl.java

@@ -2,10 +2,13 @@ package com.usky.vpp.service.impl;
 
 import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
 import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import com.usky.common.core.bean.ApiResult;
 import com.usky.common.core.bean.CommonPage;
 import com.usky.common.core.exception.BusinessException;
 import com.usky.common.core.util.UUIDUtils;
 import com.usky.common.security.utils.SecurityUtils;
+import com.usky.demo.RemoteTsdbProxyService;
+import com.usky.demo.domain.DeviceCsvImportResultVO;
 import com.usky.iot.RemoteIotTaskService;
 import com.usky.vpp.domain.DmpDeviceStatus;
 import com.usky.vpp.domain.VppDevice;
@@ -14,6 +17,7 @@ import com.usky.vpp.mapper.VppDeviceMapper;
 import com.usky.vpp.service.VppDeviceService;
 import com.usky.vpp.service.VppSiteService;
 import com.usky.vpp.service.vo.DeviceRequest;
+import com.usky.vpp.service.vo.DeviceTsdbImportResultVO;
 import com.usky.vpp.util.VppAuditHelper;
 import com.usky.vpp.util.VppPageHelper;
 import org.springframework.beans.factory.annotation.Autowired;
@@ -21,6 +25,7 @@ import org.springframework.stereotype.Service;
 import org.springframework.transaction.annotation.Transactional;
 import org.springframework.util.CollectionUtils;
 import org.springframework.util.StringUtils;
+import org.springframework.web.multipart.MultipartFile;
 
 import java.util.List;
 import java.util.Map;
@@ -37,6 +42,8 @@ public class VppDeviceServiceImpl implements VppDeviceService {
     private VppSiteService siteService;
     @Autowired
     private RemoteIotTaskService remoteIotTaskService;
+    @Autowired
+    private RemoteTsdbProxyService remoteTsdbProxyService;
 
     @Override
     public CommonPage<VppDevice> pageDevice(Map<String, Object> params) {
@@ -123,6 +130,25 @@ public class VppDeviceServiceImpl implements VppDeviceService {
         remoteIotTaskService.deleteDeviceInfo(device.getDeviceUuid());
     }
 
+    @Override
+    public DeviceTsdbImportResultVO importTsdbData(String deviceUuid, MultipartFile file) {
+        getByDeviceUuid(deviceUuid);
+        if (file == null || file.isEmpty()) {
+            throw new BusinessException("导入文件不能为空");
+        }
+        ApiResult<DeviceCsvImportResultVO> result = remoteTsdbProxyService.importDeviceCsv(deviceUuid, file);
+        if (result == null || !result.isSuccess() || result.getData() == null) {
+            throw new BusinessException(result != null && StringUtils.hasText(result.getMsg())
+                    ? result.getMsg() : "时序数据导入失败");
+        }
+        DeviceCsvImportResultVO data = result.getData();
+        DeviceTsdbImportResultVO vo = new DeviceTsdbImportResultVO();
+        vo.setDeviceUuid(data.getDeviceUuid());
+        vo.setTableName(data.getTableName());
+        vo.setImportedRows(data.getImportedRows());
+        return vo;
+    }
+
     private LambdaQueryWrapper<VppDevice> buildQueryWrapper(Map<String, Object> params) {
         LambdaQueryWrapper<VppDevice> wrapper = new LambdaQueryWrapper<VppDevice>()
                 .eq(VppDevice::getDeleteFlag, VppAuditHelper.NOT_DELETED)

+ 19 - 0
service-vpp/service-vpp-biz/src/main/java/com/usky/vpp/service/vo/DeviceTsdbImportResultVO.java

@@ -0,0 +1,19 @@
+package com.usky.vpp.service.vo;
+
+import lombok.Data;
+
+/**
+ * TDengine CSV 导入结果
+ */
+@Data
+public class DeviceTsdbImportResultVO {
+
+    /** 设备 UUID */
+    private String deviceUuid;
+
+    /** TDengine 子表名(_{deviceUuid}) */
+    private String tableName;
+
+    /** 成功导入的数据行数 */
+    private Integer importedRows;
+}