From 18d52e2ad1bbbefd8e4030e597314b63df398b41 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Sat, 3 Jan 2026 21:28:13 +0800 Subject: [PATCH] =?UTF-8?q?=E3=80=90feat=E3=80=91=20=E4=BC=98=E5=8C=96?= =?UTF-8?q?=E5=85=83=E6=95=B0=E6=8D=AE=E5=8F=98=E6=9B=B4=E5=90=8C=E6=AD=A5?= =?UTF-8?q?=E6=9C=AC=E5=9C=B0=E6=96=B9=E6=B3=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...aiIntegrationMetadataChangeController.java | 11 + .../integration/mapper/CustomMapper.java | 38 ++ .../DataiIntegrationMetadataChangeMapper.java | 12 + ...DataiIntegrationMetadataChangeService.java | 10 +- ...iIntegrationMetadataChangeServiceImpl.java | 644 +++++++++++++----- .../mapper/integration/CustomMapper.xml | 17 + .../DataiIntegrationMetadataChangeMapper.xml | 17 + 7 files changed, 567 insertions(+), 182 deletions(-) diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationMetadataChangeController.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationMetadataChangeController.java index d99cd6fa..334175fe 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationMetadataChangeController.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationMetadataChangeController.java @@ -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 = "同步元数据变更到本地数据库") @PreAuthorize("@ss.hasPermi('integration:change:sync')") diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/CustomMapper.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/CustomMapper.java index 5dbf73a6..af105dc4 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/CustomMapper.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/CustomMapper.java @@ -77,6 +77,44 @@ public interface CustomMapper { */ 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更新记录 * diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationMetadataChangeMapper.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationMetadataChangeMapper.java index 35b5920a..d9819e70 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationMetadataChangeMapper.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationMetadataChangeMapper.java @@ -86,4 +86,16 @@ public interface DataiIntegrationMetadataChangeMapper * @return 统计信息 */ public Map selectChangeStatistics(Map 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); } diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationMetadataChangeService.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationMetadataChangeService.java index f1f0567a..e99eaee7 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationMetadataChangeService.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationMetadataChangeService.java @@ -88,9 +88,17 @@ public interface IDataiIntegrationMetadataChangeService /** * 同步元数据变更到本地数据库 + * 根据元数据变更ID将指定的元数据变更同步到本地数据库 + * + * 该方法会: + * 1. 根据ID查询元数据变更记录 + * 2. 根据变更类型(OBJECT或FIELD)执行相应的同步操作 + * 3. 对于对象变更:执行对象的创建、修改或删除操作 + * 4. 对于字段变更:执行字段的创建、修改或删除操作 + * 5. 更新元数据变更记录的同步状态 * * @param id 元数据变更ID - * @return 同步结果 + * @return 同步结果,包含success(是否成功)和message(消息)字段 */ public Map syncToLocalDatabase(Long id); diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationMetadataChangeServiceImpl.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationMetadataChangeServiceImpl.java index df8a1277..a8064630 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationMetadataChangeServiceImpl.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationMetadataChangeServiceImpl.java @@ -8,6 +8,7 @@ import com.datai.common.utils.SecurityUtils; import com.datai.integration.factory.impl.SOAPConnectionFactory; import com.datai.integration.mapper.CustomMapper; import com.sforce.soap.partner.*; +import com.sforce.soap.partner.sobject.SObject; import com.sforce.ws.ConnectionException; import lombok.extern.slf4j.Slf4j; 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.IDataiIntegrationObjectService; 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; /** @@ -31,6 +34,11 @@ import com.datai.common.core.domain.model.LoginUser; @Service @Slf4j public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrationMetadataChangeService { + /** + * 大数据量对象阈值(500万) + */ + private static final int LARGE_OBJECT_THRESHOLD = 5000000; + @Autowired private DataiIntegrationMetadataChangeMapper dataiIntegrationMetadataChangeMapper; @@ -40,6 +48,12 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat @Autowired private IDataiIntegrationFieldService dataiIntegrationFieldService; + @Autowired + private IDataiIntegrationPicklistService dataiIntegrationPicklistService; + + @Autowired + private IDataiIntegrationFilterLookupService dataiIntegrationFilterLookupService; + @Autowired private SOAPConnectionFactory soapConnectionFactory; @@ -196,6 +210,7 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat Map result = new HashMap<>(); try { + // 根据ID查询元数据变更记录 DataiIntegrationMetadataChange metadataChange = selectDataiIntegrationMetadataChangeById(id); if (metadataChange == null) { result.put("success", false); @@ -203,12 +218,16 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat return result; } + // 获取变更类型和操作类型 String changeType = metadataChange.getChangeType(); String operationType = metadataChange.getOperationType(); + // 根据变更类型执行相应的同步操作 if ("OBJECT".equals(changeType)) { + // 对象级别的变更同步 syncObjectChange(metadataChange, result); } else if ("FIELD".equals(changeType)) { + // 字段级别的变更同步 syncFieldChange(metadataChange, result); } else { result.put("success", false); @@ -216,13 +235,17 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat return result; } + // 根据同步结果更新元数据变更记录的同步状态 if ((Boolean) result.get("success")) { + // 同步成功,更新同步状态为true updateSyncStatus(metadataChange.getId(), true, null); } else { + // 同步失败,更新同步状态为false并记录错误信息 updateSyncStatus(metadataChange.getId(), false, (String) result.get("message")); } } catch (Exception e) { + // 捕获异常并返回错误信息 result.put("success", false); result.put("message", "同步失败: " + e.getMessage()); updateSyncStatus(id, false, e.getMessage()); @@ -274,19 +297,53 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat try { if ("INSERT".equals(operationType)) { - DataiIntegrationObject object = new DataiIntegrationObject(); - object.setApi(objectApi); - object.setLabel(metadataChange.getObjectLabel()); - object.setIsCustom(metadataChange.getIsCustom()); - object.setCreateTime(DateUtils.getNowDate()); - object.setUpdateTime(DateUtils.getNowDate()); - - int insertResult = dataiIntegrationObjectService.insertDataiIntegrationObject(object); - result.put("success", insertResult > 0); - result.put("message", insertResult > 0 ? "对象创建成功" : "对象创建失败"); - - if (insertResult > 0) { - createOrUpdateDatabaseTable(objectApi, metadataChange.getObjectLabel(), result); + try { + log.info("开始同步对象: {}", objectApi); + + PartnerConnection connection = retryOperation(() -> soapConnectionFactory.getConnection(), 3, 1000); + log.info("成功获取源ORG连接"); + + 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); + result.put("success", insertResult > 0); + result.put("message", 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); + } 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)) { @@ -303,39 +360,34 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat DataiIntegrationObject object = existingObjects.get(0); object.setLabel(metadataChange.getObjectLabel()); + object.setIsWork(false); + object.setIsIncremental(false); object.setUpdateTime(DateUtils.getNowDate()); - + int updateResult = dataiIntegrationObjectService.updateDataiIntegrationObject(object); result.put("success", updateResult > 0); result.put("message", updateResult > 0 ? "对象更新成功" : "对象更新失败"); - if (updateResult > 0) { - updateDatabaseTable(objectApi, metadataChange.getObjectLabel(), result); - } - } else if ("DELETE".equals(operationType)) { DataiIntegrationObject queryObject = new DataiIntegrationObject(); queryObject.setApi(objectApi); - + List existingObjects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject); - + if (existingObjects.isEmpty()) { result.put("success", false); - result.put("message", "对象不存在,无法删除"); + result.put("message", "对象不存在,无法更新"); return; } - Integer[] ids = existingObjects.stream() - .map(DataiIntegrationObject::getId) - .toArray(Integer[]::new); - - int deleteResult = dataiIntegrationObjectService.deleteDataiIntegrationObjectByIds(ids); - result.put("success", deleteResult > 0); - result.put("message", deleteResult > 0 ? "对象删除成功" : "对象删除失败"); - - if (deleteResult > 0) { - dropDatabaseTable(objectApi, result); - } + DataiIntegrationObject object = existingObjects.get(0); + object.setIsWork(false); + object.setIsIncremental(false); + object.setUpdateTime(DateUtils.getNowDate()); + + int updateResult = dataiIntegrationObjectService.updateDataiIntegrationObject(object); + result.put("success", updateResult > 0); + result.put("message", updateResult > 0 ? "对象更新成功" : "对象更新失败"); } else { result.put("success", false); @@ -501,46 +553,46 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat } } - private void updateDatabaseTable(String objectApi, String objectLabel, Map result) { - PartnerConnection connection = null; - try { - connection = soapConnectionFactory.getConnection(); - DescribeSObjectResult objDetail = connection.describeSObject(objectApi); - - if (objDetail == null) { - result.put("success", false); - result.put("message", "无法获取对象元数据"); - return; - } - - List> fieldMaps = new ArrayList<>(); - - for (com.sforce.soap.partner.Field field : objDetail.getFields()) { - Map fieldMap = new HashMap<>(); - fieldMap.put("fieldName", field.getName()); - fieldMap.put("fieldType", convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null)); - fieldMap.put("fieldLength", field.getLength()); - fieldMap.put("fieldPrecision", field.getPrecision()); - fieldMap.put("fieldScale", field.getScale()); - fieldMap.put("isNullable", field.isNillable()); - fieldMap.put("isUnique", field.isUnique()); - fieldMap.put("isPrimaryKey", field.isIdLookup()); - fieldMaps.add(fieldMap); - } - - customMapper.createTable(objectApi, objectLabel, fieldMaps, new ArrayList<>()); - log.info("成功更新表结构: {}", objectApi); - - } catch (ConnectionException e) { - result.put("success", false); - result.put("message", "获取Salesforce连接失败: " + e.getMessage()); - log.error("获取Salesforce连接失败: {}", e.getMessage(), e); - } catch (Exception e) { - result.put("success", false); - result.put("message", "更新表结构失败: " + e.getMessage()); - log.error("更新表结构失败: {}", e.getMessage(), e); - } - } +// private void updateDatabaseTable(String objectApi, String objectLabel, Map result) { +// PartnerConnection connection = null; +// try { +// connection = soapConnectionFactory.getConnection(); +// DescribeSObjectResult objDetail = connection.describeSObject(objectApi); +// +// if (objDetail == null) { +// result.put("success", false); +// result.put("message", "无法获取对象元数据"); +// return; +// } +// +// List> fieldMaps = new ArrayList<>(); +// +// for (com.sforce.soap.partner.Field field : objDetail.getFields()) { +// Map fieldMap = new HashMap<>(); +// fieldMap.put("fieldName", field.getName()); +// fieldMap.put("fieldType", convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null)); +// fieldMap.put("fieldLength", field.getLength()); +// fieldMap.put("fieldPrecision", field.getPrecision()); +// fieldMap.put("fieldScale", field.getScale()); +// fieldMap.put("isNullable", field.isNillable()); +// fieldMap.put("isUnique", field.isUnique()); +// fieldMap.put("isPrimaryKey", field.isIdLookup()); +// fieldMaps.add(fieldMap); +// } +// +// customMapper.createTable(objectApi, objectLabel, fieldMaps, new ArrayList<>()); +// log.info("成功更新表结构: {}", objectApi); +// +// } catch (ConnectionException e) { +// result.put("success", false); +// result.put("message", "获取Salesforce连接失败: " + e.getMessage()); +// log.error("获取Salesforce连接失败: {}", e.getMessage(), e); +// } catch (Exception e) { +// result.put("success", false); +// result.put("message", "更新表结构失败: " + e.getMessage()); +// log.error("更新表结构失败: {}", e.getMessage(), e); +// } +// } private void dropDatabaseTable(String objectApi, Map result) { try { @@ -568,8 +620,10 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat for (com.sforce.soap.partner.Field field : objDetail.getFields()) { if (fieldApi.equals(field.getName())) { - customMapper.createField(objectApi, fieldApi); - log.info("成功添加字段: {}.{}", objectApi, fieldApi); + String mysqlType = convertSalesforceTypeToMySQL(field.getType() != null ? field.getType().toString() : null); + + customMapper.addField(objectApi, fieldApi, mysqlType, field.isNillable()); + log.info("成功添加字段: {}.{} 类型: {}", objectApi, fieldApi, mysqlType); return; } } @@ -603,10 +657,9 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat for (com.sforce.soap.partner.Field field : objDetail.getFields()) { if (fieldApi.equals(field.getName())) { 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.executeUpdate(alterSql); - log.info("成功修改字段: {}.{}", objectApi, fieldApi); + + customMapper.modifyField(objectApi, fieldApi, mysqlType, field.isNillable()); + log.info("成功修改字段: {}.{} 类型: {}", objectApi, fieldApi, mysqlType); return; } } @@ -627,8 +680,7 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat private void dropDatabaseColumn(String objectApi, String fieldApi, Map result) { try { - String alterSql = String.format("ALTER TABLE %s DROP COLUMN %s", objectApi, fieldApi); - customMapper.executeUpdate(alterSql); + customMapper.dropField(objectApi, fieldApi); log.info("成功删除字段: {}.{}", objectApi, fieldApi); } catch (Exception e) { result.put("success", false); @@ -894,50 +946,6 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat * @param objDetail Salesforce对象的详细描述信息 * @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的Field对象转换为DataiIntegrationField实体 @@ -1249,45 +1257,40 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat */ private void recordObjectChange(DataiIntegrationObject newObject, DataiIntegrationObject oldObject, String operationType) { try { - // 创建元数据变更记录实体 - DataiIntegrationMetadataChange metadataChange = new DataiIntegrationMetadataChange(); - // 设置变更类型为对象 - metadataChange.setChangeType("OBJECT"); - // 设置操作类型(INSERT、UPDATE、DELETE) - 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),则比较新旧对象的差异 + String changeReason; if (oldObject != null) { List changedFields = compareObjects(oldObject, newObject); if (!changedFields.isEmpty()) { - // 如果有具体变更字段,则记录详细变更信息 - metadataChange.setChangeReason("对象属性变更: " + String.join(", ", changedFields)); + changeReason = "对象属性变更: " + String.join(", ", changedFields); } else { - // 如果没有具体变更字段,则记录通用更新信息 - metadataChange.setChangeReason("对象属性更新"); + changeReason = "对象属性更新"; } } else { - // 如果是新增操作,则记录为新增对象 - metadataChange.setChangeReason("新增对象"); + changeReason = "新增对象"; } - // 设置变更用户为系统 + + 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"); - // 插入元数据变更记录 insertDataiIntegrationMetadataChange(metadataChange); log.debug("记录对象变更成功: {} - {}", newObject.getApi(), operationType); } catch (Exception 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, String changeReason, String operationType, String defaultReason, Boolean isCustom) { 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(); - // 设置变更类型为字段 metadataChange.setChangeType("FIELD"); - // 设置操作类型(INSERT、UPDATE、DELETE) metadataChange.setOperationType(operationType); - // 设置对象API名称 metadataChange.setObjectApi(objectApi); - // 设置对象标签名称 metadataChange.setObjectLabel(objectLabel); - // 设置字段API名称 metadataChange.setFieldApi(fieldApi); - // 设置字段标签名称 metadataChange.setFieldLabel(fieldLabel); - // 设置变更时间 metadataChange.setChangeTime(LocalDateTime.now()); - // 设置同步状态为未同步 metadataChange.setSyncStatus(false); - // 设置是否为自定义对象 metadataChange.setIsCustom(isCustom); - // 设置变更原因,如果changeReason不为null则使用它,否则使用默认原因 - metadataChange.setChangeReason(changeReason != null ? changeReason : defaultReason); - // 设置变更用户为系统 + metadataChange.setChangeReason(finalChangeReason); metadataChange.setChangeUser("SYSTEM"); - // 插入元数据变更记录 insertDataiIntegrationMetadataChange(metadataChange); log.debug("记录字段变更成功: {}.{} - {}", objectApi, fieldApi, operationType); } catch (Exception e) { - // 记录记录字段变更时的错误 log.error("记录字段变更失败: {}.{} - {}", objectApi, fieldApi, e.getMessage(), e); } } @@ -1351,41 +1350,37 @@ public class DataiIntegrationMetadataChangeServiceImpl implements IDataiIntegrat */ private void checkDeletedObjectsForMetadata(Set syncedObjectApis) { try { - // 查询数据库中所有对象 DataiIntegrationObject queryObject = new DataiIntegrationObject(); List allObjects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject); - // 遍历数据库中的所有对象,检查哪些对象已不在Salesforce中 for (DataiIntegrationObject object : allObjects) { 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(); - // 设置变更类型为对象 metadataChange.setChangeType("OBJECT"); - // 设置操作类型为删除 metadataChange.setOperationType("DELETE"); - // 设置对象API名称 metadataChange.setObjectApi(object.getApi()); - // 设置对象标签名称 metadataChange.setObjectLabel(object.getLabel()); - // 设置变更时间 metadataChange.setChangeTime(LocalDateTime.now()); - // 设置同步状态为未同步 metadataChange.setSyncStatus(false); - // 设置是否为自定义对象 metadataChange.setIsCustom(object.getIsCustom()); - // 设置变更原因 - metadataChange.setChangeReason("对象已从Salesforce中删除"); - // 设置变更用户为系统 + metadataChange.setChangeReason(changeReason); metadataChange.setChangeUser("SYSTEM"); - // 插入元数据变更记录 insertDataiIntegrationMetadataChange(metadataChange); log.warn("检测到对象已删除: {}", object.getApi()); } } } catch (Exception 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 fields = new ArrayList<>(); + List picklists = new ArrayList<>(); + List 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 newFields, + List picklists, + List filterLookups) { + try { + DataiIntegrationField queryField = new DataiIntegrationField(); + queryField.setApi(objectApi); + List existingFields = dataiIntegrationFieldService.selectDataiIntegrationFieldList(queryField); + + Map existingFieldMap = new HashMap<>(); + for (DataiIntegrationField existingField : existingFields) { + existingFieldMap.put(existingField.getField(), existingField); + } + + Map newFieldMap = new HashMap<>(); + for (DataiIntegrationField newField : newFields) { + newFieldMap.put(newField.getField(), newField); + } + + DataiIntegrationObject queryObject = new DataiIntegrationObject(); + queryObject.setApi(objectApi); + List 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 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 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); + } + } } diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/CustomMapper.xml b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/CustomMapper.xml index c0be0dfd..97f3e15d 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/CustomMapper.xml +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/CustomMapper.xml @@ -36,6 +36,23 @@ ALTER TABLE `${tableName}` ADD COLUMN `${fieldName}` varchar(255) + + + ALTER TABLE `${tableName}` ADD COLUMN `${fieldName}` ${fieldType} + NOT NULL + + + + + ALTER TABLE `${tableName}` MODIFY COLUMN `${fieldName}` ${fieldType} + NOT NULL + + + + + ALTER TABLE `${tableName}` DROP COLUMN `${fieldName}` + + + + \ No newline at end of file