【feat】 优化元数据变更同步本地方法

This commit is contained in:
Kris 2026-01-03 21:28:13 +08:00
parent decbebb4de
commit 18d52e2ad1
7 changed files with 567 additions and 182 deletions

View File

@ -198,6 +198,17 @@ public class DataiIntegrationMetadataChangeController extends BaseController
/** /**
* 同步元数据变更到本地数据库 * 同步元数据变更到本地数据库
* 根据元数据变更ID将指定的元数据变更同步到本地数据库
*
* 该方法会
* 1. 根据ID查询元数据变更记录
* 2. 根据变更类型OBJECT或FIELD执行相应的同步操作
* 3. 对于对象变更执行对象的创建修改或删除操作
* 4. 对于字段变更执行字段的创建修改或删除操作
* 5. 更新元数据变更记录的同步状态
*
* @param id 元数据变更ID
* @return 同步结果包含success是否成功和message消息字段
*/ */
@Operation(summary = "同步元数据变更到本地数据库") @Operation(summary = "同步元数据变更到本地数据库")
@PreAuthorize("@ss.hasPermi('integration:change:sync')") @PreAuthorize("@ss.hasPermi('integration:change:sync')")

View File

@ -77,6 +77,44 @@ public interface CustomMapper {
*/ */
int createField(@Param("tableName") String tableName, @Param("fieldName") String fieldName); int createField(@Param("tableName") String tableName, @Param("fieldName") String fieldName);
/**
* 添加字段通用方法支持字段类型和约束
*
* @param tableName 表名
* @param fieldName 字段名
* @param fieldType 字段类型
* @param isNullable 是否可为空
* @return 影响行数
*/
int addField(@Param("tableName") String tableName,
@Param("fieldName") String fieldName,
@Param("fieldType") String fieldType,
@Param("isNullable") Boolean isNullable);
/**
* 修改字段通用方法支持字段类型和约束
*
* @param tableName 表名
* @param fieldName 字段名
* @param fieldType 字段类型
* @param isNullable 是否可为空
* @return 影响行数
*/
int modifyField(@Param("tableName") String tableName,
@Param("fieldName") String fieldName,
@Param("fieldType") String fieldType,
@Param("isNullable") Boolean isNullable);
/**
* 删除字段通用方法
*
* @param tableName 表名
* @param fieldName 字段名
* @return 影响行数
*/
int dropField(@Param("tableName") String tableName,
@Param("fieldName") String fieldName);
/** /**
* 根据ID更新记录 * 根据ID更新记录
* *

View File

@ -86,4 +86,16 @@ public interface DataiIntegrationMetadataChangeMapper
* @return 统计信息 * @return 统计信息
*/ */
public Map<String, Object> selectChangeStatistics(Map<String, Object> params); public Map<String, Object> selectChangeStatistics(Map<String, Object> params);
/**
* 查询是否存在相同的元数据变更记录
*
* @param changeType 变更类型
* @param operationType 操作类型
* @param objectApi 对象API
* @param fieldApi 字段API可为null
* @param changeReason 变更原因
* @return 存在的记录数
*/
public int countSimilarChanges(String changeType, String operationType, String objectApi, String fieldApi, String changeReason);
} }

View File

@ -88,9 +88,17 @@ public interface IDataiIntegrationMetadataChangeService
/** /**
* 同步元数据变更到本地数据库 * 同步元数据变更到本地数据库
* 根据元数据变更ID将指定的元数据变更同步到本地数据库
*
* 该方法会
* 1. 根据ID查询元数据变更记录
* 2. 根据变更类型OBJECT或FIELD执行相应的同步操作
* 3. 对于对象变更执行对象的创建修改或删除操作
* 4. 对于字段变更执行字段的创建修改或删除操作
* 5. 更新元数据变更记录的同步状态
* *
* @param id 元数据变更ID * @param id 元数据变更ID
* @return 同步结果 * @return 同步结果包含success是否成功和message消息字段
*/ */
public Map<String, Object> syncToLocalDatabase(Long id); public Map<String, Object> syncToLocalDatabase(Long id);

View File

@ -8,6 +8,7 @@ import com.datai.common.utils.SecurityUtils;
import com.datai.integration.factory.impl.SOAPConnectionFactory; import com.datai.integration.factory.impl.SOAPConnectionFactory;
import com.datai.integration.mapper.CustomMapper; import com.datai.integration.mapper.CustomMapper;
import com.sforce.soap.partner.*; import com.sforce.soap.partner.*;
import com.sforce.soap.partner.sobject.SObject;
import com.sforce.ws.ConnectionException; import com.sforce.ws.ConnectionException;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -20,6 +21,8 @@ import com.datai.integration.model.domain.DataiIntegrationField;
import com.datai.integration.service.IDataiIntegrationMetadataChangeService; import com.datai.integration.service.IDataiIntegrationMetadataChangeService;
import com.datai.integration.service.IDataiIntegrationObjectService; import com.datai.integration.service.IDataiIntegrationObjectService;
import com.datai.integration.service.IDataiIntegrationFieldService; import com.datai.integration.service.IDataiIntegrationFieldService;
import com.datai.integration.service.IDataiIntegrationPicklistService;
import com.datai.integration.service.IDataiIntegrationFilterLookupService;
import com.datai.common.core.domain.model.LoginUser; import com.datai.common.core.domain.model.LoginUser;
/** /**
@ -31,6 +34,11 @@ import com.datai.common.core.domain.model.LoginUser;
@Service @Service
@Slf4j @Slf4j
public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrationMetadataChangeService { public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrationMetadataChangeService {
/**
* 大数据量对象阈值500万
*/
private static final int LARGE_OBJECT_THRESHOLD = 5000000;
@Autowired @Autowired
private DataiIntegrationMetadataChangeMapper dataiIntegrationMetadataChangeMapper; private DataiIntegrationMetadataChangeMapper dataiIntegrationMetadataChangeMapper;
@ -40,6 +48,12 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
@Autowired @Autowired
private IDataiIntegrationFieldService dataiIntegrationFieldService; private IDataiIntegrationFieldService dataiIntegrationFieldService;
@Autowired
private IDataiIntegrationPicklistService dataiIntegrationPicklistService;
@Autowired
private IDataiIntegrationFilterLookupService dataiIntegrationFilterLookupService;
@Autowired @Autowired
private SOAPConnectionFactory soapConnectionFactory; private SOAPConnectionFactory soapConnectionFactory;
@ -196,6 +210,7 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
Map<String, Object> result = new HashMap<>(); Map<String, Object> result = new HashMap<>();
try { try {
// 根据ID查询元数据变更记录
DataiIntegrationMetadataChange metadataChange = selectDataiIntegrationMetadataChangeById(id); DataiIntegrationMetadataChange metadataChange = selectDataiIntegrationMetadataChangeById(id);
if (metadataChange == null) { if (metadataChange == null) {
result.put("success", false); result.put("success", false);
@ -203,12 +218,16 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
return result; return result;
} }
// 获取变更类型和操作类型
String changeType = metadataChange.getChangeType(); String changeType = metadataChange.getChangeType();
String operationType = metadataChange.getOperationType(); String operationType = metadataChange.getOperationType();
// 根据变更类型执行相应的同步操作
if ("OBJECT".equals(changeType)) { if ("OBJECT".equals(changeType)) {
// 对象级别的变更同步
syncObjectChange(metadataChange, result); syncObjectChange(metadataChange, result);
} else if ("FIELD".equals(changeType)) { } else if ("FIELD".equals(changeType)) {
// 字段级别的变更同步
syncFieldChange(metadataChange, result); syncFieldChange(metadataChange, result);
} else { } else {
result.put("success", false); result.put("success", false);
@ -216,13 +235,17 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
return result; return result;
} }
// 根据同步结果更新元数据变更记录的同步状态
if ((Boolean) result.get("success")) { if ((Boolean) result.get("success")) {
// 同步成功更新同步状态为true
updateSyncStatus(metadataChange.getId(), true, null); updateSyncStatus(metadataChange.getId(), true, null);
} else { } else {
// 同步失败更新同步状态为false并记录错误信息
updateSyncStatus(metadataChange.getId(), false, (String) result.get("message")); updateSyncStatus(metadataChange.getId(), false, (String) result.get("message"));
} }
} catch (Exception e) { } catch (Exception e) {
// 捕获异常并返回错误信息
result.put("success", false); result.put("success", false);
result.put("message", "同步失败: " + e.getMessage()); result.put("message", "同步失败: " + e.getMessage());
updateSyncStatus(id, false, e.getMessage()); updateSyncStatus(id, false, e.getMessage());
@ -274,19 +297,53 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
try { try {
if ("INSERT".equals(operationType)) { if ("INSERT".equals(operationType)) {
DataiIntegrationObject object = new DataiIntegrationObject(); try {
object.setApi(objectApi); log.info("开始同步对象: {}", objectApi);
object.setLabel(metadataChange.getObjectLabel());
object.setIsCustom(metadataChange.getIsCustom()); PartnerConnection connection = retryOperation(() -> soapConnectionFactory.getConnection(), 3, 1000);
object.setCreateTime(DateUtils.getNowDate()); log.info("成功获取源ORG连接");
object.setUpdateTime(DateUtils.getNowDate());
DescribeSObjectResult objDetail = connection.describeSObject(objectApi.trim());
DataiIntegrationObject object = buildObjectMetadata(objDetail);
if (object == null) {
log.error("构建对象 {} 的元数据失败", objectApi);
result.put("success", false);
result.put("message", "构建对象元数据失败");
return;
}
int insertResult = dataiIntegrationObjectService.insertDataiIntegrationObject(object); int insertResult = dataiIntegrationObjectService.insertDataiIntegrationObject(object);
result.put("success", insertResult > 0); result.put("success", insertResult > 0);
result.put("message", insertResult > 0 ? "对象创建成功" : "对象创建失败"); result.put("message", insertResult > 0 ? "对象创建成功" : "对象创建失败");
if (insertResult > 0) { if (insertResult > 0) {
saveObjectFieldsToDataiIntegrationField(objDetail);
int objectNum = isLargeObject(connection, objectApi.trim());
log.info("对象 {} 的数据量: {}", objectApi, objectNum);
updateDataiIntegrationObjectFields(objectApi.trim(), objectNum);
if (objectNum > LARGE_OBJECT_THRESHOLD) {
log.info("对象 {} 是大数据量对象,数据量大于五百万,创建分区表", objectApi);
createOrUpdateDatabaseTable(objectApi, metadataChange.getObjectLabel(), result); createOrUpdateDatabaseTable(objectApi, metadataChange.getObjectLabel(), result);
} else {
log.info("对象 {} 是普通对象,数据量少于五百万,创建正常表", objectApi);
createOrUpdateDatabaseTable(objectApi, metadataChange.getObjectLabel(), result);
}
log.info("同步对象 {} 成功", objectApi);
}
} catch (ConnectionException e) {
log.error("获取对象 {} 的元数据失败: {}", objectApi, e.getMessage(), e);
result.put("success", false);
result.put("message", "获取对象元数据失败: " + e.getMessage());
} catch (Exception e) {
log.error("处理对象 {} 时出错: {}", objectApi, e.getMessage(), e);
result.put("success", false);
result.put("message", "处理对象时出错: " + e.getMessage());
} }
} else if ("UPDATE".equals(operationType)) { } else if ("UPDATE".equals(operationType)) {
@ -303,16 +360,14 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
DataiIntegrationObject object = existingObjects.get(0); DataiIntegrationObject object = existingObjects.get(0);
object.setLabel(metadataChange.getObjectLabel()); object.setLabel(metadataChange.getObjectLabel());
object.setIsWork(false);
object.setIsIncremental(false);
object.setUpdateTime(DateUtils.getNowDate()); object.setUpdateTime(DateUtils.getNowDate());
int updateResult = dataiIntegrationObjectService.updateDataiIntegrationObject(object); int updateResult = dataiIntegrationObjectService.updateDataiIntegrationObject(object);
result.put("success", updateResult > 0); result.put("success", updateResult > 0);
result.put("message", updateResult > 0 ? "对象更新成功" : "对象更新失败"); result.put("message", updateResult > 0 ? "对象更新成功" : "对象更新失败");
if (updateResult > 0) {
updateDatabaseTable(objectApi, metadataChange.getObjectLabel(), result);
}
} else if ("DELETE".equals(operationType)) { } else if ("DELETE".equals(operationType)) {
DataiIntegrationObject queryObject = new DataiIntegrationObject(); DataiIntegrationObject queryObject = new DataiIntegrationObject();
queryObject.setApi(objectApi); queryObject.setApi(objectApi);
@ -321,21 +376,18 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
if (existingObjects.isEmpty()) { if (existingObjects.isEmpty()) {
result.put("success", false); result.put("success", false);
result.put("message", "对象不存在,无法删除"); result.put("message", "对象不存在,无法更新");
return; return;
} }
Integer[] ids = existingObjects.stream() DataiIntegrationObject object = existingObjects.get(0);
.map(DataiIntegrationObject::getId) object.setIsWork(false);
.toArray(Integer[]::new); object.setIsIncremental(false);
object.setUpdateTime(DateUtils.getNowDate());
int deleteResult = dataiIntegrationObjectService.deleteDataiIntegrationObjectByIds(ids); int updateResult = dataiIntegrationObjectService.updateDataiIntegrationObject(object);
result.put("success", deleteResult > 0); result.put("success", updateResult > 0);
result.put("message", deleteResult > 0 ? "对象删除成功" : "对象删除失败"); result.put("message", updateResult > 0 ? "对象更新成功" : "对象更新失败");
if (deleteResult > 0) {
dropDatabaseTable(objectApi, result);
}
} else { } else {
result.put("success", false); result.put("success", false);
@ -501,46 +553,46 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
} }
} }
private void updateDatabaseTable(String objectApi, String objectLabel, Map<String, Object> result) { // private void updateDatabaseTable(String objectApi, String objectLabel, Map<String, Object> result) {
PartnerConnection connection = null; // PartnerConnection connection = null;
try { // try {
connection = soapConnectionFactory.getConnection(); // connection = soapConnectionFactory.getConnection();
DescribeSObjectResult objDetail = connection.describeSObject(objectApi); // DescribeSObjectResult objDetail = connection.describeSObject(objectApi);
//
if (objDetail == null) { // if (objDetail == null) {
result.put("success", false); // result.put("success", false);
result.put("message", "无法获取对象元数据"); // result.put("message", "无法获取对象元数据");
return; // return;
} // }
//
List<Map<String, Object>> fieldMaps = new ArrayList<>(); // List<Map<String, Object>> fieldMaps = new ArrayList<>();
//
for (com.sforce.soap.partner.Field field : objDetail.getFields()) { // for (com.sforce.soap.partner.Field field : objDetail.getFields()) {
Map<String, Object> fieldMap = new HashMap<>(); // Map<String, Object> fieldMap = new HashMap<>();
fieldMap.put("fieldName", field.getName()); // fieldMap.put("fieldName", field.getName());
fieldMap.put("fieldType", convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null)); // fieldMap.put("fieldType", convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null));
fieldMap.put("fieldLength", field.getLength()); // fieldMap.put("fieldLength", field.getLength());
fieldMap.put("fieldPrecision", field.getPrecision()); // fieldMap.put("fieldPrecision", field.getPrecision());
fieldMap.put("fieldScale", field.getScale()); // fieldMap.put("fieldScale", field.getScale());
fieldMap.put("isNullable", field.isNillable()); // fieldMap.put("isNullable", field.isNillable());
fieldMap.put("isUnique", field.isUnique()); // fieldMap.put("isUnique", field.isUnique());
fieldMap.put("isPrimaryKey", field.isIdLookup()); // fieldMap.put("isPrimaryKey", field.isIdLookup());
fieldMaps.add(fieldMap); // fieldMaps.add(fieldMap);
} // }
//
customMapper.createTable(objectApi, objectLabel, fieldMaps, new ArrayList<>()); // customMapper.createTable(objectApi, objectLabel, fieldMaps, new ArrayList<>());
log.info("成功更新表结构: {}", objectApi); // log.info("成功更新表结构: {}", objectApi);
//
} catch (ConnectionException e) { // } catch (ConnectionException e) {
result.put("success", false); // result.put("success", false);
result.put("message", "获取Salesforce连接失败: " + e.getMessage()); // result.put("message", "获取Salesforce连接失败: " + e.getMessage());
log.error("获取Salesforce连接失败: {}", e.getMessage(), e); // log.error("获取Salesforce连接失败: {}", e.getMessage(), e);
} catch (Exception e) { // } catch (Exception e) {
result.put("success", false); // result.put("success", false);
result.put("message", "更新表结构失败: " + e.getMessage()); // result.put("message", "更新表结构失败: " + e.getMessage());
log.error("更新表结构失败: {}", e.getMessage(), e); // log.error("更新表结构失败: {}", e.getMessage(), e);
} // }
} // }
private void dropDatabaseTable(String objectApi, Map<String, Object> result) { private void dropDatabaseTable(String objectApi, Map<String, Object> result) {
try { try {
@ -568,8 +620,10 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
for (com.sforce.soap.partner.Field field : objDetail.getFields()) { for (com.sforce.soap.partner.Field field : objDetail.getFields()) {
if (fieldApi.equals(field.getName())) { if (fieldApi.equals(field.getName())) {
customMapper.createField(objectApi, fieldApi); String mysqlType = convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null);
log.info("成功添加字段: {}.{}", objectApi, fieldApi);
customMapper.addField(objectApi, fieldApi, mysqlType, field.isNillable());
log.info("成功添加字段: {}.{} 类型: {}", objectApi, fieldApi, mysqlType);
return; return;
} }
} }
@ -603,10 +657,9 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
for (com.sforce.soap.partner.Field field : objDetail.getFields()) { for (com.sforce.soap.partner.Field field : objDetail.getFields()) {
if (fieldApi.equals(field.getName())) { if (fieldApi.equals(field.getName())) {
String mysqlType = convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null); String mysqlType = convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null);
String alterSql = String.format("ALTER TABLE %s MODIFY COLUMN %s %s",
objectApi, fieldApi, mysqlType); customMapper.modifyField(objectApi, fieldApi, mysqlType, field.isNillable());
customMapper.executeUpdate(alterSql); log.info("成功修改字段: {}.{} 类型: {}", objectApi, fieldApi, mysqlType);
log.info("成功修改字段: {}.{}", objectApi, fieldApi);
return; return;
} }
} }
@ -627,8 +680,7 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
private void dropDatabaseColumn(String objectApi, String fieldApi, Map<String, Object> result) { private void dropDatabaseColumn(String objectApi, String fieldApi, Map<String, Object> result) {
try { try {
String alterSql = String.format("ALTER TABLE %s DROP COLUMN %s", objectApi, fieldApi); customMapper.dropField(objectApi, fieldApi);
customMapper.executeUpdate(alterSql);
log.info("成功删除字段: {}.{}", objectApi, fieldApi); log.info("成功删除字段: {}.{}", objectApi, fieldApi);
} catch (Exception e) { } catch (Exception e) {
result.put("success", false); result.put("success", false);
@ -894,50 +946,6 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
* @param objDetail Salesforce对象的详细描述信息 * @param objDetail Salesforce对象的详细描述信息
* @return 构建的DataiIntegrationObject实体如果构建失败返回null * @return 构建的DataiIntegrationObject实体如果构建失败返回null
*/ */
private DataiIntegrationObject buildObjectMetadata(DescribeSObjectResult objDetail) {
try {
// 创建对象元数据实体
DataiIntegrationObject object = new DataiIntegrationObject();
// 设置对象API名称
object.setApi(objDetail.getName());
// 设置对象标签名称
object.setLabel(objDetail.getLabel());
// 设置对象复数标签名称
object.setLabelPlural(objDetail.getLabelPlural());
// 设置对象键前缀
object.setKeyPrefix(objDetail.getKeyPrefix());
// 设置是否为自定义对象
object.setIsCustom(objDetail.isCustom());
// 设置是否为自定义设置对象
object.setIsCustomSetting(objDetail.isCustomSetting());
// 设置是否可查询
object.setIsQueryable(objDetail.isQueryable());
// 设置是否可创建
object.setIsCreateable(objDetail.isCreateable());
// 设置是否可更新
object.setIsUpdateable(objDetail.isUpdateable());
// 设置是否可删除
object.setIsDeletable(objDetail.isDeletable());
// 设置是否可复制
object.setIsReplicateable(objDetail.isReplicateable());
// 设置是否可检索
object.setIsRetrieveable(objDetail.isRetrieveable());
// 设置是否可搜索
object.setIsSearchable(objDetail.isSearchable());
// 设置是否工作状态启用
object.setIsWork(false);
// 设置是否启用增量更新
object.setIsIncremental(false);
// 设置最后同步时间
object.setLastSyncDate(LocalDateTime.now());
return object;
} catch (Exception e) {
// 记录构建对象元数据时的错误
log.error("构建对象 {} 元数据时出错: {}", objDetail.getName(), e.getMessage(), e);
return null;
}
}
/** /**
* 根据Salesforce字段描述信息构建字段元数据实体 * 根据Salesforce字段描述信息构建字段元数据实体
* 将Salesforce的Field对象转换为DataiIntegrationField实体 * 将Salesforce的Field对象转换为DataiIntegrationField实体
@ -1249,45 +1257,40 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
*/ */
private void recordObjectChange(DataiIntegrationObject newObject, DataiIntegrationObject oldObject, String operationType) { private void recordObjectChange(DataiIntegrationObject newObject, DataiIntegrationObject oldObject, String operationType) {
try { try {
// 创建元数据变更记录实体 String changeReason;
DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange();
// 设置变更类型为对象
metadataChange.setChangeType("OBJECT");
// 设置操作类型INSERTUPDATEDELETE
metadataChange.setOperationType(operationType);
// 设置对象API名称
metadataChange.setObjectApi(newObject.getApi());
// 设置对象标签名称
metadataChange.setObjectLabel(newObject.getLabel());
// 设置变更时间
metadataChange.setChangeTime(LocalDateTime.now());
// 设置同步状态为未同步
metadataChange.setSyncStatus(false);
// 设置是否为自定义对象
metadataChange.setIsCustom(newObject.getIsCustom());
// 如果是更新操作oldObject不为null则比较新旧对象的差异
if (oldObject != null) { if (oldObject != null) {
List<String> changedFields = compareObjects(oldObject, newObject); List<String> changedFields = compareObjects(oldObject, newObject);
if (!changedFields.isEmpty()) { if (!changedFields.isEmpty()) {
// 如果有具体变更字段则记录详细变更信息 changeReason = "对象属性变更: " + String.join(", ", changedFields);
metadataChange.setChangeReason("对象属性变更: " + String.join(", ", changedFields));
} else { } else {
// 如果没有具体变更字段则记录通用更新信息 changeReason = "对象属性更新";
metadataChange.setChangeReason("对象属性更新");
} }
} else { } else {
// 如果是新增操作则记录为新增对象 changeReason = "新增对象";
metadataChange.setChangeReason("新增对象");
} }
// 设置变更用户为系统
int similarCount = dataiIntegrationMetadataChangeMapper.countSimilarChanges(
"OBJECT", operationType, newObject.getApi(), null, changeReason);
if (similarCount > 0) {
log.debug("发现相似的未同步对象变更记录,跳过重复记录: {} - {}", newObject.getApi(), operationType);
return;
}
DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange();
metadataChange.setChangeType("OBJECT");
metadataChange.setOperationType(operationType);
metadataChange.setObjectApi(newObject.getApi());
metadataChange.setObjectLabel(newObject.getLabel());
metadataChange.setChangeTime(LocalDateTime.now());
metadataChange.setSyncStatus(false);
metadataChange.setIsCustom(newObject.getIsCustom());
metadataChange.setChangeReason(changeReason);
metadataChange.setChangeUser("SYSTEM"); metadataChange.setChangeUser("SYSTEM");
// 插入元数据变更记录
insertDataiIntegrationMetadataChange(metadataChange); insertDataiIntegrationMetadataChange(metadataChange);
log.debug("记录对象变更成功: {} - {}", newObject.getApi(), operationType); log.debug("记录对象变更成功: {} - {}", newObject.getApi(), operationType);
} catch (Exception e) { } catch (Exception e) {
// 记录记录对象变更时的错误
log.error("记录对象变更失败: {} - {}", newObject.getApi(), e.getMessage(), e); log.error("记录对象变更失败: {} - {}", newObject.getApi(), e.getMessage(), e);
} }
} }
@ -1308,36 +1311,32 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
private void recordFieldChange(String objectApi, String objectLabel, String fieldApi, String fieldLabel, private void recordFieldChange(String objectApi, String objectLabel, String fieldApi, String fieldLabel,
String changeReason, String operationType, String defaultReason, Boolean isCustom) { String changeReason, String operationType, String defaultReason, Boolean isCustom) {
try { try {
// 创建元数据变更记录实体 String finalChangeReason = changeReason != null ? changeReason : defaultReason;
int similarCount = dataiIntegrationMetadataChangeMapper.countSimilarChanges(
"FIELD", operationType, objectApi, fieldApi, finalChangeReason);
if (similarCount > 0) {
log.debug("发现相似的未同步字段变更记录,跳过重复记录: {}.{} - {}", objectApi, fieldApi, operationType);
return;
}
DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange(); DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange();
// 设置变更类型为字段
metadataChange.setChangeType("FIELD"); metadataChange.setChangeType("FIELD");
// 设置操作类型INSERTUPDATEDELETE
metadataChange.setOperationType(operationType); metadataChange.setOperationType(operationType);
// 设置对象API名称
metadataChange.setObjectApi(objectApi); metadataChange.setObjectApi(objectApi);
// 设置对象标签名称
metadataChange.setObjectLabel(objectLabel); metadataChange.setObjectLabel(objectLabel);
// 设置字段API名称
metadataChange.setFieldApi(fieldApi); metadataChange.setFieldApi(fieldApi);
// 设置字段标签名称
metadataChange.setFieldLabel(fieldLabel); metadataChange.setFieldLabel(fieldLabel);
// 设置变更时间
metadataChange.setChangeTime(LocalDateTime.now()); metadataChange.setChangeTime(LocalDateTime.now());
// 设置同步状态为未同步
metadataChange.setSyncStatus(false); metadataChange.setSyncStatus(false);
// 设置是否为自定义对象
metadataChange.setIsCustom(isCustom); metadataChange.setIsCustom(isCustom);
// 设置变更原因如果changeReason不为null则使用它否则使用默认原因 metadataChange.setChangeReason(finalChangeReason);
metadataChange.setChangeReason(changeReason != null ? changeReason : defaultReason);
// 设置变更用户为系统
metadataChange.setChangeUser("SYSTEM"); metadataChange.setChangeUser("SYSTEM");
// 插入元数据变更记录
insertDataiIntegrationMetadataChange(metadataChange); insertDataiIntegrationMetadataChange(metadataChange);
log.debug("记录字段变更成功: {}.{} - {}", objectApi, fieldApi, operationType); log.debug("记录字段变更成功: {}.{} - {}", objectApi, fieldApi, operationType);
} catch (Exception e) { } catch (Exception e) {
// 记录记录字段变更时的错误
log.error("记录字段变更失败: {}.{} - {}", objectApi, fieldApi, e.getMessage(), e); log.error("记录字段变更失败: {}.{} - {}", objectApi, fieldApi, e.getMessage(), e);
} }
} }
@ -1351,41 +1350,37 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
*/ */
private void checkDeletedObjectsForMetadata(Set<String> syncedObjectApis) { private void checkDeletedObjectsForMetadata(Set<String> syncedObjectApis) {
try { try {
// 查询数据库中所有对象
DataiIntegrationObject queryObject = new DataiIntegrationObject(); DataiIntegrationObject queryObject = new DataiIntegrationObject();
List<DataiIntegrationObject> allObjects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject); List<DataiIntegrationObject> allObjects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject);
// 遍历数据库中的所有对象检查哪些对象已不在Salesforce中
for (DataiIntegrationObject object : allObjects) { for (DataiIntegrationObject object : allObjects) {
if (!syncedObjectApis.contains(object.getApi())) { if (!syncedObjectApis.contains(object.getApi())) {
// 创建元数据变更记录标记为删除操作 String changeReason = "对象已从Salesforce中删除";
int similarCount = dataiIntegrationMetadataChangeMapper.countSimilarChanges(
"OBJECT", "DELETE", object.getApi(), null, changeReason);
if (similarCount > 0) {
log.debug("发现相似的未同步对象删除记录,跳过重复记录: {}", object.getApi());
continue;
}
DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange(); DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange();
// 设置变更类型为对象
metadataChange.setChangeType("OBJECT"); metadataChange.setChangeType("OBJECT");
// 设置操作类型为删除
metadataChange.setOperationType("DELETE"); metadataChange.setOperationType("DELETE");
// 设置对象API名称
metadataChange.setObjectApi(object.getApi()); metadataChange.setObjectApi(object.getApi());
// 设置对象标签名称
metadataChange.setObjectLabel(object.getLabel()); metadataChange.setObjectLabel(object.getLabel());
// 设置变更时间
metadataChange.setChangeTime(LocalDateTime.now()); metadataChange.setChangeTime(LocalDateTime.now());
// 设置同步状态为未同步
metadataChange.setSyncStatus(false); metadataChange.setSyncStatus(false);
// 设置是否为自定义对象
metadataChange.setIsCustom(object.getIsCustom()); metadataChange.setIsCustom(object.getIsCustom());
// 设置变更原因 metadataChange.setChangeReason(changeReason);
metadataChange.setChangeReason("对象已从Salesforce中删除");
// 设置变更用户为系统
metadataChange.setChangeUser("SYSTEM"); metadataChange.setChangeUser("SYSTEM");
// 插入元数据变更记录
insertDataiIntegrationMetadataChange(metadataChange); insertDataiIntegrationMetadataChange(metadataChange);
log.warn("检测到对象已删除: {}", object.getApi()); log.warn("检测到对象已删除: {}", object.getApi());
} }
} }
} catch (Exception e) { } catch (Exception e) {
// 记录检查已删除对象时的错误
log.error("检查已删除对象时出错: {}", e.getMessage(), e); log.error("检查已删除对象时出错: {}", e.getMessage(), e);
} }
} }
@ -1429,4 +1424,291 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat
} }
} }
} }
/**
* 构建对象元数据
*
* @param objDetail 对象详情
* @return 构建的对象元数据
*/
private DataiIntegrationObject buildObjectMetadata(DescribeSObjectResult objDetail) {
try {
DataiIntegrationObject object = new DataiIntegrationObject();
// 基本信息
object.setApi(objDetail.getName());
object.setLabel(objDetail.getLabel());
object.setLabelPlural(objDetail.getLabelPlural());
object.setKeyPrefix(objDetail.getKeyPrefix());
// 对象属性
object.setIsCustom(objDetail.isCustom());
object.setIsCustomSetting(objDetail.isCustomSetting());
// 对象权限
object.setIsQueryable(objDetail.isQueryable());
object.setIsCreateable(objDetail.isCreateable());
object.setIsUpdateable(objDetail.isUpdateable());
object.setIsDeletable(objDetail.isDeletable());
object.setIsReplicateable(objDetail.isReplicateable());
object.setIsRetrieveable(objDetail.isRetrieveable());
object.setIsSearchable(objDetail.isSearchable());
// 同步设置
object.setIsWork(true);
object.setIsIncremental(true);
// 同步时间
object.setLastSyncDate(LocalDateTime.now());
return object;
} catch (Exception e) {
log.error("构建对象 {} 元数据时出错: {}", objDetail.getName(), e.getMessage(), e);
return null;
}
}
/**
* 保存对象字段信息
*
* @param objDetail 对象详情
*/
private void saveObjectFieldsToDataiIntegrationField(DescribeSObjectResult objDetail) {
String objectApi = objDetail.getName();
List<DataiIntegrationField> fields = new ArrayList<>();
List<com.datai.integration.model.domain.DataiIntegrationPicklist> picklists = new ArrayList<>();
List<com.datai.integration.model.domain.DataiIntegrationFilterLookup> filterLookups = new ArrayList<>();
for (Field field : objDetail.getFields()) {
try {
DataiIntegrationField fieldEntity = new DataiIntegrationField();
fieldEntity.setApi(objectApi);
fieldEntity.setField(field.getName());
fieldEntity.setLabel(field.getLabel());
fieldEntity.setIsCreateable(field.isCreateable());
fieldEntity.setIsNillable(field.isNillable());
fieldEntity.setIsUpdateable(field.isUpdateable());
fieldEntity.setIsDefaultedOnCreate(field.isDefaultedOnCreate());
fieldEntity.setIsUnique(field.isUnique());
fieldEntity.setIsFilterable(field.isFilterable());
fieldEntity.setIsSortable(field.isSortable());
fieldEntity.setIsAggregatable(field.isAggregatable());
fieldEntity.setIsGroupable(field.isGroupable());
fieldEntity.setIsPolymorphicForeignKey(field.isPolymorphicForeignKey());
fieldEntity.setPolymorphicForeignField(field.getName() + "_type");
fieldEntity.setIsExternalId(field.isExternalId());
fieldEntity.setIsCustom(field.isCustom());
fieldEntity.setIsCalculated(field.isCalculated());
fieldEntity.setIsAutoNumber(field.isAutoNumber());
fieldEntity.setIsCaseSensitive(field.isCaseSensitive());
fieldEntity.setIsEncrypted(field.isEncrypted());
fieldEntity.setIsHtmlFormatted(field.isHtmlFormatted());
fieldEntity.setIsIdLookup(field.isIdLookup());
fieldEntity.setIsPermissionable(field.isPermissionable());
fieldEntity.setIsRestrictedPicklist(field.isRestrictedPicklist());
fieldEntity.setIsRestrictedDelete(field.isRestrictedDelete());
fieldEntity.setIsWriteRequiresMasterRead(field.isWriteRequiresMasterRead());
fieldEntity.setFieldDataType(field.getType() != null ? field.getType().toString() : null);
fieldEntity.setFieldLength(field.getLength());
fieldEntity.setFieldPrecision(field.getPrecision());
fieldEntity.setFieldScale(field.getScale());
fieldEntity.setFieldByteLength(field.getByteLength());
fieldEntity.setDefaultValue(field.getDefaultValueFormula());
fieldEntity.setCalculatedFormula(field.getCalculatedFormula());
fieldEntity.setInlineHelpText(field.getInlineHelpText());
fieldEntity.setRelationshipName(field.getRelationshipName());
fieldEntity.setRelationshipOrder(field.getRelationshipOrder());
fieldEntity.setReferenceTargetField(field.getReferenceTo() != null && field.getReferenceTo().length > 0 ? field.getReferenceTo()[0] : null);
if (field.getReferenceTo() != null && field.getReferenceTo().length > 0) {
fieldEntity.setReferenceTo(String.join(",", field.getReferenceTo()));
fieldEntity.setReferenceTargetField(field.getReferenceTo()[0]);
}
if (field.getType() != null) {
String fieldType = field.getType().toString().toLowerCase();
if ("picklist".equals(fieldType) || "multipicklist".equals(fieldType)) {
PicklistEntry[] picklistValues = field.getPicklistValues();
if (picklistValues != null && picklistValues.length > 0) {
for (PicklistEntry picklistValue : picklistValues) {
com.datai.integration.model.domain.DataiIntegrationPicklist picklist = new com.datai.integration.model.domain.DataiIntegrationPicklist();
picklist.setApi(objectApi);
picklist.setField(field.getName());
picklist.setPicklistLabel(picklistValue.getLabel());
picklist.setPicklistValue(picklistValue.getValue());
picklist.setIsActive(picklistValue.isActive());
picklist.setIsDefault(picklistValue.isDefaultValue());
picklists.add(picklist);
}
}
}
}
if (field.getType() != null && "reference".equals(field.getType().toString().toLowerCase())) {
FilteredLookupInfo filteredLookupInfo = field.getFilteredLookupInfo();
if (filteredLookupInfo != null) {
com.datai.integration.model.domain.DataiIntegrationFilterLookup filterLookup = new com.datai.integration.model.domain.DataiIntegrationFilterLookup();
filterLookup.setApi(objectApi);
filterLookup.setField(field.getName());
String[] controllingFields = filteredLookupInfo.getControllingFields();
if (controllingFields != null) {
filterLookup.setControllingField(String.join(",", controllingFields));
} else {
filterLookup.setControllingField(null);
}
filterLookup.setDependent(filteredLookupInfo.getDependent());
filterLookups.add(filterLookup);
}
}
fields.add(fieldEntity);
} catch (Exception e) {
log.error("保存对象 {} 的字段 {} 信息时出错: {}", objectApi, field.getName(), e.getMessage());
}
}
syncFieldsAndRecordChanges(objectApi, fields, picklists, filterLookups);
}
private void syncFieldsAndRecordChanges(String objectApi, List<DataiIntegrationField> newFields,
List<com.datai.integration.model.domain.DataiIntegrationPicklist> picklists,
List<com.datai.integration.model.domain.DataiIntegrationFilterLookup> filterLookups) {
try {
DataiIntegrationField queryField = new DataiIntegrationField();
queryField.setApi(objectApi);
List<DataiIntegrationField> existingFields = dataiIntegrationFieldService.selectDataiIntegrationFieldList(queryField);
Map<String, DataiIntegrationField> existingFieldMap = new HashMap<>();
for (DataiIntegrationField existingField : existingFields) {
existingFieldMap.put(existingField.getField(), existingField);
}
Map<String, DataiIntegrationField> newFieldMap = new HashMap<>();
for (DataiIntegrationField newField : newFields) {
newFieldMap.put(newField.getField(), newField);
}
DataiIntegrationObject queryObject = new DataiIntegrationObject();
queryObject.setApi(objectApi);
List<DataiIntegrationObject> objects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject);
String objectLabel = objects.isEmpty() ? objectApi : objects.get(0).getLabel();
Boolean isCustom = objects.isEmpty() ? false : objects.get(0).getIsCustom();
for (DataiIntegrationField newField : newFields) {
DataiIntegrationField existingField = existingFieldMap.get(newField.getField());
if (existingField == null) {
dataiIntegrationFieldService.insertDataiIntegrationField(newField);
recordFieldChange(objectApi, objectLabel, newField.getField(), newField.getLabel(),
null, "INSERT", "新增字段", isCustom);
log.debug("新增字段: {}.{}", objectApi, newField.getField());
} else {
List<String> changedFields = compareFields(existingField, newField);
if (!changedFields.isEmpty()) {
newField.setId(existingField.getId());
dataiIntegrationFieldService.updateDataiIntegrationField(newField);
recordFieldChange(objectApi, objectLabel, newField.getField(), newField.getLabel(),
"字段属性变更: " + String.join(", ", changedFields),
"UPDATE", "字段属性更新", isCustom);
log.debug("更新字段: {}.{} - 变更: {}", objectApi, newField.getField(),
String.join(", ", changedFields));
}
}
}
if (!picklists.isEmpty()) {
for (com.datai.integration.model.domain.DataiIntegrationPicklist picklist : picklists) {
dataiIntegrationPicklistService.insertDataiIntegrationPicklist(picklist);
}
}
if (!filterLookups.isEmpty()) {
for (com.datai.integration.model.domain.DataiIntegrationFilterLookup filterLookup : filterLookups) {
dataiIntegrationFilterLookupService.insertDataiIntegrationFilterLookup(filterLookup);
}
}
} catch (Exception e) {
log.error("同步对象 {} 的字段时出错: {}", objectApi, e.getMessage(), e);
}
}
/**
* 检查对象是否为大数据量对象
*
* @param connection Salesforce连接
* @param objectApi 对象API名称
* @return 对象数据量
*/
private int isLargeObject(PartnerConnection connection, String objectApi) {
try {
QueryResult queryResult = connection.queryAll("SELECT COUNT() FROM " + objectApi);
if (queryResult != null && queryResult.getRecords() != null && queryResult.getRecords().length > 0) {
SObject record = queryResult.getRecords()[0];
Object countValue = record.getField("expr0");
if (countValue != null) {
String countStr = countValue.toString();
long countLong = Long.parseLong(countStr);
return Math.min(Math.toIntExact(countLong), Integer.MAX_VALUE);
}
}
} catch (NumberFormatException e) {
log.error("将对象 {} 的数据量转换为数字时出错: {}", objectApi, e.getMessage());
} catch (ArithmeticException e) {
log.error("对象 {} 的数据量超过int最大值返回int最大值: {}", objectApi, e.getMessage());
return Integer.MAX_VALUE;
} catch (Exception e) {
log.error("检查对象 {} 的数据量时出错: {}", objectApi, e.getMessage());
}
return 0;
}
/**
* 更新对象的字段信息
*
* @param objectApi 对象API名称
* @param objectNum 对象数据量
*/
private void updateDataiIntegrationObjectFields(String objectApi, int objectNum) {
if (objectApi == null || objectApi.trim().isEmpty()) {
log.warn("对象API为空无法更新字段信息");
return;
}
try {
String dateField = dataiIntegrationFieldService.getDateField(objectApi);
String blobField = dataiIntegrationFieldService.getBlobField(objectApi);
boolean hasDeletedField = dataiIntegrationFieldService.isDeletedFieldExists(objectApi);
boolean isPartitioned = objectNum > LARGE_OBJECT_THRESHOLD;
log.info("对象 {} 字段信息,数据量: {}, 是否分区: {}, 日期字段: {}, 二进制字段: {}, 是否有删除字段: {}",
objectApi, objectNum, isPartitioned, dateField, blobField, hasDeletedField);
DataiIntegrationObject queryObject = new DataiIntegrationObject();
queryObject.setApi(objectApi);
List<DataiIntegrationObject> objects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject);
if (!objects.isEmpty()) {
DataiIntegrationObject object = objects.get(0);
object.setLastSyncDate(LocalDateTime.now());
int updateResult = dataiIntegrationObjectService.updateDataiIntegrationObject(object);
if (updateResult > 0) {
log.debug("成功更新对象 {} 的最后同步时间", objectApi);
} else {
log.warn("更新对象 {} 的最后同步时间失败", objectApi);
}
}
} catch (Exception e) {
log.error("处理对象 {} 的字段信息时出错: {}", objectApi, e.getMessage(), e);
}
}
} }

View File

@ -36,6 +36,23 @@
ALTER TABLE `${tableName}` ADD COLUMN `${fieldName}` varchar(255) ALTER TABLE `${tableName}` ADD COLUMN `${fieldName}` varchar(255)
</update> </update>
<!-- 添加字段(通用方法,支持字段类型和约束) -->
<update id="addField">
ALTER TABLE `${tableName}` ADD COLUMN `${fieldName}` ${fieldType}
<if test="!isNullable"> NOT NULL</if>
</update>
<!-- 修改字段(通用方法,支持字段类型和约束) -->
<update id="modifyField">
ALTER TABLE `${tableName}` MODIFY COLUMN `${fieldName}` ${fieldType}
<if test="!isNullable"> NOT NULL</if>
</update>
<!-- 删除字段(通用方法) -->
<update id="dropField">
ALTER TABLE `${tableName}` DROP COLUMN `${fieldName}`
</update>
<!-- 获取存在的ID列表 --> <!-- 获取存在的ID列表 -->
<select id="getIds" resultType="String"> <select id="getIds" resultType="String">
SELECT id SELECT id

View File

@ -241,4 +241,21 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
<if test="endTime != null"> and change_time &lt;= #{endTime}</if> <if test="endTime != null"> and change_time &lt;= #{endTime}</if>
</where> </where>
</select> </select>
<select id="countSimilarChanges" resultType="int">
select count(*) from datai_integration_metadata_change
where change_type = #{changeType}
and operation_type = #{operationType}
and object_api = #{objectApi}
<if test="fieldApi != null and fieldApi != ''">
and field_api = #{fieldApi}
</if>
<if test="fieldApi == null or fieldApi == ''">
and field_api is null
</if>
and change_reason = #{changeReason}
and sync_status = false
order by change_time desc
limit 1
</select>
</mapper> </mapper>