| | |
| | | package com.trafficaudit.reportexport.service; |
| | | |
| | | import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; |
| | | import com.trafficaudit.dataimport.entity.ImportBatch; |
| | | import com.trafficaudit.dataimport.mapper.ImportBatchMapper; |
| | | import lombok.extern.slf4j.Slf4j; |
| | | import org.springframework.stereotype.Service; |
| | | |
| | | import javax.annotation.PreDestroy; |
| | | import javax.annotation.Resource; |
| | | import javax.sql.DataSource; |
| | | import java.io.File; |
| | | import java.nio.charset.StandardCharsets; |
| | | import java.security.MessageDigest; |
| | | import java.sql.Connection; |
| | | import java.sql.ResultSet; |
| | | import java.sql.Statement; |
| | | import java.util.ArrayList; |
| | | import java.util.Arrays; |
| | | import java.util.Comparator; |
| | | import java.util.Date; |
| | | import java.util.HashSet; |
| | | import java.util.Iterator; |
| | | import java.util.LinkedHashMap; |
| | | import java.util.List; |
| | | import java.util.Map; |
| | | import java.util.Set; |
| | | import java.util.UUID; |
| | | import java.util.concurrent.ConcurrentHashMap; |
| | | import java.util.concurrent.ExecutorService; |
| | |
| | | * 2) 生成过程不占用 HTTP 请求线程,避免多人同时点把容器拖死(工作线程固定 1 个,多余请求显示「排队中」); |
| | | * 3) 同一报表期的结果可复用(**落盘保存,进程重启后依然有效**),避免每次重算整本约 4 万格公式。 |
| | | * |
| | | * 缓存正确性守卫:只有当「期间没有任何新的数据导入(import_batch 最大 id 未变)」且「结果不超过 7 天」时才复用; |
| | | * 导入即失效,避免把过期报表给用户。前端仍提供「重新生成」强制重算(跳过缓存)。 |
| | | * 落盘位置:{summary.template-dir}/_cache/<key>_v<数据版本>.xlsx |
| | | * 缓存正确性守卫=**数据指纹**:指纹由「汇总大表依赖的源表内容校验和」+「母版等文件身份(大小/修改时间)」算出, |
| | | * 只有指纹一致(且结果不超过 7 天)才复用。因此: |
| | | * - 投资/能耗/审核等**与汇总大表无关**的导入不再让缓存失效; |
| | | * - 换母版、改参照件、动相关源表(含原地 UPDATE)都会让缓存失效,不会把旧表当新表发出。 |
| | | * 前端仍提供「重新生成」强制重算(跳过缓存)。落盘位置:{summary.template-dir}/_cache/<key>_v<指纹>.xlsx |
| | | */ |
| | | @Slf4j |
| | | @Service |
| | |
| | | @Resource |
| | | private ReportExportService reportService; |
| | | @Resource |
| | | private ImportBatchMapper importBatchMapper; |
| | | private DataSource dataSource; |
| | | |
| | | /** 汇总工作簿目录(与 ReportExportService 同配置),落盘缓存放在其 _cache 子目录 */ |
| | | @org.springframework.beans.factory.annotation.Value("${summary.template-dir:docs/生成汇总大表}") |
| | | private String summaryTemplateDir; |
| | | |
| | | /** 城市客运目录:网约车订单及全省总量.xlsx(库里缺该月时的兜底输入) */ |
| | | @org.springframework.beans.factory.annotation.Value("${city-passenger.template-dir:docs/城市客运}") |
| | | private String cityPassengerTemplateDir; |
| | | |
| | | /** 计算指纹时**排除**的表:投资线、能耗线、审核线、系统表、导入日志。 |
| | | * 这些表变化不影响汇总大表内容,没必要让缓存失效;其余表一律纳入(宁可多算不可出错)。 */ |
| | | private static final Set<String> FINGERPRINT_EXCLUDE = new HashSet<>(Arrays.asList( |
| | | "investment_project", "investment_monthly", "investment_system", |
| | | "h204_vehicle_quarterly", "h204_auth_vehicle", |
| | | "audit_result", "audit_run", "audit_rule", "audit_explanation", |
| | | "sys_user", "sys_role", "sys_user_role", "sys_dept", "sys_operation_log", |
| | | "llm_desensitize_map", "import_batch")); |
| | | |
| | | /** 落盘缓存目录;定位策略与模板一致(配置目录 / user.dir 相对 / 上级目录) */ |
| | | private File cacheDir() { |
| | |
| | | private final String mode; |
| | | private final Integer toMonth; |
| | | private volatile long createdAt = System.currentTimeMillis(); |
| | | /** 生成时的数据版本(import_batch 最大 id),用于判断缓存是否过期 */ |
| | | private final long importVersion; |
| | | /** 生成时的数据指纹(源表校验和 + 母版文件身份),用于判断缓存是否过期 */ |
| | | private final String fingerprint; |
| | | private volatile String status = "PENDING"; // PENDING/RUNNING/DONE/FAILED |
| | | private volatile int percent = 0; |
| | | private volatile String step = "排队中"; |
| | |
| | | private volatile File cacheFile; |
| | | private volatile String fileName; |
| | | private volatile boolean fromCache; |
| | | /** 本次「没命中缓存」是因为存在旧结果但数据/母版已变(用于给用户一句解释) */ |
| | | private volatile boolean staleCache; |
| | | private volatile long finishedAt; |
| | | |
| | | Task(String type, String period, String mode, Integer toMonth, long importVersion) { |
| | | Task(String type, String period, String mode, Integer toMonth, String fingerprint) { |
| | | this.type = type; |
| | | this.period = period; |
| | | this.mode = mode; |
| | | this.toMonth = toMonth; |
| | | this.importVersion = importVersion; |
| | | this.fingerprint = fingerprint; |
| | | } |
| | | |
| | | void progress(int pct, String stepText) { |
| | |
| | | m.put("percent", "DONE".equals(status) ? 100 : percent); |
| | | m.put("step", step); |
| | | m.put("fromCache", fromCache); |
| | | m.put("staleCache", staleCache); |
| | | m.put("fileName", fileName); |
| | | m.put("error", error); |
| | | long end = finishedAt > 0 ? finishedAt : System.currentTimeMillis(); |
| | |
| | | |
| | | /** |
| | | * 提交导出任务。forceRefresh=true 时忽略缓存重新生成; |
| | | * 否则在「同一报表期、10 分钟内、期间无新导入」的前提下直接复用上次结果。 |
| | | * 否则在「同一报表期 + 数据指纹一致 + 7 天内」的前提下直接复用上次结果。 |
| | | * 同一报表期已有任务在跑时直接返回该任务,避免重复排队。 |
| | | */ |
| | | public Task submit(String type, String period, String mode, Integer toMonth, boolean forceRefresh) { |
| | | cleanupExpired(); |
| | | String key = cacheKey(type, period, mode, toMonth); |
| | | String fingerprint = dataFingerprint(); |
| | | if (!forceRefresh) { |
| | | Task hit = cache.get(key); |
| | | if (cacheUsable(hit)) { |
| | | if (cacheUsable(hit, fingerprint)) { |
| | | hit.fromCache = true; |
| | | hit.staleCache = false; // 命中缓存时不能残留上一次「结果已失效」的提示 |
| | | log.info("报表导出任务命中内存缓存:{} {}(生成于 {})", type, period, new Date(hit.finishedAt)); |
| | | return hit; |
| | | } |
| | | Task disk = loadFromDisk(key, type, period, mode, toMonth); |
| | | Task disk = loadFromDisk(key, type, period, mode, toMonth, fingerprint); |
| | | if (disk != null) { |
| | | tasks.put(disk.id, disk); |
| | | cache.put(key, disk); |
| | |
| | | log.info("报表导出任务已在执行,复用进行中的任务:{} {}({})", type, period, running.id); |
| | | return running; |
| | | } |
| | | Task task = new Task(type, period, mode, toMonth, importVersion()); |
| | | Task task = new Task(type, period, mode, toMonth, fingerprint); |
| | | // 有旧结果但没命中(数据或母版变了)→ 前端提示一句,避免用户以为"缓存坏了" |
| | | task.staleCache = !forceRefresh && hasStaleCache(type, period, mode, toMonth, fingerprint); |
| | | tasks.put(task.id, task); |
| | | pool.submit(() -> run(task)); |
| | | return task; |
| | |
| | | return null; |
| | | } |
| | | |
| | | private boolean cacheUsable(Task t) { |
| | | if (t == null || !"DONE".equals(t.status) || t.data == null) return false; |
| | | private boolean cacheUsable(Task t, String fingerprint) { |
| | | if (t == null || !"DONE".equals(t.status)) return false; |
| | | if (t.data == null && (t.cacheFile == null || !t.cacheFile.isFile())) return false; |
| | | if (System.currentTimeMillis() - t.finishedAt > CACHE_MAX_AGE_MS) return false; |
| | | long now = importVersion(); |
| | | if (now < 0 || t.importVersion < 0) return false; // 取不到数据版本时不冒风险复用 |
| | | return now == t.importVersion; |
| | | if (fingerprint == null || t.fingerprint == null) return false; // 算不出指纹时不冒风险复用 |
| | | return fingerprint.equals(t.fingerprint); |
| | | } |
| | | |
| | | /** 数据版本=import_batch 最大 id:任何一次导入都会让它变化,取不到返回 -1(此时禁用缓存) */ |
| | | private long importVersion() { |
| | | /** |
| | | * 汇总大表数据指纹=「依赖的源表内容校验和」+「母版等文件身份」。 |
| | | * - 源表:库内除 FINGERPRINT_EXCLUDE(投资/能耗/审核/系统/导入日志)以外的全部表,用 CHECKSUM TABLE 取校验和, |
| | | * 能捕捉行数变化、新增删除、**原地 UPDATE**; |
| | | * - 文件:母版目录下全部文件(母版 / 参照件 / _备份_ 回填件)+《网约车订单及全省总量.xlsx》,取「名称+大小+修改时间」。 |
| | | * 任一项变化即指纹变化 → 缓存失效。计算失败返回 null(此时禁用缓存,只重算不复用)。 |
| | | */ |
| | | private String dataFingerprint() { |
| | | try { |
| | | List<Object> v = importBatchMapper.selectObjs(new QueryWrapper<ImportBatch>().select("IFNULL(MAX(id),0)")); |
| | | if (v == null || v.isEmpty() || v.get(0) == null) return 0L; |
| | | return ((Number) v.get(0)).longValue(); |
| | | List<String> parts = new ArrayList<>(); |
| | | try (Connection conn = dataSource.getConnection()) { |
| | | List<String> tables = new ArrayList<>(); |
| | | try (Statement st = conn.createStatement(); |
| | | ResultSet rs = st.executeQuery("SELECT table_name FROM information_schema.tables " |
| | | + "WHERE table_schema = DATABASE() AND table_type = 'BASE TABLE' ORDER BY table_name")) { |
| | | while (rs.next()) tables.add(rs.getString(1)); |
| | | } |
| | | for (String t : tables) { |
| | | if (FINGERPRINT_EXCLUDE.contains(t)) continue; |
| | | try (Statement st = conn.createStatement(); |
| | | ResultSet rs = st.executeQuery("CHECKSUM TABLE `" + t + "`")) { |
| | | if (rs.next()) parts.add("T:" + t + ":" + rs.getString(2)); |
| | | } |
| | | } |
| | | } |
| | | for (File f : listIdentityFiles(resolveBaseDir(summaryTemplateDir))) { |
| | | parts.add("F:" + f.getName() + ":" + f.length() + ":" + f.lastModified()); |
| | | } |
| | | File wyc = new File(resolveBaseDir(cityPassengerTemplateDir), "网约车订单及全省总量.xlsx"); |
| | | if (wyc.isFile()) parts.add("F:" + wyc.getName() + ":" + wyc.length() + ":" + wyc.lastModified()); |
| | | return sha256Hex(String.join("|", parts)).substring(0, 16); |
| | | } catch (Exception e) { |
| | | log.warn("读取导入数据版本失败,本次禁用导出缓存:{}", e.getMessage()); |
| | | return -1L; |
| | | log.warn("计算数据指纹失败,本次禁用导出缓存:{}", e.getMessage()); |
| | | return null; |
| | | } |
| | | } |
| | | |
| | | /** 母版目录等基础目录定位(与模板解析同策略:配置目录 / user.dir 相对 / 上级目录) */ |
| | | private File resolveBaseDir(String dir) { |
| | | if (dir == null) return null; |
| | | String rel = dir; |
| | | while (rel.startsWith("./")) rel = rel.substring(2); |
| | | String[] roots = {rel, System.getProperty("user.dir") + "/" + rel, System.getProperty("user.dir") + "/../" + rel}; |
| | | for (String root : roots) { |
| | | if (root == null || root.trim().isEmpty()) continue; |
| | | File f = new File(root); |
| | | if (f.isDirectory()) return f; |
| | | } |
| | | return null; |
| | | } |
| | | |
| | | /** 目录下的文件(按名字排序,排除子目录,保证指纹稳定) */ |
| | | private List<File> listIdentityFiles(File dir) { |
| | | List<File> out = new ArrayList<>(); |
| | | if (dir == null) return out; |
| | | File[] fs = dir.listFiles(File::isFile); |
| | | if (fs == null) return out; |
| | | Arrays.sort(fs, Comparator.comparing(File::getName)); |
| | | out.addAll(Arrays.asList(fs)); |
| | | return out; |
| | | } |
| | | |
| | | private String sha256Hex(String text) throws Exception { |
| | | MessageDigest md = MessageDigest.getInstance("SHA-256"); |
| | | byte[] d = md.digest(text.getBytes(StandardCharsets.UTF_8)); |
| | | StringBuilder sb = new StringBuilder(); |
| | | for (int i = 0; i < 8; i++) sb.append(String.format("%02x", d[i])); // 取前 8 字节=16 位十六进制,够用且文件名短 |
| | | return sb.toString(); |
| | | } |
| | | |
| | | /** 该报表期是否存在「与当前指纹不一致」的旧结果(用于提示「上次结果已失效,正在按最新数据重算」) */ |
| | | private boolean hasStaleCache(String type, String period, String mode, Integer toMonth, String fingerprint) { |
| | | try { |
| | | File dir = cacheDir(); |
| | | if (dir == null) return false; |
| | | String prefix = cacheFilePrefix(type, period, mode, toMonth); |
| | | File[] fs = dir.listFiles((d, n) -> n.startsWith(prefix + "_v") && n.endsWith(".xlsx")); |
| | | if (fs == null || fs.length == 0) return false; |
| | | String current = prefix + "_v" + fingerprint + ".xlsx"; |
| | | for (File f : fs) { |
| | | if (!f.getName().equals(current)) return true; |
| | | } |
| | | return false; |
| | | } catch (Exception e) { |
| | | return false; |
| | | } |
| | | } |
| | | |
| | |
| | | try { |
| | | File dir = cacheDir(); |
| | | if (dir == null) return; |
| | | long version = t.importVersion; |
| | | if (version < 0) return; |
| | | String fingerprint = t.fingerprint; |
| | | if (fingerprint == null) return; |
| | | String prefix = cacheFilePrefix(t.type, t.period, t.mode, t.toMonth); |
| | | File target = new File(dir, prefix + "_v" + version + ".xlsx"); |
| | | File target = new File(dir, prefix + "_v" + fingerprint + ".xlsx"); |
| | | java.nio.file.Files.write(target.toPath(), data); |
| | | File[] siblings = dir.listFiles((d, n) -> n.startsWith(prefix + "_v") && n.endsWith(".xlsx")); |
| | | if (siblings != null) { |
| | |
| | | } |
| | | } |
| | | |
| | | /** 从落盘缓存里找「同一报表期 + 同一数据版本 + 未超期」的结果;没有返回 null */ |
| | | private Task loadFromDisk(String key, String type, String period, String mode, Integer toMonth) { |
| | | /** 从落盘缓存里找「同一报表期 + 同一数据指纹 + 未超期」的结果;没有返回 null */ |
| | | private Task loadFromDisk(String key, String type, String period, String mode, Integer toMonth, String fingerprint) { |
| | | try { |
| | | if (fingerprint == null) return null; // 算不出指纹时不冒风险复用 |
| | | File dir = cacheDir(); |
| | | if (dir == null) return null; |
| | | long version = importVersion(); |
| | | if (version < 0) return null; // 取不到数据版本时不冒风险复用 |
| | | File f = new File(dir, cacheFilePrefix(type, period, mode, toMonth) + "_v" + version + ".xlsx"); |
| | | File f = new File(dir, cacheFilePrefix(type, period, mode, toMonth) + "_v" + fingerprint + ".xlsx"); |
| | | if (!f.isFile()) return null; |
| | | long finished = f.lastModified(); |
| | | if (finished <= 0 || System.currentTimeMillis() - finished > CACHE_MAX_AGE_MS) return null; |
| | | Task t = new Task(type, period, mode, toMonth, version); |
| | | Task t = new Task(type, period, mode, toMonth, fingerprint); |
| | | t.cacheFile = f; |
| | | t.createdAt = finished; // 复用历史结果:以生成时间作为起点,避免 elapsedMs 出现负数 |
| | | t.fileName = "生成_道路运输量汇总表_" + period + ".xlsx"; |