Ver código fonte

流程调整。超过六个月的已完成流程进行冷热数据分离操作

徐滕 4 dias atrás
pai
commit
c9631f8a3e

+ 90 - 0
jeeplus-modules/jeeplus-flowable/src/main/java/com/jeeplus/flowable/controller/FlowArchiveController.java

@@ -0,0 +1,90 @@
+/**
+ * Copyright &copy; 2021-2026 <a href="http://www.jeeplus.org/">JeePlus</a> All rights reserved.
+ */
+package com.jeeplus.flowable.controller;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * 流程数据归档控制器
+ * 提供REST接口供XXL-Job或其他定时任务调用,触发历史数据从热表迁移到冷表
+ *
+ * @author jeeplus
+ * @version 2024-01-01
+ */
+@Slf4j
+@RestController
+@RequestMapping("/flow/archive")
+public class FlowArchiveController {
+
+    @Autowired
+    private JdbcTemplate jdbcTemplate;
+
+    /**
+     * 触发数据归档迁移
+     * 调用存储过程 sp_flow_archive,将3个月前已结束的流程数据从热表迁移到冷表
+     * XXL-Job配置: HTTP GET /flow/archive/trigger
+     *
+     * @return 迁移结果
+     */
+    @GetMapping("/trigger")
+    public Map<String, Object> triggerArchive() {
+        Map<String, Object> result = new HashMap<>();
+        try {
+            log.info("开始执行流程数据归档...");
+            jdbcTemplate.execute("CALL sp_flow_archive()");
+            log.info("流程数据归档执行完成");
+            result.put("code", 200);
+            result.put("msg", "归档执行成功");
+        } catch (Exception e) {
+            log.error("流程数据归档执行失败", e);
+            result.put("code", 500);
+            result.put("msg", "归档执行失败: " + e.getMessage());
+        }
+        return result;
+    }
+
+    /**
+     * 查看归档统计信息
+     */
+    @GetMapping("/stats")
+    public Map<String, Object> getArchiveStats() {
+        Map<String, Object> result = new HashMap<>();
+        try {
+            Map<String, Object> stats = new HashMap<>();
+            stats.put("hot_procinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_procinst", Integer.class));
+            stats.put("archive_procinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_procinst_archive", Integer.class));
+            stats.put("hot_actinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_actinst", Integer.class));
+            stats.put("archive_actinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_actinst_archive", Integer.class));
+            stats.put("hot_varinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_varinst", Integer.class));
+            stats.put("archive_varinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_varinst_archive", Integer.class));
+            stats.put("hot_taskinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_taskinst", Integer.class));
+            stats.put("archive_taskinst", jdbcTemplate.queryForObject("SELECT COUNT(*) FROM act_hi_taskinst_archive", Integer.class));
+
+            // 最近一次归档日志
+            Map<String, Object> lastLog = null;
+            try {
+                lastLog = jdbcTemplate.queryForList("SELECT * FROM flow_archive_log ORDER BY archive_time DESC LIMIT 1").stream()
+                        .findFirst().orElse(null);
+            } catch (Exception ignored) {
+                // flow_archive_log表可能还未创建
+            }
+
+            result.put("code", 200);
+            result.put("stats", stats);
+            result.put("lastArchiveLog", lastLog);
+        } catch (Exception e) {
+            result.put("code", 500);
+            result.put("msg", "查询失败: " + e.getMessage());
+        }
+        return result;
+    }
+}

+ 133 - 0
jeeplus-modules/jeeplus-flowable/src/main/java/com/jeeplus/flowable/service/FlowArchiveService.java

@@ -0,0 +1,133 @@
+/**
+ * Copyright &copy; 2021-2026 <a href="http://www.jeeplus.org/">JeePlus</a> All rights reserved.
+ */
+package com.jeeplus.flowable.service;
+
+import lombok.extern.slf4j.Slf4j;
+import org.flowable.engine.impl.persistence.entity.HistoricActivityInstanceEntityImpl;
+import org.flowable.engine.history.HistoricActivityInstance;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.jdbc.core.RowMapper;
+import org.springframework.stereotype.Service;
+
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Timestamp;
+import java.util.*;
+
+/**
+ * 流程归档数据查询服务(冷表回查)
+ * 当热表中查不到数据时,从归档表中查询历史数据
+ *
+ * @author jeeplus
+ * @version 2024-01-01
+ */
+@Slf4j
+@Service
+public class FlowArchiveService {
+
+    @Autowired
+    private JdbcTemplate jdbcTemplate;
+
+    /**
+     * 从归档表查询活动实例列表
+     */
+    public List<HistoricActivityInstance> findActivityInstances(String processInstanceId) {
+        String sql = "SELECT * FROM act_hi_actinst_archive WHERE PROC_INST_ID_ = ? ORDER BY START_TIME_ ASC";
+        return jdbcTemplate.query(sql, new Object[]{processInstanceId}, new HistoricActivityInstanceRowMapper());
+    }
+
+    /**
+     * 从归档表查询已结束的活动实例列表
+     */
+    public List<HistoricActivityInstance> findFinishedActivityInstances(String processInstanceId) {
+        String sql = "SELECT * FROM act_hi_actinst_archive WHERE PROC_INST_ID_ = ? AND END_TIME_ IS NOT NULL ORDER BY END_TIME_ ASC";
+        return jdbcTemplate.query(sql, new Object[]{processInstanceId}, new HistoricActivityInstanceRowMapper());
+    }
+
+    /**
+     * 从归档表查询历史流程实例
+     */
+    public Map<String, Object> findHistoricProcessInstance(String processInstanceId) {
+        String sql = "SELECT * FROM act_hi_procinst_archive WHERE PROC_INST_ID_ = ?";
+        List<Map<String, Object>> result = jdbcTemplate.queryForList(sql, processInstanceId);
+        return result.isEmpty() ? null : result.get(0);
+    }
+
+    /**
+     * 从归档表查询历史变量
+     */
+    public List<Map<String, Object>> findHistoricVariables(String processInstanceId) {
+        String sql = "SELECT * FROM act_hi_varinst_archive WHERE PROC_INST_ID_ = ?";
+        return jdbcTemplate.queryForList(sql, processInstanceId);
+    }
+
+    /**
+     * 从归档表查询指定流程变量的值
+     */
+    public Map<String, Object> findHistoricVariable(String processInstanceId, String variableName) {
+        String sql = "SELECT * FROM act_hi_varinst_archive WHERE PROC_INST_ID_ = ? AND NAME_ = ?";
+        List<Map<String, Object>> result = jdbcTemplate.queryForList(sql, processInstanceId, variableName);
+        return result.isEmpty() ? null : result.get(0);
+    }
+
+    /**
+     * 从归档表查询任务评论
+     */
+    public List<Map<String, Object>> findTaskComments(String taskId) {
+        String sql = "SELECT * FROM act_hi_comment_archive WHERE TASK_ID_ = ? ORDER BY TIME_ DESC";
+        return jdbcTemplate.queryForList(sql, taskId);
+    }
+
+    /**
+     * 从归档表查询历史任务实例
+     */
+    public List<Map<String, Object>> findHistoricTaskInstances(String processInstanceId) {
+        String sql = "SELECT * FROM act_hi_taskinst_archive WHERE PROC_INST_ID_ = ?";
+        return jdbcTemplate.queryForList(sql, processInstanceId);
+    }
+
+    /**
+     * 检查流程实例是否在归档表中
+     */
+    public boolean isArchived(String processInstanceId) {
+        String sql = "SELECT COUNT(1) FROM act_hi_procinst_archive WHERE PROC_INST_ID_ = ?";
+        Integer count = jdbcTemplate.queryForObject(sql, Integer.class, processInstanceId);
+        return count != null && count > 0;
+    }
+
+    /**
+     * HistoricActivityInstance 行映射器
+     * 将归档表数据映射为 Flowable 的 HistoricActivityInstance 对象
+     */
+    private static class HistoricActivityInstanceRowMapper implements RowMapper<HistoricActivityInstance> {
+        @Override
+        public HistoricActivityInstance mapRow(ResultSet rs, int rowNum) throws SQLException {
+            HistoricActivityInstanceEntityImpl entity = new HistoricActivityInstanceEntityImpl();
+            entity.setId(rs.getString("ID_"));
+            entity.setRevision(rs.getInt("REV_"));
+            entity.setProcessDefinitionId(rs.getString("PROC_DEF_ID_"));
+            entity.setProcessInstanceId(rs.getString("PROC_INST_ID_"));
+            entity.setExecutionId(rs.getString("EXECUTION_ID_"));
+            entity.setActivityId(rs.getString("ACT_ID_"));
+            entity.setTaskId(rs.getString("TASK_ID_"));
+            // entity.setCallProcessInstanceId - 跳过,非关键字段
+            entity.setActivityName(rs.getString("ACT_NAME_"));
+            entity.setActivityType(rs.getString("ACT_TYPE_"));
+            entity.setAssignee(rs.getString("ASSIGNEE_"));
+            Timestamp startTime = rs.getTimestamp("START_TIME_");
+            if (startTime != null) {
+                entity.setStartTime(new Date(startTime.getTime()));
+            }
+            Timestamp endTime = rs.getTimestamp("END_TIME_");
+            if (endTime != null) {
+                entity.setEndTime(new Date(endTime.getTime()));
+            }
+            entity.setDurationInMillis(rs.getLong("DURATION_"));
+            entity.setDeleteReason(rs.getString("DELETE_REASON_"));
+            entity.setTenantId(rs.getString("TENANT_ID_"));
+            return entity;
+        }
+    }
+}