fix: 优化逻辑
This commit is contained in:
parent
24515bf559
commit
be328e3107
@ -71,6 +71,7 @@ public class SecurityConfig {
|
||||
// .requestMatchers("/system/**").permitAll()
|
||||
// .requestMatchers("/qgcExport/**").permitAll()
|
||||
.requestMatchers("/eng/**").permitAll()
|
||||
.requestMatchers("/system/**").permitAll()
|
||||
// .requestMatchers("/eq/**").permitAll()
|
||||
// .requestMatchers("/env/**").permitAll()
|
||||
// .requestMatchers("/warn/**").permitAll()
|
||||
|
||||
@ -36,5 +36,11 @@ public class DynamicDataSource extends AbstractRoutingDataSource {
|
||||
contextHolder.remove();
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* 根据数据源 key 获取目标数据源实例(如 dm-master / oracle-master)
|
||||
*/
|
||||
public DataSource getDataSourceByKey(String key) {
|
||||
Map<Object, DataSource> resolved = getResolvedDataSources();
|
||||
return resolved == null ? null : resolved.get(key);
|
||||
}
|
||||
}
|
||||
|
||||
@ -2592,7 +2592,7 @@ public class FpRunServiceImpl implements FpRunService {
|
||||
" t2.pycnt AS pycnt, " +
|
||||
" NULL AS pyRate, " +
|
||||
" t.STTP AS sttpCode, " +
|
||||
" t1.planFtp AS planFtp, " +
|
||||
" t.ZYGYDXMS AS planFtp, " +
|
||||
" t.ZYGYDX AS fallbackFishJson, " +
|
||||
" NVL(t.ZYGYDXMS, t.JGGYDXMS) AS fallbackFishText " +
|
||||
"FROM SD_FPSS_B_H t " +
|
||||
@ -3502,7 +3502,7 @@ public class FpRunServiceImpl implements FpRunService {
|
||||
|
||||
private String buildFpssrlQdrtpSOrderBySql(List<DataSourceRequest.SortDescriptor> sortList) {
|
||||
if (CollUtil.isEmpty(sortList)) {
|
||||
return " ORDER BY t.YEAR DESC, t.MONTH DESC, t.DR DESC";
|
||||
return " ORDER BY t.YEAR DESC, t.MONTH DESC";
|
||||
}
|
||||
List<String> orders = new ArrayList<>();
|
||||
for (DataSourceRequest.SortDescriptor sort : sortList) {
|
||||
|
||||
@ -0,0 +1,93 @@
|
||||
package com.yfd.platform.system.controller;
|
||||
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import com.yfd.platform.config.ResponseResult;
|
||||
import com.yfd.platform.system.domain.DataSyncRequest;
|
||||
import com.yfd.platform.utils.DataSyncUtil;
|
||||
import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.PostMapping;
|
||||
import org.springframework.web.bind.annotation.RequestBody;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
/**
|
||||
* 数据同步控制器:支持 DM / Oracle 之间数据同步(MERGE INTO 方式)
|
||||
*/
|
||||
@RestController
|
||||
@RequestMapping("/system/dataSync")
|
||||
@Tag(name = "数据同步")
|
||||
@Validated
|
||||
public class DataSyncController {
|
||||
|
||||
@Resource
|
||||
private DataSyncUtil dataSyncUtil;
|
||||
|
||||
/**
|
||||
* DM 数据同步到 Oracle(使用调用方提供的连接信息)
|
||||
* <p>
|
||||
* 请求示例(把 DM 库 SD_FPSS_R 同步到 Oracle 库 SD_FPSS_R):
|
||||
* <pre>
|
||||
* {
|
||||
* "sourceTable": "SD_FPSS_R",
|
||||
* "targetTable": "SD_FPSS_R",
|
||||
* "keyColumns": ["ID"],
|
||||
* "clearTarget": true,
|
||||
* "batchSize": 500,
|
||||
* "sourceConn": {
|
||||
* "url": "jdbc:dm://172.16.21.143:5236/QGC_REFA_TEST",
|
||||
* "username": "QGC_REFA_TEST",
|
||||
* "password": "xxxx"
|
||||
* },
|
||||
* "targetConn": {
|
||||
* "url": "jdbc:oracle:thin:@172.16.31.172:1521/SDLYZ",
|
||||
* "username": "SDLY_QX",
|
||||
* "password": "xxxx"
|
||||
* }
|
||||
* }
|
||||
* </pre>
|
||||
*/
|
||||
@PostMapping("/dmToOracle")
|
||||
@Operation(summary = "DM 数据同步到 Oracle(自定义连接,MERGE INTO)")
|
||||
public ResponseResult dmToOracle(@RequestBody DataSyncRequest request) {
|
||||
DataSyncUtil.DataSyncResult result = dataSyncUtil.syncWithConnections(buildConfig(request),
|
||||
request.getSourceConn(), request.getTargetConn());
|
||||
return ResponseResult.successData(result);
|
||||
}
|
||||
|
||||
/**
|
||||
* 通用数据同步(自定义连接,可指定任意方向)
|
||||
*/
|
||||
@PostMapping("/sync")
|
||||
@Operation(summary = "通用数据同步(自定义连接,MERGE INTO)")
|
||||
public ResponseResult sync(@RequestBody DataSyncRequest request) {
|
||||
DataSyncUtil.DataSyncResult result = dataSyncUtil.syncWithConnections(buildConfig(request),
|
||||
request.getSourceConn(), request.getTargetConn());
|
||||
return ResponseResult.successData(result);
|
||||
}
|
||||
|
||||
private DataSyncUtil.DataSyncConfig buildConfig(DataSyncRequest request) {
|
||||
if (request == null) {
|
||||
throw new IllegalArgumentException("请求参数不能为空");
|
||||
}
|
||||
if (StrUtil.isBlank(request.getSourceTable())) {
|
||||
throw new IllegalArgumentException("源表(sourceTable)不能为空");
|
||||
}
|
||||
if (StrUtil.isBlank(request.getTargetTable())) {
|
||||
throw new IllegalArgumentException("目标表(targetTable)不能为空");
|
||||
}
|
||||
DataSyncUtil.DataSyncConfig config = new DataSyncUtil.DataSyncConfig();
|
||||
config.setSourceTable(request.getSourceTable());
|
||||
config.setTargetTable(request.getTargetTable());
|
||||
config.setColumns(request.getColumns());
|
||||
config.setColumnMapping(request.getColumnMapping());
|
||||
config.setKeyColumns(request.getKeyColumns());
|
||||
config.setConditionSql(request.getConditionSql());
|
||||
config.setOrderColumn(request.getOrderColumn());
|
||||
config.setBatchSize(request.getBatchSize());
|
||||
config.setClearTarget(request.getClearTarget());
|
||||
return config;
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,49 @@
|
||||
package com.yfd.platform.system.domain;
|
||||
|
||||
import com.yfd.platform.utils.DataSyncUtil;
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* 数据同步请求
|
||||
*/
|
||||
@Data
|
||||
@Schema(description = "数据同步请求")
|
||||
public class DataSyncRequest {
|
||||
|
||||
@Schema(description = "源表名(DM/Oracle)")
|
||||
private String sourceTable;
|
||||
|
||||
@Schema(description = "目标表名(可与源表名不同)")
|
||||
private String targetTable;
|
||||
|
||||
@Schema(description = "源表同步列(不指定则自动读取源表全列)")
|
||||
private List<String> columns;
|
||||
|
||||
@Schema(description = "字段映射:源列名 -> 目标列名")
|
||||
private Map<String, String> columnMapping;
|
||||
|
||||
@Schema(description = "MERGE 匹配主键列(目标表列名,必填)")
|
||||
private List<String> keyColumns;
|
||||
|
||||
@Schema(description = "增量同步条件 SQL 片段(基于源表列)")
|
||||
private String conditionSql;
|
||||
|
||||
@Schema(description = "分页排序列(源表列,不指定则取第一列)")
|
||||
private String orderColumn;
|
||||
|
||||
@Schema(description = "每批同步条数,默认 500")
|
||||
private Integer batchSize;
|
||||
|
||||
@Schema(description = "同步前是否清空目标表(全量覆盖)")
|
||||
private Boolean clearTarget;
|
||||
|
||||
@Schema(description = "源库连接信息(driverClassName 不填时按 URL 自动识别)")
|
||||
private DataSyncUtil.DbConnectionInfo sourceConn;
|
||||
|
||||
@Schema(description = "目标库连接信息(driverClassName 不填时按 URL 自动识别)")
|
||||
private DataSyncUtil.DbConnectionInfo targetConn;
|
||||
}
|
||||
637
backend/src/main/java/com/yfd/platform/utils/DataSyncUtil.java
Normal file
637
backend/src/main/java/com/yfd/platform/utils/DataSyncUtil.java
Normal file
@ -0,0 +1,637 @@
|
||||
package com.yfd.platform.utils;
|
||||
|
||||
import cn.hutool.core.collection.CollUtil;
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import com.alibaba.druid.pool.DruidDataSource;
|
||||
import com.yfd.platform.datasource.DataSourceKeys;
|
||||
import com.yfd.platform.datasource.DynamicDataSource;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import javax.sql.DataSource;
|
||||
import java.sql.ResultSetMetaData;
|
||||
import java.sql.SQLException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* 数据同步控制类:支持 Oracle 与达梦(DM) 之间双向同步,使用 MERGE INTO 方式合并数据。
|
||||
* <p>
|
||||
* 支持两种数据源来源:
|
||||
* <ul>
|
||||
* <li>项目配置的数据源:{@link #syncTable(String, String, DataSyncConfig)}(dm-master / oracle-master 等 key)</li>
|
||||
* <li>调用方自定义连接:{@link #syncWithConnections(DataSyncConfig, DbConnectionInfo, DbConnectionInfo)}(URL/账号/密码)</li>
|
||||
* </ul>
|
||||
* </p>
|
||||
* <p>
|
||||
* 支持:源表名与目标表名不同、源字段与目标字段不同(通过 {@link DataSyncConfig#setColumnMapping} 配置映射)。
|
||||
* </p>
|
||||
*/
|
||||
@Component
|
||||
public class DataSyncUtil {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(DataSyncUtil.class);
|
||||
|
||||
@Resource
|
||||
private DynamicDataSource dynamicDataSource;
|
||||
|
||||
// ==================== 项目数据源入口 ====================
|
||||
|
||||
/**
|
||||
* 达梦(DM) → Oracle 数据同步(使用项目配置的 dm-master / oracle-master 数据源)
|
||||
*/
|
||||
public DataSyncResult syncDmToOracle(DataSyncConfig config) {
|
||||
return syncTable(DataSourceKeys.DM_MASTER, DataSourceKeys.ORACLE_MASTER, config);
|
||||
}
|
||||
|
||||
/**
|
||||
* Oracle → 达梦(DM) 数据同步(使用项目配置的 oracle-master / dm-master 数据源)
|
||||
*/
|
||||
public DataSyncResult syncOracleToDm(DataSyncConfig config) {
|
||||
return syncTable(DataSourceKeys.ORACLE_MASTER, DataSourceKeys.DM_MASTER, config);
|
||||
}
|
||||
|
||||
/**
|
||||
* 通用数据同步:从源数据源读取,MERGE INTO 写入目标数据源(使用项目配置的数据源 key)
|
||||
*
|
||||
* @param sourceDsKey 源数据源 key(如 dm-master / oracle-master)
|
||||
* @param targetDsKey 目标数据源 key
|
||||
* @param config 同步配置
|
||||
*/
|
||||
public DataSyncResult syncTable(String sourceDsKey, String targetDsKey, DataSyncConfig config) {
|
||||
DataSource sourceDataSource = dynamicDataSource.getDataSourceByKey(sourceDsKey);
|
||||
DataSource targetDataSource = dynamicDataSource.getDataSourceByKey(targetDsKey);
|
||||
if (sourceDataSource == null) {
|
||||
throw new IllegalArgumentException("未找到数据源: " + sourceDsKey);
|
||||
}
|
||||
if (targetDataSource == null) {
|
||||
throw new IllegalArgumentException("未找到数据源: " + targetDsKey);
|
||||
}
|
||||
return doSync(new JdbcTemplate(sourceDataSource), new JdbcTemplate(targetDataSource), config);
|
||||
}
|
||||
|
||||
// ==================== 自定义连接入口 ====================
|
||||
|
||||
/**
|
||||
* 使用调用方提供的连接信息同步数据(DM → Oracle / Oracle → DM 均适用)
|
||||
* <p>
|
||||
* 典型场景:把 DM 库 {@code SD_FPSS_R} 同步到 Oracle 库 {@code SD_FPSS_R},
|
||||
* 连接信息由调用方在请求中提供。
|
||||
* </p>
|
||||
*
|
||||
* @param config 同步配置
|
||||
* @param sourceConn 源库连接信息(driverClassName 不填时按 URL 自动识别)
|
||||
* @param targetConn 目标库连接信息
|
||||
*/
|
||||
public DataSyncResult syncWithConnections(DataSyncConfig config, DbConnectionInfo sourceConn, DbConnectionInfo targetConn) {
|
||||
if (sourceConn == null || StrUtil.isBlank(sourceConn.getUrl())) {
|
||||
throw new IllegalArgumentException("源库连接信息(url)不能为空");
|
||||
}
|
||||
if (targetConn == null || StrUtil.isBlank(targetConn.getUrl())) {
|
||||
throw new IllegalArgumentException("目标库连接信息(url)不能为空");
|
||||
}
|
||||
DruidDataSource sourceDataSource = buildDataSource(sourceConn);
|
||||
DruidDataSource targetDataSource = buildDataSource(targetConn);
|
||||
try {
|
||||
return doSync(new JdbcTemplate(sourceDataSource), new JdbcTemplate(targetDataSource), config);
|
||||
} finally {
|
||||
sourceDataSource.close();
|
||||
targetDataSource.close();
|
||||
}
|
||||
}
|
||||
|
||||
// ==================== 核心同步逻辑 ====================
|
||||
|
||||
/**
|
||||
* 核心同步:从源 JdbcTemplate 分页读取,MERGE INTO 写入目标 JdbcTemplate
|
||||
*/
|
||||
private DataSyncResult doSync(JdbcTemplate sourceJdbc, JdbcTemplate targetJdbc, DataSyncConfig config) {
|
||||
long startTime = System.currentTimeMillis();
|
||||
DataSyncResult result = new DataSyncResult();
|
||||
result.setStartTime(startTime);
|
||||
try {
|
||||
validateConfig(config);
|
||||
|
||||
// 1. 源表列(未指定时自动读取源表结构)
|
||||
List<String> sourceColumns = normalizeColumns(config.getColumns());
|
||||
if (CollUtil.isEmpty(sourceColumns)) {
|
||||
sourceColumns = loadTableColumns(sourceJdbc, config.getSourceTable());
|
||||
if (CollUtil.isEmpty(sourceColumns)) {
|
||||
throw new IllegalArgumentException("未从源表 [" + config.getSourceTable() + "] 获取到任何列");
|
||||
}
|
||||
}
|
||||
|
||||
// 2. 目标表列:按 columnMapping 映射,未配置映射的列名保持一致
|
||||
Map<String, String> columnMapping = normalizeMapping(config.getColumnMapping());
|
||||
List<String> targetColumns = new ArrayList<>(sourceColumns.size());
|
||||
for (String sourceColumn : sourceColumns) {
|
||||
targetColumns.add(columnMapping.getOrDefault(sourceColumn, sourceColumn));
|
||||
}
|
||||
|
||||
// 3. 主键列:keyColumns 配置为目标表列名,反查对应的源表列名用于取值
|
||||
List<String> keyColumns = normalizeColumns(config.getKeyColumns());
|
||||
Map<String, String> reverseMapping = new HashMap<>();
|
||||
for (Map.Entry<String, String> entry : columnMapping.entrySet()) {
|
||||
reverseMapping.put(entry.getValue(), entry.getKey());
|
||||
}
|
||||
List<String> keySourceColumns = new ArrayList<>(keyColumns.size());
|
||||
for (String keyColumn : keyColumns) {
|
||||
keySourceColumns.add(reverseMapping.getOrDefault(keyColumn, keyColumn));
|
||||
}
|
||||
|
||||
// 4. 可选:同步前清空目标表(全量覆盖)
|
||||
if (Boolean.TRUE.equals(config.getClearTarget())) {
|
||||
clearTargetTable(targetJdbc, config.getTargetTable());
|
||||
}
|
||||
|
||||
// 5. 分页读取源数据 → 批量 MERGE 写入目标库
|
||||
int batchSize = config.getBatchSize() != null && config.getBatchSize() > 0
|
||||
? config.getBatchSize() : 500;
|
||||
long totalRead = 0;
|
||||
long totalProcessed = 0;
|
||||
int offset = 0;
|
||||
|
||||
while (true) {
|
||||
List<Map<String, Object>> rows = readPage(sourceJdbc, config, sourceColumns, offset, batchSize);
|
||||
if (CollUtil.isEmpty(rows)) {
|
||||
break;
|
||||
}
|
||||
int[] effects = mergeInto(targetJdbc, config, sourceColumns, targetColumns,
|
||||
keyColumns, keySourceColumns, rows);
|
||||
for (int effect : effects) {
|
||||
totalProcessed += effect;
|
||||
}
|
||||
totalRead += rows.size();
|
||||
offset += rows.size();
|
||||
if (rows.size() < batchSize) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
result.setSuccess(true);
|
||||
result.setTotalRead(totalRead);
|
||||
result.setTotalProcessed(totalProcessed);
|
||||
result.setEndTime(System.currentTimeMillis());
|
||||
result.setElapsedMillis(result.getEndTime() - startTime);
|
||||
result.setMessage("同步完成:读取 " + totalRead + " 条,MERGE 处理 " + totalProcessed + " 条");
|
||||
log.info("数据同步完成 sourceTable={} targetTable={} {}", config.getSourceTable(),
|
||||
config.getTargetTable(), result.getMessage());
|
||||
} catch (Exception e) {
|
||||
result.setSuccess(false);
|
||||
result.setEndTime(System.currentTimeMillis());
|
||||
result.setElapsedMillis(result.getEndTime() - startTime);
|
||||
result.setMessage("数据同步失败:" + e.getMessage());
|
||||
log.error("数据同步失败 sourceTable={} targetTable={}",
|
||||
config == null ? "" : config.getSourceTable(),
|
||||
config == null ? "" : config.getTargetTable(), e);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
// ==================== 内部实现 ====================
|
||||
|
||||
private void validateConfig(DataSyncConfig config) {
|
||||
if (config == null) {
|
||||
throw new IllegalArgumentException("同步配置(config)不能为空");
|
||||
}
|
||||
if (StrUtil.isBlank(config.getSourceTable())) {
|
||||
throw new IllegalArgumentException("源表(sourceTable)不能为空");
|
||||
}
|
||||
if (StrUtil.isBlank(config.getTargetTable())) {
|
||||
throw new IllegalArgumentException("目标表(targetTable)不能为空");
|
||||
}
|
||||
if (CollUtil.isEmpty(config.getKeyColumns())) {
|
||||
throw new IllegalArgumentException("必须配置主键列(keyColumns)用于 MERGE 匹配");
|
||||
}
|
||||
}
|
||||
|
||||
private List<String> normalizeColumns(List<String> columns) {
|
||||
if (CollUtil.isEmpty(columns)) {
|
||||
return new ArrayList<>();
|
||||
}
|
||||
return columns.stream()
|
||||
.filter(StrUtil::isNotBlank)
|
||||
.map(String::trim)
|
||||
.map(String::toUpperCase)
|
||||
.distinct()
|
||||
.collect(Collectors.toList());
|
||||
}
|
||||
|
||||
private Map<String, String> normalizeMapping(Map<String, String> mapping) {
|
||||
Map<String, String> normalized = new HashMap<>();
|
||||
if (mapping == null || mapping.isEmpty()) {
|
||||
return normalized;
|
||||
}
|
||||
for (Map.Entry<String, String> entry : mapping.entrySet()) {
|
||||
if (StrUtil.isNotBlank(entry.getKey()) && StrUtil.isNotBlank(entry.getValue())) {
|
||||
normalized.put(entry.getKey().trim().toUpperCase(), entry.getValue().trim().toUpperCase());
|
||||
}
|
||||
}
|
||||
return normalized;
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据连接信息构建 Druid 数据源,driverClassName 未填时按 URL 自动识别
|
||||
*/
|
||||
private DruidDataSource buildDataSource(DbConnectionInfo conn) {
|
||||
DruidDataSource dataSource = new DruidDataSource();
|
||||
dataSource.setDriverClassName(resolveDriverClassName(conn));
|
||||
dataSource.setUrl(conn.getUrl());
|
||||
dataSource.setUsername(conn.getUsername());
|
||||
dataSource.setPassword(conn.getPassword());
|
||||
dataSource.setInitialSize(1);
|
||||
dataSource.setMinIdle(1);
|
||||
dataSource.setMaxActive(5);
|
||||
dataSource.setMaxWait(30000);
|
||||
dataSource.setValidationQuery("SELECT 1 FROM DUAL");
|
||||
dataSource.setTestWhileIdle(true);
|
||||
dataSource.setBreakAfterAcquireFailure(true);
|
||||
dataSource.setConnectionErrorRetryAttempts(0);
|
||||
try {
|
||||
dataSource.init();
|
||||
} catch (SQLException e) {
|
||||
dataSource.close();
|
||||
throw new IllegalStateException("初始化源数据源失败: " + conn.getUrl(), e);
|
||||
}
|
||||
return dataSource;
|
||||
}
|
||||
|
||||
private String resolveDriverClassName(DbConnectionInfo conn) {
|
||||
if (StrUtil.isNotBlank(conn.getDriverClassName())) {
|
||||
return conn.getDriverClassName();
|
||||
}
|
||||
String url = conn.getUrl() == null ? "" : conn.getUrl().toLowerCase();
|
||||
if (url.startsWith("jdbc:dm")) {
|
||||
return "dm.jdbc.driver.DmDriver";
|
||||
}
|
||||
if (url.startsWith("jdbc:oracle")) {
|
||||
return "oracle.jdbc.OracleDriver";
|
||||
}
|
||||
throw new IllegalArgumentException("无法根据 URL 识别驱动类型,请显式传入 driverClassName: " + conn.getUrl());
|
||||
}
|
||||
|
||||
/**
|
||||
* 读取源表全部列名(大写)
|
||||
*/
|
||||
private List<String> loadTableColumns(JdbcTemplate sourceJdbc, String table) {
|
||||
List<String> columns = new ArrayList<>();
|
||||
sourceJdbc.query("SELECT * FROM " + table + " WHERE 1 = 0", rs -> {
|
||||
ResultSetMetaData metaData = rs.getMetaData();
|
||||
for (int i = 1; i <= metaData.getColumnCount(); i++) {
|
||||
columns.add(metaData.getColumnName(i).toUpperCase());
|
||||
}
|
||||
});
|
||||
return columns;
|
||||
}
|
||||
|
||||
/**
|
||||
* 分页读取源数据(Oracle/DM 均支持 ROWNUM 分页)
|
||||
*/
|
||||
private List<Map<String, Object>> readPage(JdbcTemplate sourceJdbc, DataSyncConfig config,
|
||||
List<String> sourceColumns, int offset, int batchSize) {
|
||||
String columnList = String.join(", ", sourceColumns);
|
||||
String orderBy = StrUtil.isBlank(config.getOrderColumn())
|
||||
? sourceColumns.get(0) : config.getOrderColumn().trim().toUpperCase();
|
||||
String where = StrUtil.isBlank(config.getConditionSql())
|
||||
? "" : " AND " + config.getConditionSql();
|
||||
String sql = "SELECT * FROM (SELECT t.*, ROWNUM rn FROM ("
|
||||
+ "SELECT " + columnList + " FROM " + config.getSourceTable()
|
||||
+ " WHERE 1 = 1" + where
|
||||
+ " ORDER BY " + orderBy
|
||||
+ ") t WHERE ROWNUM <= " + (offset + batchSize)
|
||||
+ ") WHERE rn > " + offset;
|
||||
return sourceJdbc.queryForList(sql);
|
||||
}
|
||||
|
||||
/**
|
||||
* 批量 MERGE INTO 写入目标库
|
||||
*
|
||||
* @param sourceColumns 源表列(USING 子句占位别名)
|
||||
* @param targetColumns 目标表列(与 sourceColumns 一一对应)
|
||||
* @param keyColumns 目标表主键列(ON 匹配 t 侧)
|
||||
* @param keySourceColumns 源表主键列(ON 匹配 s 侧,与 keyColumns 一一对应)
|
||||
*/
|
||||
private int[] mergeInto(JdbcTemplate targetJdbc, DataSyncConfig config,
|
||||
List<String> sourceColumns, List<String> targetColumns,
|
||||
List<String> keyColumns, List<String> keySourceColumns,
|
||||
List<Map<String, Object>> rows) {
|
||||
String targetTable = config.getTargetTable();
|
||||
|
||||
// USING 子句:SELECT ? AS SRC_COL1, ? AS SRC_COL2 FROM DUAL
|
||||
String usingSelect = sourceColumns.stream()
|
||||
.map(c -> "? AS " + c)
|
||||
.collect(Collectors.joining(", "));
|
||||
|
||||
// ON 匹配条件:t.KEY_TARGET = s.KEY_SOURCE
|
||||
List<String> onItems = new ArrayList<>(keyColumns.size());
|
||||
for (int i = 0; i < keyColumns.size(); i++) {
|
||||
onItems.add("t." + keyColumns.get(i) + " = s." + keySourceColumns.get(i));
|
||||
}
|
||||
String onCondition = String.join(" AND ", onItems);
|
||||
|
||||
// 更新列(排除主键列):t.TARGET_COL = s.SOURCE_COL
|
||||
List<String> updateSetItems = new ArrayList<>();
|
||||
for (int i = 0; i < sourceColumns.size(); i++) {
|
||||
String targetColumn = targetColumns.get(i);
|
||||
if (keyColumns.contains(targetColumn)) {
|
||||
continue;
|
||||
}
|
||||
updateSetItems.add("t." + targetColumn + " = s." + sourceColumns.get(i));
|
||||
}
|
||||
if (updateSetItems.isEmpty()) {
|
||||
throw new IllegalArgumentException("目标表除主键列外没有可更新的列,无法构造 MERGE 的 WHEN MATCHED 分支");
|
||||
}
|
||||
String updateSet = String.join(", ", updateSetItems);
|
||||
|
||||
// 插入列:INSERT (TARGET_COL1, TARGET_COL2) VALUES (s.SRC_COL1, s.SRC_COL2)
|
||||
String insertColumns = String.join(", ", targetColumns);
|
||||
String insertValues = sourceColumns.stream()
|
||||
.map(c -> "s." + c)
|
||||
.collect(Collectors.joining(", "));
|
||||
|
||||
String sql = "MERGE INTO " + targetTable + " t "
|
||||
+ "USING (SELECT " + usingSelect + " FROM DUAL) s "
|
||||
+ "ON (" + onCondition + ") "
|
||||
+ "WHEN MATCHED THEN UPDATE SET " + updateSet + " "
|
||||
+ "WHEN NOT MATCHED THEN INSERT (" + insertColumns + ") VALUES (" + insertValues + ")";
|
||||
|
||||
List<Object[]> batchArgs = new ArrayList<>(rows.size());
|
||||
for (Map<String, Object> row : rows) {
|
||||
Object[] args = new Object[sourceColumns.size()];
|
||||
for (int i = 0; i < sourceColumns.size(); i++) {
|
||||
args[i] = row.get(sourceColumns.get(i));
|
||||
}
|
||||
batchArgs.add(args);
|
||||
}
|
||||
return targetJdbc.batchUpdate(sql, batchArgs);
|
||||
}
|
||||
|
||||
/**
|
||||
* 清空目标表(TRUNCATE,DDL 自动提交)
|
||||
*/
|
||||
private void clearTargetTable(JdbcTemplate targetJdbc, String targetTable) {
|
||||
targetJdbc.execute("TRUNCATE TABLE " + targetTable);
|
||||
log.info("已清空目标表 {}", targetTable);
|
||||
}
|
||||
|
||||
// ==================== 配置与结果 ====================
|
||||
|
||||
/**
|
||||
* 数据库连接信息(调用方自定义连接时使用)
|
||||
*/
|
||||
public static class DbConnectionInfo {
|
||||
|
||||
/** JDBC URL,如 jdbc:dm://127.0.0.1:5236/DMDB 或 jdbc:oracle:thin:@127.0.0.1:1521/ORCL */
|
||||
private String url;
|
||||
|
||||
/** 用户名 */
|
||||
private String username;
|
||||
|
||||
/** 密码 */
|
||||
private String password;
|
||||
|
||||
/** 驱动类名,不填时按 URL 自动识别(dm.jdbc.driver.DmDriver / oracle.jdbc.OracleDriver) */
|
||||
private String driverClassName;
|
||||
|
||||
public String getUrl() {
|
||||
return url;
|
||||
}
|
||||
|
||||
public void setUrl(String url) {
|
||||
this.url = url;
|
||||
}
|
||||
|
||||
public String getUsername() {
|
||||
return username;
|
||||
}
|
||||
|
||||
public void setUsername(String username) {
|
||||
this.username = username;
|
||||
}
|
||||
|
||||
public String getPassword() {
|
||||
return password;
|
||||
}
|
||||
|
||||
public void setPassword(String password) {
|
||||
this.password = password;
|
||||
}
|
||||
|
||||
public String getDriverClassName() {
|
||||
return driverClassName;
|
||||
}
|
||||
|
||||
public void setDriverClassName(String driverClassName) {
|
||||
this.driverClassName = driverClassName;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 数据同步配置
|
||||
*/
|
||||
public static class DataSyncConfig {
|
||||
|
||||
/** 源表名(DM/Oracle) */
|
||||
private String sourceTable;
|
||||
|
||||
/** 目标表名(可与源表名不同) */
|
||||
private String targetTable;
|
||||
|
||||
/** 源表同步列(不指定则自动读取源表全列) */
|
||||
private List<String> columns;
|
||||
|
||||
/**
|
||||
* 字段映射:源列名 -> 目标列名。
|
||||
* 未在映射中的列默认目标列名与源列名一致;
|
||||
* 若两边字段名完全相同,可不配置。
|
||||
*/
|
||||
private Map<String, String> columnMapping;
|
||||
|
||||
/** MERGE 匹配主键列(目标表列名,必填) */
|
||||
private List<String> keyColumns;
|
||||
|
||||
/** 增量同步条件 SQL 片段(基于源表列),如 MODIFY_TIME >= TO_DATE('2025-01-01','YYYY-MM-DD') */
|
||||
private String conditionSql;
|
||||
|
||||
/** 分页排序列(源表列,不指定则取第一列) */
|
||||
private String orderColumn;
|
||||
|
||||
/** 每批同步条数,默认 500 */
|
||||
private Integer batchSize;
|
||||
|
||||
/** 同步前是否清空目标表(全量覆盖) */
|
||||
private Boolean clearTarget;
|
||||
|
||||
public String getSourceTable() {
|
||||
return sourceTable;
|
||||
}
|
||||
|
||||
public void setSourceTable(String sourceTable) {
|
||||
this.sourceTable = sourceTable;
|
||||
}
|
||||
|
||||
public String getTargetTable() {
|
||||
return targetTable;
|
||||
}
|
||||
|
||||
public void setTargetTable(String targetTable) {
|
||||
this.targetTable = targetTable;
|
||||
}
|
||||
|
||||
public List<String> getColumns() {
|
||||
return columns;
|
||||
}
|
||||
|
||||
public void setColumns(List<String> columns) {
|
||||
this.columns = columns;
|
||||
}
|
||||
|
||||
public Map<String, String> getColumnMapping() {
|
||||
return columnMapping;
|
||||
}
|
||||
|
||||
public void setColumnMapping(Map<String, String> columnMapping) {
|
||||
this.columnMapping = columnMapping;
|
||||
}
|
||||
|
||||
public List<String> getKeyColumns() {
|
||||
return keyColumns;
|
||||
}
|
||||
|
||||
public void setKeyColumns(List<String> keyColumns) {
|
||||
this.keyColumns = keyColumns;
|
||||
}
|
||||
|
||||
public String getConditionSql() {
|
||||
return conditionSql;
|
||||
}
|
||||
|
||||
public void setConditionSql(String conditionSql) {
|
||||
this.conditionSql = conditionSql;
|
||||
}
|
||||
|
||||
public String getOrderColumn() {
|
||||
return orderColumn;
|
||||
}
|
||||
|
||||
public void setOrderColumn(String orderColumn) {
|
||||
this.orderColumn = orderColumn;
|
||||
}
|
||||
|
||||
public Integer getBatchSize() {
|
||||
return batchSize;
|
||||
}
|
||||
|
||||
public void setBatchSize(Integer batchSize) {
|
||||
this.batchSize = batchSize;
|
||||
}
|
||||
|
||||
public Boolean getClearTarget() {
|
||||
return clearTarget;
|
||||
}
|
||||
|
||||
public void setClearTarget(Boolean clearTarget) {
|
||||
this.clearTarget = clearTarget;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 数据同步结果
|
||||
*/
|
||||
public static class DataSyncResult {
|
||||
|
||||
/** 是否成功 */
|
||||
private boolean success;
|
||||
|
||||
/** 读取记录数 */
|
||||
private long totalRead;
|
||||
|
||||
/** MERGE 处理记录数(含新增+更新) */
|
||||
private long totalProcessed;
|
||||
|
||||
/** 开始时间 */
|
||||
private long startTime;
|
||||
|
||||
/** 结束时间 */
|
||||
private long endTime;
|
||||
|
||||
/** 耗时(毫秒) */
|
||||
private long elapsedMillis;
|
||||
|
||||
/** 结果描述 */
|
||||
private String message;
|
||||
|
||||
public boolean isSuccess() {
|
||||
return success;
|
||||
}
|
||||
|
||||
public void setSuccess(boolean success) {
|
||||
this.success = success;
|
||||
}
|
||||
|
||||
public long getTotalRead() {
|
||||
return totalRead;
|
||||
}
|
||||
|
||||
public void setTotalRead(long totalRead) {
|
||||
this.totalRead = totalRead;
|
||||
}
|
||||
|
||||
public long getTotalProcessed() {
|
||||
return totalProcessed;
|
||||
}
|
||||
|
||||
public void setTotalProcessed(long totalProcessed) {
|
||||
this.totalProcessed = totalProcessed;
|
||||
}
|
||||
|
||||
public long getStartTime() {
|
||||
return startTime;
|
||||
}
|
||||
|
||||
public void setStartTime(long startTime) {
|
||||
this.startTime = startTime;
|
||||
}
|
||||
|
||||
public long getEndTime() {
|
||||
return endTime;
|
||||
}
|
||||
|
||||
public void setEndTime(long endTime) {
|
||||
this.endTime = endTime;
|
||||
}
|
||||
|
||||
public long getElapsedMillis() {
|
||||
return elapsedMillis;
|
||||
}
|
||||
|
||||
public void setElapsedMillis(long elapsedMillis) {
|
||||
this.elapsedMillis = elapsedMillis;
|
||||
}
|
||||
|
||||
public String getMessage() {
|
||||
return message;
|
||||
}
|
||||
|
||||
public void setMessage(String message) {
|
||||
this.message = message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "DataSyncResult{" +
|
||||
"success=" + success +
|
||||
", totalRead=" + totalRead +
|
||||
", totalProcessed=" + totalProcessed +
|
||||
", elapsedMillis=" + elapsedMillis +
|
||||
", message='" + message + '\'' +
|
||||
'}';
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user