diff --git a/README.md b/README.md index 6747a8b0..d56dc933 100644 --- a/README.md +++ b/README.md @@ -17,8 +17,8 @@ 本人的其他两个推荐搭配的项目 -1. [DataI-App-Geek: 这是若依极客生态的小程序版本 (gitee.com)](https://gitee.com/geek-xd/geek-uniapp-vue3-uview-plus-uchart) -2. [DataI-Vue3-Geek: 这是若依极客生态的Vue3版本 (gitee.com)](https://gitee.com/geek-xd/ruo-yi-vue3-geek) +1. [DataI-App-Geek: 这是集成系统生态的小程序版本 (gitee.com)](https://gitee.com/geek-xd/geek-uniapp-vue3-uview-plus-uchart) +2. [DataI-Vue3-Geek: 这是集成系统生态的Vue3版本 (gitee.com)](https://gitee.com/geek-xd/ruo-yi-vue3-geek) 与本项目同为一个作者开发,兼容性最好,学习成本最低。 diff --git a/datai-admin/src/main/java/com/datai/DataIApplication.java b/datai-admin/src/main/java/com/datai/DataIApplication.java index 9a1daaa4..013cf3cf 100644 --- a/datai-admin/src/main/java/com/datai/DataIApplication.java +++ b/datai-admin/src/main/java/com/datai/DataIApplication.java @@ -17,19 +17,9 @@ import org.springframework.core.env.Environment; @SpringBootApplication(exclude = { DataSourceAutoConfiguration.class }) public class DataIApplication { public static void main(String[] args) throws UnknownHostException { - // System.setProperty("spring.devtools.restart.enabled", "false"); + System.setProperty("spring.devtools.restart.enabled", "false"); ConfigurableApplicationContext application = SpringApplication.run(DataIApplication.class, args); - System.out.println("(♥◠‿◠)ノ゙ 若依极客启动成功 ლ(´ڡ`ლ)゙ \n" + - " .-------. ____ __ \n" + - " | _ _ \\ \\ \\ / / \n" + - " | ( ' ) | \\ _. / ' \n" + - " |(_ o _) / _( )_ .' \n" + - " | (_,_).' __ ___(_ o _)' " + " ____ _ " + "\n" + - " | |\\ \\ | || |(_,_)' " + " / ___| ___ ___| | __ " + "\n" + - " | | \\ `' /| `-' / " + "| | _ / _ \\/ _ \\ |/ / " + "\n" + - " | | \\ / \\ / " + " | |_| | __/ __/ < " + "\n" + - " ''-' `'-' `-..-' " + "\\____|\\___|\\___|_|\\_\\"); - + System.out.println("(♥◠‿◠)ノ゙ 集成系统启动成功 ლ(´ڡ`ლ)゙ \n"); Environment env = application.getEnvironment(); String ip = InetAddress.getLocalHost().getHostAddress(); String port = env.getProperty("server.port"); diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationObjectController.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationObjectController.java index ba051d3f..e06836d2 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationObjectController.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationObjectController.java @@ -24,10 +24,12 @@ import com.datai.common.enums.BusinessType; import com.datai.integration.model.domain.DataiIntegrationObject; import com.datai.integration.model.dto.DataiIntegrationObjectDto; import com.datai.integration.service.IDataiIntegrationObjectService; +import com.datai.integration.service.IDataiIntegrationFieldService; import com.datai.common.utils.poi.ExcelUtil; import com.datai.common.core.page.TableDataInfo; import io.swagger.v3.oas.annotations.tags.Tag; import io.swagger.v3.oas.annotations.Operation; +import lombok.extern.slf4j.Slf4j; /** * 对象同步控制Controller @@ -38,6 +40,7 @@ import io.swagger.v3.oas.annotations.Operation; @RestController @RequestMapping("/integration/object") @Tag(name = "【对象同步控制】管理") +@Slf4j public class DataiIntegrationObjectController extends BaseController { @Autowired @@ -214,4 +217,45 @@ public class DataiIntegrationObjectController extends BaseController return error((String) statistics.get("message")); } } + + /** + * 同步单对象数据到本地数据库 + * + * 该方法用于触发指定对象的单次数据同步操作,将Salesforce对象的数据同步到本地数据库 + * 同步操作包括全量同步和增量同步两种模式,具体模式由对象的配置决定 + * + * @param id 对象ID,用于标识需要同步的Salesforce对象 + * @return AjaxResult 同步结果,包含成功/失败状态、同步数据量、耗时等信息 + * - 成功时返回:success=true, message="对象数据同步成功", 以及详细的同步信息 + * - 失败时返回:success=false, message=错误信息 + */ + @Operation(summary = "同步单对象数据到本地数据库") + @PreAuthorize("@ss.hasPermi('integration:object:syncData')") + @Log(title = "对象同步控制", businessType = BusinessType.UPDATE) + @PostMapping("/{id}/syncData") + public AjaxResult syncObjectData(@PathVariable("id") Integer id) + { + if (id == null) { + log.error("对象ID为空,无法同步数据"); + return error("对象ID不能为空"); + } + + try { + log.info("开始同步对象数据,对象ID: {}", id); + + Map result = dataiIntegrationObjectService.syncSingleObjectData(id); + + if ((Boolean) result.get("success")) { + log.info("对象数据同步成功,对象ID: {}", id); + return success(result); + } else { + log.error("对象数据同步失败,对象ID: {}, 错误信息: {}", id, result.get("message")); + return error((String) result.get("message")); + } + + } catch (Exception e) { + log.error("同步对象数据时发生异常,对象ID: {}", id, e); + return error("同步对象数据时发生异常: " + e.getMessage()); + } + } } diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationObjectMapper.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationObjectMapper.java index ed405eb5..46a09477 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationObjectMapper.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationObjectMapper.java @@ -20,6 +20,14 @@ public interface DataiIntegrationObjectMapper */ public DataiIntegrationObject selectDataiIntegrationObjectById(Integer id); + /** + * 根据API查询对象同步控制 + * + * @param api 对象API + * @return 对象同步控制 + */ + public DataiIntegrationObject selectDataiIntegrationObjectByApi(String api); + /** * 查询对象同步控制列表 * diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationObjectService.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationObjectService.java index 847d2307..8806f10c 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationObjectService.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/IDataiIntegrationObjectService.java @@ -108,4 +108,12 @@ public interface IDataiIntegrationObjectService * @return 统计信息 */ public Map getObjectStatistics(); + + /** + * 同步单对象数据到本地数据库 + * + * @param id 对象ID + * @return 同步结果 + */ + public Map syncSingleObjectData(Integer id); } diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationBatchServiceImpl.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationBatchServiceImpl.java index 5b5a0ba5..119cb5aa 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationBatchServiceImpl.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationBatchServiceImpl.java @@ -1,5 +1,6 @@ package com.datai.integration.service.impl; +import java.time.LocalDateTime; import java.util.List; import java.util.Map; import java.util.HashMap; @@ -21,6 +22,7 @@ import com.datai.common.core.domain.model.LoginUser; import com.datai.integration.model.domain.DataiIntegrationBatchHistory; import com.datai.integration.service.IDataiIntegrationBatchHistoryService; import com.datai.integration.service.IDataiIntegrationFieldService; +import com.datai.integration.service.IDataiIntegrationSyncLogService; import com.datai.integration.mapper.CustomMapper; import com.datai.integration.factory.impl.SOAPConnectionFactory; import com.sforce.soap.partner.DescribeSObjectResult; @@ -56,6 +58,9 @@ public class DataiIntegrationBatchServiceImpl implements IDataiIntegrationBatchS @Autowired private SOAPConnectionFactory soapConnectionFactory; + @Autowired + private IDataiIntegrationSyncLogService dataiIntegrationSyncLogService; + /** * 查询数据批次 * @@ -338,7 +343,18 @@ public class DataiIntegrationBatchServiceImpl implements IDataiIntegrationBatchS return false; } + DataiIntegrationBatch batch = null; + DataiIntegrationBatchHistory batchHistory = new DataiIntegrationBatchHistory(); + long startTime = System.currentTimeMillis(); + LocalDateTime syncStartTime = LocalDateTime.now(); + try { + batch = selectDataiIntegrationBatchById(Integer.parseInt(batchId)); + if (batch == null) { + log.error("批次不存在,批次ID: {}", batchId); + return false; + } + PartnerConnection connection = retryOperation(() -> soapConnectionFactory.getConnection(), 3, 1000); log.info("成功获取Salesforce SOAP连接"); @@ -355,14 +371,83 @@ public class DataiIntegrationBatchServiceImpl implements IDataiIntegrationBatchS int totalCount = executeQueryAndProcessData(connection, param, fieldList); log.info("对象 {} 批次 {} 数据同步完成,共处理 {} 条记录", objectApi, batchId, totalCount); + + long endTime = System.currentTimeMillis(); + long duration = endTime - startTime; + LocalDateTime syncEndTime = LocalDateTime.now(); + + batchHistory.setApi(objectApi); + batchHistory.setLabel(batch.getLabel()); + batchHistory.setBatchId(Integer.parseInt(batchId)); + batchHistory.setBatchField(batch.getBatchField()); + batchHistory.setSyncNum(totalCount); + batchHistory.setSyncType(batch.getSyncType()); + batchHistory.setSyncStatus(true); + batchHistory.setStartTime(syncStartTime); + batchHistory.setEndTime(syncEndTime); + batchHistory.setCost(duration); + batchHistory.setSyncStartTime(syncStartTime); + batchHistory.setSyncEndTime(syncEndTime); + batchHistoryService.insertDataiIntegrationBatchHistory(batchHistory); + log.info("批次 {} 历史记录已保存,同步数据量: {}, 耗时: {}ms", batchId, totalCount, duration); + insertSyncLog(objectApi, batch.getLabel(), batch.getSyncType(), true, null, duration); + return true; } catch (Exception e) { log.error("同步Salesforce对象指定批次数据失败,对象API: {}, 批次ID: {}", objectApi, batchId, e); + + long endTime = System.currentTimeMillis(); + long duration = endTime - startTime; + LocalDateTime syncEndTime = LocalDateTime.now(); + + batchHistory.setApi(objectApi); + batchHistory.setLabel(batch != null ? batch.getLabel() : ""); + batchHistory.setBatchId(Integer.parseInt(batchId)); + batchHistory.setBatchField(batch != null ? batch.getBatchField() : ""); + batchHistory.setSyncNum(0); + batchHistory.setSyncType(batch != null ? batch.getSyncType() : ""); + batchHistory.setSyncStatus(false); + batchHistory.setStartTime(syncStartTime); + batchHistory.setEndTime(syncEndTime); + batchHistory.setCost(duration); + batchHistory.setSyncStartTime(syncStartTime); + batchHistory.setSyncEndTime(syncEndTime); + batchHistoryService.insertDataiIntegrationBatchHistory(batchHistory); + log.info("批次 {} 失败历史记录已保存,耗时: {}ms", batchId, duration); + + insertSyncLog(objectApi, batch != null ? batch.getLabel() : "", batch != null ? batch.getSyncType() : "", false, e.getMessage(), duration); + return false; } } + private void insertSyncLog(String objectApi, String label, String syncType, boolean success, String errorMessage, long duration) { + try { + com.datai.integration.model.domain.DataiIntegrationSyncLog syncLog = new com.datai.integration.model.domain.DataiIntegrationSyncLog(); + syncLog.setObjectApi(objectApi); + syncLog.setOperationType("FULL".equals(syncType) ? "全量同步" : "增量同步"); + syncLog.setOperationStatus(success ? "成功" : "失败"); + syncLog.setErrorMessage(errorMessage); + syncLog.setExecutionTime(new java.math.BigDecimal(duration / 1000.0)); + syncLog.setCreateTime(DateUtils.getNowDate()); + + try { + LoginUser loginUser = SecurityUtils.getLoginUser(); + if (loginUser != null && loginUser.getDeptId() != null) { + syncLog.setDeptId(loginUser.getDeptId()); + } + } catch (Exception ex) { + log.warn("获取部门ID失败: {}", ex.getMessage()); + } + + dataiIntegrationSyncLogService.insertDataiIntegrationSyncLog(syncLog); + log.info("同步日志已记录,对象API: {}, 操作类型: {}, 状态: {}", objectApi, syncLog.getOperationType(), syncLog.getOperationStatus()); + } catch (Exception e) { + log.error("插入同步日志失败: {}", e.getMessage(), e); + } + } + /** * 重试操作 * diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationObjectServiceImpl.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationObjectServiceImpl.java index 2360778d..511d3c70 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationObjectServiceImpl.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/service/impl/DataiIntegrationObjectServiceImpl.java @@ -1,10 +1,9 @@ package com.datai.integration.service.impl; +import java.text.SimpleDateFormat; import java.time.LocalDateTime; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import java.util.*; +import java.util.concurrent.CopyOnWriteArrayList; import com.datai.common.core.domain.model.LoginUser; import com.datai.common.utils.DateUtils; @@ -12,12 +11,20 @@ import com.datai.common.utils.SecurityUtils; import com.datai.integration.factory.impl.SOAPConnectionFactory; import com.datai.integration.mapper.CustomMapper; import com.datai.integration.mapper.DataiIntegrationObjectMapper; +import com.datai.integration.model.domain.DataiIntegrationBatch; import com.datai.integration.model.domain.DataiIntegrationObject; +import com.datai.integration.service.IDataiIntegrationBatchService; +import com.datai.integration.service.IDataiIntegrationFieldService; import com.datai.integration.service.IDataiIntegrationObjectService; +import com.datai.salesforce.common.param.SalesforceParam; +import com.datai.setting.future.SalesforceExecutor; import com.sforce.soap.partner.DescribeSObjectResult; import com.sforce.soap.partner.Field; import com.sforce.soap.partner.PartnerConnection; +import com.sforce.soap.partner.QueryResult; +import com.sforce.soap.partner.sobject.SObject; import com.sforce.ws.ConnectionException; +import org.apache.commons.lang3.StringUtils; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -42,6 +49,15 @@ public class DataiIntegrationObjectServiceImpl implements IDataiIntegrationObjec @Autowired private SOAPConnectionFactory soapConnectionFactory; + @Autowired + private IDataiIntegrationFieldService dataiIntegrationFieldService; + + @Autowired + private IDataiIntegrationBatchService dataiIntegrationBatchService; + + @Autowired + private SalesforceExecutor salesforceExecutor; + /** * 查询对象同步控制 * @@ -305,6 +321,40 @@ public class DataiIntegrationObjectServiceImpl implements IDataiIntegrationObjec return result; } + /** + * 将Salesforce数据类型转换为MySQL数据类型 + * + * 该方法负责将Salesforce对象字段的数据类型映射为MySQL数据库对应的数据类型, + * 用于自动创建本地数据表时确定字段的SQL数据类型定义。 + * + * 类型映射规则: + * - id → VARCHAR(18):Salesforce的ID字段为18位字符串 + * - string/email/url/phone → VARCHAR(255):文本类型字段 + * - textarea → TEXT:长文本字段 + * - boolean → TINYINT(1):布尔值,MySQL中用0和1表示 + * - int → INT:整数字段 + * - double → DOUBLE:双精度浮点数 + * - currency → DECIMAL(18,4):货币类型,保留4位小数 + * - date → DATE:日期类型 + * - datetime → DATETIME:日期时间类型 + * - time → TIME:时间类型 + * - percent → DECIMAL(10,2):百分比,保留2位小数 + * - reference → VARCHAR(18):引用字段,存储其他对象的ID + * - picklist/multipicklist → VARCHAR(255):下拉列表字段 + * - combobox → VARCHAR(255):组合框字段 + * - base64 → LONGBLOB:二进制大对象,用于存储附件等 + * - address → TEXT:地址字段 + * - location → VARCHAR(255):位置字段 + * - 其他未知类型 → VARCHAR(255):默认类型 + * + * 使用场景: + * - 自动创建数据表时确定字段类型 + * - 动态生成DDL语句 + * - 确保数据类型兼容性 + * + * @param typeStr Salesforce字段类型字符串,如"string"、"int"等 + * @return 对应的MySQL数据类型字符串,如"VARCHAR(255)"、"INT"等 + */ private String convertSalesforceTypeToMySQL(String typeStr) { if (typeStr == null || typeStr.isEmpty()) { return "VARCHAR(255)"; @@ -579,4 +629,825 @@ public class DataiIntegrationObjectServiceImpl implements IDataiIntegrationObjec return statistics; } + + /** + * 同步单个Salesforce对象的数据到本地数据库 + * + * @param id 对象ID,用于标识需要同步的Salesforce对象 + * @return Map 同步结果,包含success、message、objectId、objectApi、totalCount、duration、syncType、lastFullSyncDate等字段 + * @throws ConnectionException 当获取Salesforce连接失败时抛出 + * @throws RuntimeException 当同步过程中发生其他异常时抛出 + */ + @Override + public Map syncSingleObjectData(Integer id) + { + Map result = new HashMap<>(); + long startTime = System.currentTimeMillis(); + DataiIntegrationObject object = null; + + try { + if (id == null) { + log.error("对象ID不能为空"); + result.put("success", false); + result.put("message", "对象ID不能为空"); + return result; + } + + log.info("开始同步单对象数据,对象ID: {}", id); + + object = dataiIntegrationObjectMapper.selectDataiIntegrationObjectById(id); + if (object == null) { + log.error("对象不存在,对象ID: {}", id); + result.put("success", false); + result.put("message", "对象不存在"); + return result; + } + + if (object.getApi() == null || object.getApi().trim().isEmpty()) { + log.error("对象API为空,对象ID: {}", id); + result.put("success", false); + result.put("message", "对象API不能为空"); + return result; + } + + if (!object.getIsWork()) { + log.warn("对象未启用同步,对象ID: {}, 对象API: {}", id, object.getApi()); + result.put("success", false); + result.put("message", "对象未启用同步"); + return result; + } + + String objectApi = object.getApi().trim(); + + log.info("准备同步对象数据,对象API: {}", objectApi); + + LocalDateTime lastFullSyncDate = object.getLastFullSyncDate(); + Map syncResult; + + if (lastFullSyncDate == null) { + log.info("对象 {} 的全量拉取时间为空,执行全量数据拉取", objectApi); + syncResult = syncFullData(object, startTime); + } else { + log.info("对象 {} 的全量拉取时间不为空({}),执行增量数据拉取", objectApi, lastFullSyncDate); + syncResult = syncIncrementalData(object, startTime); + } + + result.putAll(syncResult); + + } catch (Exception e) { + log.error("同步单对象数据失败,对象ID: {}", id, e); + if (object != null) { + object.setErrorMessage("同步单对象数据失败: " + e.getMessage()); + object.setSyncStatus(false); + object.setUpdateTime(DateUtils.getNowDate()); + updateDataiIntegrationObject(object); + } + result.put("success", false); + result.put("message", "同步单对象数据失败: " + e.getMessage()); + } + + return result; + } + + /** + * 全量数据拉取 + * 查询所有批次并使用多线程并行拉取 + * + * @param object 对象信息 + * @param startTime 开始时间 + * @return 同步结果 + */ + private Map syncFullData(DataiIntegrationObject object, long startTime) { + Map result = new HashMap<>(); + String objectApi = object.getApi().trim(); + + try { + DataiIntegrationBatch queryBatch = new DataiIntegrationBatch(); + queryBatch.setApi(objectApi); + queryBatch.setSyncType("FULL"); + List batches = dataiIntegrationBatchService.selectDataiIntegrationBatchList(queryBatch); + + if (batches.isEmpty()) { + log.warn("对象 {} 没有找到全量同步批次,无法执行全量数据拉取", objectApi); + result.put("success", false); + result.put("message", "没有找到全量同步批次"); + return result; + } + + log.info("对象 {} 共找到 {} 个全量同步批次,准备使用多线程并行拉取", objectApi, batches.size()); + + List> futures = new ArrayList<>(); + List> batchResults = new CopyOnWriteArrayList<>(); + int totalSyncCount = 0; + int successBatchCount = 0; + int failedBatchCount = 0; + + for (int i = 0; i < batches.size(); i++) { + DataiIntegrationBatch batch = batches.get(i); + final int batchIndex = i; + + java.util.concurrent.Future future = salesforceExecutor.execute(() -> { + try { + log.info("开始拉取批次 {} (批次ID: {}, 批次索引: {}/{})", + batch.getSyncStartDate(), batch.getId(), batchIndex + 1, batches.size()); + + Map batchResult = dataiIntegrationBatchService.syncBatchData(batch.getId()); + batchResults.add(batchResult); + + if ((Boolean) batchResult.get("success")) { + log.info("批次 {} 数据拉取成功,同步数据量: {}", + batch.getSyncStartDate(), batchResult.get("syncNum")); + } else { + log.error("批次 {} 数据拉取失败: {}", + batch.getSyncStartDate(), batchResult.get("message")); + } + + } catch (Exception e) { + log.error("拉取批次 {} 时发生异常", batch.getSyncStartDate(), e); + Map errorResult = new HashMap<>(); + errorResult.put("success", false); + errorResult.put("message", e.getMessage()); + batchResults.add(errorResult); + } + }, 0, i); + + futures.add(future); + } + + for (java.util.concurrent.Future future : futures) { + try { + future.get(); + } catch (Exception e) { + log.error("获取批次同步结果时发生异常", e); + failedBatchCount++; + } + } + + for (Map batchResult : batchResults) { + if ((Boolean) batchResult.get("success")) { + successBatchCount++; + if (batchResult.get("syncNum") != null) { + totalSyncCount += (Integer) batchResult.get("syncNum"); + } + } else { + failedBatchCount++; + } + } + + long endTime = System.currentTimeMillis(); + long duration = endTime - startTime; + LocalDateTime now = LocalDateTime.now(); + + object.setLastFullSyncDate(now); + object.setLastSyncDate(now); + object.setTotalRows(totalSyncCount); + object.setSyncStatus(failedBatchCount == 0); + object.setUpdateTime(DateUtils.getNowDate()); + updateDataiIntegrationObject(object); + + result.put("success", true); + result.put("message", "全量数据拉取完成"); + result.put("objectId", object.getId()); + result.put("objectApi", objectApi); + result.put("totalCount", totalSyncCount); + result.put("duration", duration); + result.put("syncType", "full"); + result.put("lastFullSyncDate", object.getLastFullSyncDate()); + result.put("totalBatchCount", batches.size()); + result.put("successBatchCount", successBatchCount); + result.put("failedBatchCount", failedBatchCount); + + log.info("对象 {} 全量数据拉取完成,共处理 {} 个批次,成功 {} 个,失败 {} 个,总记录数: {},耗时: {}ms", + objectApi, batches.size(), successBatchCount, failedBatchCount, totalSyncCount, duration); + + } catch (Exception e) { + log.error("全量数据拉取失败,对象API: {}", objectApi, e); + result.put("success", false); + result.put("message", "全量数据拉取失败: " + e.getMessage()); + } + + return result; + } + + /** + * 增量数据拉取 + * 创建增量批次并使用多线程并行拉取 + * + * @param object 对象信息 + * @param startTime 开始时间 + * @return 同步结果 + */ + private Map syncIncrementalData(DataiIntegrationObject object, long startTime) { + Map result = new HashMap<>(); + String objectApi = object.getApi().trim(); + LocalDateTime lastFullSyncDate = object.getLastFullSyncDate(); + + try { + if (!Boolean.TRUE.equals(object.getIsIncremental())) { + log.warn("对象 {} 未启用增量更新,不执行增量数据拉取", objectApi); + result.put("success", false); + result.put("message", "对象未启用增量更新"); + return result; + } + + PartnerConnection connection = soapConnectionFactory.getConnection(); + if (connection == null) { + log.error("无法获取Salesforce连接"); + result.put("success", false); + result.put("message", "无法获取Salesforce连接"); + return result; + } + + List fieldList = getSalesforceObjectFields(connection, objectApi); + String batchField = object.getBatchField(); + + if (batchField == null || batchField.trim().isEmpty()) { + log.warn("对象 {} 没有设置批次字段,使用默认日期字段", objectApi); + batchField = dataiIntegrationFieldService.getDateField(objectApi); + } + + if (batchField == null || batchField.trim().isEmpty()) { + log.error("对象 {} 无法确定批次字段,无法创建增量批次", objectApi); + result.put("success", false); + result.put("message", "无法确定批次字段"); + return result; + } + + DataiIntegrationBatch incrementalBatch = new DataiIntegrationBatch(); + incrementalBatch.setApi(objectApi); + incrementalBatch.setLabel(object.getLabel()); + incrementalBatch.setBatchField(batchField); + incrementalBatch.setSyncType("INCREMENTAL"); + incrementalBatch.setSyncStatus(false); + incrementalBatch.setSyncStartDate(lastFullSyncDate); + incrementalBatch.setSyncEndDate(LocalDateTime.now()); + incrementalBatch.setCreateTime(DateUtils.getNowDate()); + incrementalBatch.setUpdateTime(DateUtils.getNowDate()); + + int batchId = dataiIntegrationBatchService.insertDataiIntegrationBatch(incrementalBatch); + log.info("为对象 {} 创建增量批次,批次ID: {}", objectApi, batchId); + + Map batchResult = dataiIntegrationBatchService.syncBatchData(batchId); + + long endTime = System.currentTimeMillis(); + long duration = endTime - startTime; + LocalDateTime now = LocalDateTime.now(); + + if ((Boolean) batchResult.get("success")) { + object.setLastSyncDate(now); + Integer syncNum = (Integer) batchResult.get("syncNum"); + if (object.getTotalRows() != null) { + object.setTotalRows(object.getTotalRows() + syncNum); + } else { + object.setTotalRows(syncNum); + } + object.setSyncStatus(true); + object.setUpdateTime(DateUtils.getNowDate()); + updateDataiIntegrationObject(object); + + result.put("success", true); + result.put("message", "增量数据拉取成功"); + result.put("objectId", object.getId()); + result.put("objectApi", objectApi); + result.put("totalCount", syncNum); + result.put("duration", duration); + result.put("syncType", "incremental"); + result.put("lastFullSyncDate", object.getLastFullSyncDate()); + result.put("lastSyncDate", object.getLastSyncDate()); + + log.info("对象 {} 增量数据拉取成功,同步数据量: {},耗时: {}ms", objectApi, syncNum, duration); + } else { + object.setSyncStatus(false); + object.setUpdateTime(DateUtils.getNowDate()); + updateDataiIntegrationObject(object); + + result.put("success", false); + result.put("message", "增量数据拉取失败: " + batchResult.get("message")); + + log.error("对象 {} 增量数据拉取失败: {}", objectApi, batchResult.get("message")); + } + + } catch (Exception e) { + log.error("增量数据拉取失败,对象API: {}", objectApi, e); + result.put("success", false); + result.put("message", "增量数据拉取失败: " + e.getMessage()); + } + + return result; + } + + private List getSalesforceObjectFields(PartnerConnection connection, String objectApi) throws ConnectionException + { + List fields = new ArrayList<>(); + + DescribeSObjectResult describeResult = connection.describeSObject(objectApi); + Field[] objectFields = describeResult.getFields(); + + for (Field field : objectFields) { + if (field.getType() != null && "base64".equalsIgnoreCase(field.getType().toString())) { + continue; + } + fields.add(field.getName()); + } + + if (!fields.contains("Id")) { + fields.add("Id"); + } + if (!fields.contains("Name")) { + fields.add("Name"); + } + + return fields; + } + + private int executeQueryAndProcessData(PartnerConnection connection, SalesforceParam param, List fieldList) + { + int totalCount = 0; + int failCount = 0; + int batch = 0; + Date lastCreatedDate = null; + String maxId = null; + final int MAX_FAIL_COUNT = 3; + boolean isFirstQuery = true; + + try { + while (true) { + try { + String query = buildDynamicQuery(param, lastCreatedDate, maxId, isFirstQuery); + log.info("执行查询,批次: {}, SQL: {}", batch, query); + + QueryResult queryResult = connection.query(query); + + if (queryResult == null || queryResult.getRecords() == null || queryResult.getRecords().length == 0) { + log.info("批次 {} 查询完成,无更多数据", batch); + break; + } + + int currentBatchCount = processQueryResult(param.getApi(), queryResult, fieldList); + totalCount += currentBatchCount; + + SObject[] records = queryResult.getRecords(); + lastCreatedDate = getLastCreatedDate(records, param.getDateField()); + maxId = getMaxId(records, param.getDateField(), lastCreatedDate); + + log.info("批次 {} 处理完成,本批次记录数: {}, 总记录数: {}, 最后时间: {}, 最大ID: {}", + batch, currentBatchCount, totalCount, lastCreatedDate, maxId); + + failCount = 0; + batch++; + isFirstQuery = false; + + }catch (Exception e) { + failCount++; + log.error("批次 {} 查询失败,失败次数: {}/{},错误信息: {}", + batch, failCount, MAX_FAIL_COUNT, e.getMessage(), e); + + if (failCount > MAX_FAIL_COUNT) { + log.error("达到最大失败次数 {},停止查询", MAX_FAIL_COUNT); + throw new RuntimeException("查询失败次数过多: " + e.getMessage(), e); + } + + } + } + + log.info("查询完成,总批次: {}, 总记录数: {}", batch, totalCount); + + } catch (Exception e) { + log.error("执行查询并处理数据失败: {}", e.getMessage(), e); + } + + return totalCount; + } + + /** + * 从查询结果中获取最后一条记录的日期字段值 + * + * @param records Salesforce查询返回的记录数组 + * @param dateField 日期字段名称,用于提取时间戳 + * @return 最后一条记录的日期字段值,如果获取失败则返回null + */ + private Date getLastCreatedDate(SObject[] records, String dateField) + { + // 检查记录数组是否为空 + if (records == null || records.length == 0) { + return null; + } + + // 获取最后一条记录 + SObject lastRecord = records[records.length - 1]; + if (lastRecord == null) { + return null; + } + + try { + // 提取指定的日期字段值 + Object dateValue = lastRecord.getField(dateField); + if (dateValue instanceof Calendar) { + // 将Calendar类型转换为Date类型 + return ((Calendar) dateValue).getTime(); + } + } catch (Exception e) { + log.warn("获取最后创建时间失败: {}", e.getMessage()); + } + + return null; + } + + /** + * 从查询结果中获取指定日期对应的最大ID值 + * + * @param records Salesforce查询返回的记录数组 + * @param dateField 日期字段名称,用于比较日期值 + * @param lastCreatedDate 日期参考值,只查找该日期的记录 + * @return 指定日期对应的最大ID值,如果没有匹配记录则返回null + */ + private String getMaxId(SObject[] records, String dateField, Date lastCreatedDate) + { + // 检查记录数组和日期值是否为空,避免空指针异常 + if (records == null || records.length == 0 || lastCreatedDate == null) { + return null; + } + + // 初始化最大ID变量,用于存储遍历过程中找到的最大ID值 + String maxId = null; + + // 遍历所有记录,查找日期字段值等于lastCreatedDate的记录 + for (SObject record : records) { + // 跳过空记录,避免空指针异常 + if (record == null) { + continue; + } + + try { + // 提取记录的日期字段值 + Object dateValue = record.getField(dateField); + if (dateValue instanceof Calendar) { + // 将Calendar类型转换为Date类型,便于比较 + Date recordDate = ((Calendar) dateValue).getTime(); + // 只处理与参考日期相同的记录 + if (recordDate.equals(lastCreatedDate)) { + // 获取记录的ID + String recordId = record.getId(); + // 更新最大ID:如果当前记录ID大于已知的最大ID,则更新 + if (recordId != null && (maxId == null || recordId.compareTo(maxId) > 0)) { + maxId = recordId; + } + } + } + } catch (Exception e) { + // 字段获取失败时记录警告日志,继续处理其他记录 + log.warn("获取记录ID失败: {}", e.getMessage()); + } + } + + // 返回指定日期对应的最大ID值,该值将作为下一批次查询的起始ID + return maxId; + } + + /** + * 构建动态SOQL查询语句 + * + * @param param Salesforce查询参数对象 + * @param lastCreatedDate 最后创建日期,用于增量查询 + * @param maxId 最大ID值,用于分页查询 + * @param isFirstQuery 是否为首次查询 + * @return 构建好的SOQL查询语句 + */ + private String buildDynamicQuery(SalesforceParam param, Date lastCreatedDate, String maxId, boolean isFirstQuery) + { + // 使用StringBuilder构建查询语句,提高字符串拼接效率 + StringBuilder queryBuilder = new StringBuilder(); + // 构建SELECT和FROM子句 + queryBuilder.append("SELECT ").append(param.getSelect()) + .append(" FROM ").append(param.getApi()); + + // 创建条件列表,用于存储所有WHERE条件 + List conditions = new ArrayList<>(); + + // 如果配置了查询已删除记录,添加IsDeleted条件 + if (param.getIsDeleted() != null && param.getIsDeleted()) { + conditions.add("IsDeleted = true"); + } + + // 首次查询时,如果配置了开始修改日期,添加日期范围条件 + if (isFirstQuery && param.getBeginModifyDate() != null) { + String dateField = param.getDateField(); + if (StringUtils.isNotEmpty(dateField)) { + // 创建UTC时区的日期格式化器 + SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss'.000Z'"); + sdf.setTimeZone(TimeZone.getTimeZone("UTC")); + // 格式化日期为Salesforce支持的格式 + String dateStr = sdf.format(param.getBeginModifyDate()); + // 添加日期大于等于条件 + conditions.add(dateField + " >= " + dateStr); + } + } + + // 非首次查询时,使用最后创建日期作为查询条件 + if (!isFirstQuery && lastCreatedDate != null) { + String dateField = param.getDateField(); + if (StringUtils.isNotEmpty(dateField)) { + // 创建UTC时区的日期格式化器 + SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss'.000Z'"); + sdf.setTimeZone(TimeZone.getTimeZone("UTC")); + // 格式化日期为Salesforce支持的格式 + String dateStr = sdf.format(lastCreatedDate); + // 添加日期大于等于条件 + conditions.add(dateField + " >= " + dateStr); + } + } + + // 如果配置了最大ID,添加ID大于条件,实现精确分页 + if (StringUtils.isNotEmpty(maxId)) { + conditions.add("Id > '" + maxId + "'"); + } + + // 如果有条件,构建WHERE子句 + if (!conditions.isEmpty()) { + queryBuilder.append(" WHERE ").append(String.join(" AND ", conditions)); + } + + // 添加ORDER BY子句,按日期字段和ID排序,确保数据顺序一致 + queryBuilder.append(" ORDER BY ").append(param.getDateField()).append(", Id"); + + // 如果配置了查询限制,添加LIMIT子句 + if (param.getLimit() != null && param.getLimit() > 0) { + queryBuilder.append(" LIMIT ").append(param.getLimit()); + } + + // 返回构建好的SOQL查询语句 + return queryBuilder.toString(); + } + + /** + * 处理查询结果并将数据保存到数据库 + * + * @param api Salesforce对象的API名称,用于标识目标表 + * @param result Salesforce查询返回的结果集,包含所有需要同步的记录 + * @param fieldList 需要同步的字段列表,用于控制数据转换和存储的字段范围 + * @return 成功处理的记录数量 + */ + private int processQueryResult(String api, QueryResult result, List fieldList) + { + // 初始化计数器,用于统计成功处理的记录数量 + int count = 0; + + // 检查查询结果是否为空,避免空指针异常 + if (result == null || result.getRecords() == null) { + log.info("处理API {} 的查询结果,共 0 条记录", api); + return count; + } + + // 获取查询结果中的所有记录 + SObject[] records = result.getRecords(); + log.info("处理API {} 的查询结果,共 {} 条记录", api, records.length); + + // 检查目标表是否为分区表,决定数据保存策略 + boolean isPartitioned = checkIfTablePartitioned(api); + // 初始化批次字段变量,用于分区表的数据分配 + String batchField = null; + + // 如果是分区表,获取批次字段配置 + if (isPartitioned) { + // 根据API名称查询对象配置 + DataiIntegrationObject object = dataiIntegrationObjectMapper.selectDataiIntegrationObjectByApi(api); + if (object != null && StringUtils.isNotEmpty(object.getBatchField())) { + // 使用对象配置的批次字段 + batchField = object.getBatchField(); + log.info("对象 {} 配置了批次字段: {}", api, batchField); + } else { + // 如果对象未配置批次字段,使用日期字段作为默认批次字段 + batchField = dataiIntegrationFieldService.getDateField(api); + log.info("对象 {} 未配置批次字段,使用日期字段: {}", api, batchField); + } + } + + // 创建分区数据Map,用于存储按分区名分组的数据 + Map>> partitionedData = new HashMap<>(); + // 创建普通数据列表,用于存储非分区表的数据 + List> normalData = new ArrayList<>(); + + // 遍历所有记录,进行数据转换和分区分配 + for (SObject record : records) { + // 跳过空记录,避免空指针异常 + if (record == null) { + continue; + } + + try { + // 将SObject对象转换为Map格式,便于后续处理 + Map recordMap = convertSObjectToMap(record, fieldList); + + // 如果是分区表且配置了批次字段,进行分区分配 + if (isPartitioned && StringUtils.isNotEmpty(batchField)) { + // 初始化分区名称为默认分区 + String partitionName = "p_default"; + + // 检查记录中是否包含批次字段 + if (recordMap.containsKey(batchField.toLowerCase())) { + // 获取批次字段的值 + Object batchValue = recordMap.get(batchField.toLowerCase()); + if (batchValue instanceof Date) { + // 如果批次字段值为Date类型,使用年份作为分区标识(如p2025) + Calendar calendar = Calendar.getInstance(); + calendar.setTime((Date) batchValue); + int year = calendar.get(Calendar.YEAR); + partitionName = "p" + year; + log.debug("记录根据批次字段 {} 的值 {} 分配到分区 {}", batchField, batchValue, partitionName); + } else if (batchValue != null) { + // 如果批次字段值为其他类型,使用哈希值作为分区标识 + partitionName = "p_" + String.valueOf(batchValue).hashCode(); + log.debug("记录根据批次字段 {} 的值 {} 分配到分区 {}", batchField, batchValue, partitionName); + } + } + + // 将记录添加到对应分区的数据列表中 + partitionedData.computeIfAbsent(partitionName, k -> new ArrayList<>()).add(recordMap); + } else { + // 非分区表,直接添加到普通数据列表 + normalData.add(recordMap); + } + + // 增加处理计数 + count++; + + } catch (Exception e) { + // 记录处理异常,但继续处理其他记录 + String recordId = record.getId() != null ? record.getId() : "未知"; + log.error("处理记录时发生异常,记录ID: {}", recordId, e); + } + } + + try { + // 如果是分区表且有分区数据,批量保存到各个分区 + if (isPartitioned && !partitionedData.isEmpty()) { + log.info("开始批量保存分区数据,共 {} 个分区", partitionedData.size()); + // 遍历所有分区,保存数据 + for (Map.Entry>> entry : partitionedData.entrySet()) { + String partitionName = entry.getKey(); + List> dataList = entry.getValue(); + // 保存数据到指定分区 + saveDataToPartition(api, partitionName, dataList); + log.info("成功保存 {} 条记录到分区 {}", dataList.size(), partitionName); + } + } else if (!normalData.isEmpty()) { + // 非分区表,直接保存数据到目标表 + saveDataToTable(api, normalData); + log.info("成功保存 {} 条记录到表 {}", normalData.size(), api); + } + } catch (Exception e) { + // 记录批量保存失败的异常 + log.error("批量保存数据到数据库失败: {}", e.getMessage(), e); + } + + // 返回成功处理的记录数量 + return count; + } + + /** + * 将Salesforce的SObject对象转换为Map格式 + * + * @param record Salesforce的SObject对象,包含从Salesforce查询返回的单条记录数据 + * @param fieldList 需要提取的字段列表,控制转换过程中包含哪些字段 + * @return 包含所有指定字段值的Map对象,字段名统一为小写格式 + */ + private Map convertSObjectToMap(SObject record, List fieldList) + { + // 创建Map对象用于存储转换后的记录数据 + Map recordMap = new HashMap<>(); + + // 添加记录ID字段,Id字段总是被包含在结果中 + if (record.getId() != null) { + recordMap.put("Id", record.getId()); + } + + // 遍历所有需要提取的字段 + for (String field : fieldList) { + // 跳过Id字段,因为已经单独处理过 + if ("Id".equalsIgnoreCase(field)) { + continue; + } + + try { + // 从SObject中获取字段值 + Object value = record.getField(field); + if (value != null) { + // 检查字段值是否为Calendar类型(日期类型) + if (value instanceof java.util.Calendar) { + // 将Calendar类型转换为Date类型,便于数据库存储 + recordMap.put(field.toLowerCase(), ((Calendar) value).getTime()); + } else { + // 其他类型直接存储,字段名转换为小写以适应数据库命名规范 + recordMap.put(field.toLowerCase(), value); + } + } else { + // 空值字段也保留在结果中,值为null + recordMap.put(field.toLowerCase(), null); + } + } catch (Exception e) { + // 字段获取异常时记录警告日志,但不中断处理流程 + log.warn("获取字段 {} 的值时发生异常", field, e); + } + } + + // 返回转换后的Map对象 + return recordMap; + } + + /** + * 检查指定表是否为分区表 + * + * @param tableName 需要检查的表名,对应Salesforce对象的API名称 + * @return 如果表是分区表返回true,否则返回false + */ + private boolean checkIfTablePartitioned(String tableName) + { + try { + // 调用CustomMapper查询数据库元数据,判断表是否为分区表 + return customMapper.isPartitioned(tableName); + } catch (Exception e) { + // 查询失败时记录警告日志,默认按非分区表处理 + log.warn("检查表 {} 是否分区失败: {}", tableName, e.getMessage()); + return false; + } + } + + /** + * 批量保存数据到非分区表 + * + * @param tableName 目标表名,对应Salesforce对象的API名称 + * @param dataList 需要保存的数据列表,每条记录是一个Map,键为字段名,值为字段值 + */ + private void saveDataToTable(String tableName, List> dataList) + { + // 检查数据列表是否为空,为空则直接返回 + if (dataList == null || dataList.isEmpty()) { + return; + } + + try { + // 创建字段名列表,用于存储表的所有列名 + List keys = new ArrayList<>(); + // 创建字段值列表,用于存储所有记录的字段值 + List> values = new ArrayList<>(); + + // 遍历所有数据记录 + for (Map dataMap : dataList) { + // 如果是第一条记录,提取所有字段名作为列名 + if (keys.isEmpty()) { + keys.addAll(dataMap.keySet()); + } + // 添加当前记录的字段值到值列表 + values.add(dataMap.values()); + } + + // 调用CustomMapper执行批量插入操作 + customMapper.saveBatch(tableName, keys, values); + log.info("成功保存 {} 条记录到表 {}", dataList.size(), tableName); + } catch (Exception e) { + // 记录错误日志并重新抛出异常,由调用方处理 + log.error("批量保存数据到表 {} 失败: {}", tableName, e.getMessage(), e); + throw e; + } + } + + /** + * 批量保存数据到分区表的指定分区 + * + * @param tableName 目标表名,对应Salesforce对象的API名称 + * @param partitionName 分区名称,标识数据应该保存到哪个分区(如p2025) + * @param dataList 需要保存的数据列表,每条记录是一个Map,键为字段名,值为字段值 + */ + private void saveDataToPartition(String tableName, String partitionName, List> dataList) + { + // 检查数据列表是否为空,为空则直接返回 + if (dataList == null || dataList.isEmpty()) { + return; + } + + try { + // 创建字段名列表,用于存储表的所有列名 + List keys = new ArrayList<>(); + // 创建字段值列表,用于存储所有记录的字段值 + List> values = new ArrayList<>(); + + // 遍历所有数据记录 + for (Map dataMap : dataList) { + // 如果是第一条记录,提取所有字段名作为列名 + if (keys.isEmpty()) { + keys.addAll(dataMap.keySet()); + } + // 添加当前记录的字段值到值列表 + values.add(dataMap.values()); + } + + // 调用CustomMapper执行批量插入操作到指定分区 + customMapper.saveBatchToPartition(tableName, partitionName, keys, values); + log.info("成功保存 {} 条记录到分区 {}", dataList.size(), partitionName); + } catch (Exception e) { + // 记录错误日志并重新抛出异常,由调用方处理 + log.error("批量保存数据到分区 {} 失败: {}", partitionName, e.getMessage(), e); + throw e; + } + } } diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/doc/同步单对象数据接口文档.md b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/doc/同步单对象数据接口文档.md new file mode 100644 index 00000000..0e8e44a6 --- /dev/null +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/doc/同步单对象数据接口文档.md @@ -0,0 +1,671 @@ +# 同步单对象数据接口文档 + +## 1. 接口概述 + +### 1.1 接口名称 +同步单对象数据到本地数据库 + +### 1.2 接口描述 +该接口用于触发指定对象的单次数据同步操作,将Salesforce对象的数据同步到本地数据库。同步操作包括全量同步和增量同步两种模式,具体模式由对象的配置决定: +- **全量同步**:当对象的全量拉取时间(lastFullSyncDate)为空时,执行全量数据拉取,查询所有全量同步批次并使用多线程并行拉取 +- **增量同步**:当对象的全量拉取时间不为空时,执行增量数据拉取,创建增量批次并拉取增量数据 + +### 1.3 接口地址 +``` +POST /integration/object/{id}/syncData +``` + +### 1.4 请求方式 +POST + +### 1.5 权限要求 +需要 `integration:object:syncData` 权限 + +### 1.6 日志记录 +- 操作类型:UPDATE(更新) +- 日志标题:对象同步控制 + +## 2. 请求参数 + +### 2.1 请求头 +``` +Content-Type: application/json +Authorization: Bearer {token} +``` + +### 2.2 路径参数 + +| 参数名 | 类型 | 是否必填 | 描述 | 示例值 | +|--------|------|----------|------|--------| +| id | Integer | 是 | 对象ID,用于标识需要同步的Salesforce对象 | 1 | + +### 2.3 请求体 +该接口无需请求体参数 + +### 2.4 请求示例 + +#### 2.4.1 cURL 示例 +```bash +curl -X POST "http://localhost:8080/integration/object/1/syncData" \ + -H "Content-Type: application/json" \ + -H "Authorization: Bearer {your_token}" +``` + +#### 2.4.2 HTTP 示例 +```http +POST /integration/object/1/syncData HTTP/1.1 +Host: localhost:8080 +Content-Type: application/json +Authorization: Bearer {your_token} +``` + +#### 2.4.3 JavaScript (fetch) 示例 +```javascript +fetch('http://localhost:8080/integration/object/1/syncData', { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'Authorization': 'Bearer {your_token}' + } +}) + .then(response => response.json()) + .then(data => console.log(data)) + .catch(error => console.error('Error:', error)); +``` + +## 3. 响应格式 + +### 3.1 响应结构 +```json +{ + "code": 200, + "msg": "操作成功", + "data": { + "success": true, + "message": "全量数据拉取完成", + "objectId": 1, + "objectApi": "Account", + "totalCount": 12345, + "duration": 15234, + "syncType": "full", + "lastFullSyncDate": "2025-12-27T10:30:00", + "totalBatchCount": 5, + "successBatchCount": 5, + "failedBatchCount": 0 + } +} +``` + +### 3.2 响应字段说明 + +#### 3.2.1 通用响应字段 + +| 字段名 | 类型 | 说明 | +|--------|------|------| +| code | Integer | 响应状态码,200表示成功,其他表示失败 | +| msg | String | 响应消息 | +| data | Object | 响应数据对象,失败时为null | + +#### 3.2.2 业务数据字段(data对象) + +| 字段名 | 类型 | 说明 | +|--------|------|------| +| success | Boolean | 操作是否成功,true表示成功,false表示失败 | +| message | String | 详细消息说明 | +| objectId | Integer | 对象ID | +| objectApi | String | 对象API名称 | +| totalCount | Integer | 同步的总记录数 | +| duration | Long | 同步耗时(毫秒) | +| syncType | String | 同步类型,"full"表示全量同步,"incremental"表示增量同步 | +| lastFullSyncDate | LocalDateTime | 最后全量同步时间 | +| lastSyncDate | LocalDateTime | 最后同步时间(增量同步时返回) | +| totalBatchCount | Integer | 总批次数量(全量同步时返回) | +| successBatchCount | Integer | 成功批次数量(全量同步时返回) | +| failedBatchCount | Integer | 失败批次数量(全量同步时返回) | + +### 3.3 成功响应示例 + +#### 3.3.1 全量同步成功 +```json +{ + "code": 200, + "msg": "操作成功", + "data": { + "success": true, + "message": "全量数据拉取完成", + "objectId": 1, + "objectApi": "Account", + "totalCount": 12345, + "duration": 15234, + "syncType": "full", + "lastFullSyncDate": "2025-12-27T10:30:00", + "totalBatchCount": 5, + "successBatchCount": 5, + "failedBatchCount": 0 + } +} +``` + +#### 3.3.2 增量同步成功 +```json +{ + "code": 200, + "msg": "操作成功", + "data": { + "success": true, + "message": "增量数据拉取成功", + "objectId": 1, + "objectApi": "Account", + "totalCount": 234, + "duration": 3456, + "syncType": "incremental", + "lastFullSyncDate": "2025-12-26T10:30:00", + "lastSyncDate": "2025-12-27T10:30:00" + } +} +``` + +#### 3.3.3 部分批次失败(全量同步) +```json +{ + "code": 200, + "msg": "操作成功", + "data": { + "success": true, + "message": "全量数据拉取完成", + "objectId": 1, + "objectApi": "Account", + "totalCount": 10000, + "duration": 12000, + "syncType": "full", + "lastFullSyncDate": "2025-12-27T10:30:00", + "totalBatchCount": 5, + "successBatchCount": 4, + "failedBatchCount": 1 + } +} +``` + +### 3.4 失败响应示例 + +#### 3.4.1 对象ID为空 +```json +{ + "code": 500, + "msg": "对象ID不能为空", + "data": null +} +``` + +#### 3.4.2 对象不存在 +```json +{ + "code": 500, + "msg": "对象不存在", + "data": null +} +``` + +#### 3.4.3 对象API为空 +```json +{ + "code": 500, + "msg": "对象API不能为空", + "data": null +} +``` + +#### 3.4.4 对象未启用同步 +```json +{ + "code": 500, + "msg": "对象未启用同步", + "data": null +} +``` + +#### 3.4.5 没有找到全量同步批次 +```json +{ + "code": 500, + "msg": "没有找到全量同步批次", + "data": null +} +``` + +#### 3.4.6 对象未启用增量更新 +```json +{ + "code": 500, + "msg": "对象未启用增量更新", + "data": null +} +``` + +#### 3.4.7 无法获取Salesforce连接 +```json +{ + "code": 500, + "msg": "无法获取Salesforce连接", + "data": null +} +``` + +#### 3.4.8 无法确定批次字段 +```json +{ + "code": 500, + "msg": "无法确定批次字段", + "data": null +} +``` + +#### 3.4.9 同步异常 +```json +{ + "code": 500, + "msg": "同步对象数据时发生异常: Connection timeout", + "data": null +} +``` + +## 4. 业务逻辑说明 + +### 4.1 处理流程 + +#### 4.1.1 全量同步流程 +1. **参数校验**:检查对象ID是否为空 +2. **对象查询**:根据对象ID查询对象配置信息 +3. **对象验证**: + - 检查对象是否存在 + - 检查对象API是否配置 + - 检查对象是否启用同步 +4. **同步模式判断**:判断对象的全量拉取时间(lastFullSyncDate)是否为空 +5. **全量同步执行**: + - 查询所有 `syncType="FULL"` 的批次 + - 使用 `SalesforceExecutor` 线程池多线程并行拉取多个批次 + - 统计成功和失败的批次数量 + - 统计总同步记录数 +6. **对象状态更新**: + - 更新对象的 `lastFullSyncDate` 为当前时间 + - 更新对象的 `lastSyncDate` 为当前时间 + - 更新对象的 `totalRows` 为总记录数 + - 更新对象的 `syncStatus`(全部成功为true,否则为false) +7. **同步日志记录**:记录同步操作日志,包括操作类型、状态、执行时间等 +8. **返回结果**:返回同步结果,包括成功状态、同步数量、耗时等信息 + +#### 4.1.2 增量同步流程 +1. **参数校验**:检查对象ID是否为空 +2. **对象查询**:根据对象ID查询对象配置信息 +3. **对象验证**: + - 检查对象是否存在 + - 检查对象API是否配置 + - 检查对象是否启用同步 +4. **同步模式判断**:判断对象的全量拉取时间(lastFullSyncDate)是否为空 +5. **增量同步执行**: + - 检查对象是否启用增量更新(isIncremental) + - 获取Salesforce连接 + - 获取对象字段列表 + - 确定批次字段(batchField) + - 创建增量批次(syncType="INCREMENTAL") + - 设置批次的开始时间为上次全量同步时间,结束时间为当前时间 + - 调用批次同步服务拉取增量数据 +6. **对象状态更新**: + - 更新对象的 `lastSyncDate` 为当前时间 + - 更新对象的 `totalRows`(累加增量记录数) + - 更新对象的 `syncStatus` +7. **同步日志记录**:记录同步操作日志,包括操作类型、状态、执行时间等 +8. **返回结果**:返回同步结果,包括成功状态、同步数量、耗时等信息 + +### 4.2 同步模式判断规则 + +| 条件 | 同步模式 | 说明 | +|------|----------|------| +| lastFullSyncDate == null | 全量同步 | 对象首次同步,执行全量数据拉取 | +| lastFullSyncDate != null | 增量同步 | 对象已执行过全量同步,执行增量数据拉取 | + +### 4.3 全量同步特性 + +#### 4.3.1 批次查询 +- 查询条件:`api = objectApi` 且 `syncType = "FULL"` +- 返回所有匹配的全量同步批次 + +#### 4.3.2 多线程并行拉取 +- 使用 `SalesforceExecutor` 线程池执行批次同步任务 +- 每个批次作为一个独立的任务提交到线程池 +- 使用 `CopyOnWriteArrayList` 线程安全集合存储批次同步结果 +- 等待所有批次任务完成后统计结果 + +#### 4.3.3 结果统计 +- **totalBatchCount**:总批次数量 +- **successBatchCount**:成功批次数量 +- **failedBatchCount**:失败批次数量 +- **totalCount**:所有成功批次同步的记录数总和 + +### 4.4 增量同步特性 + +#### 4.4.1 增量更新检查 +- 必须启用增量更新(isIncremental = true) +- 如果未启用增量更新,返回失败 + +#### 4.4.2 批次字段确定 +1. 优先使用对象配置的 `batchField` +2. 如果未配置,使用 `dataiIntegrationFieldService.getDateField(objectApi)` 获取默认日期字段 +3. 如果无法确定批次字段,返回失败 + +#### 4.4.3 增量批次创建 +- `api`:对象API名称 +- `label`:对象标签名称 +- `batchField`:批次字段 +- `syncType`:INCREMENTAL +- `syncStatus`:false +- `syncStartDate`:上次全量同步时间 +- `syncEndDate`:当前时间 +- `createTime`:当前时间 +- `updateTime`:当前时间 + +#### 4.4.4 记录数累加 +- 增量同步成功后,将增量记录数累加到对象的 `totalRows` 中 +- 如果对象的 `totalRows` 为空,直接设置为增量记录数 + +### 4.5 同步日志记录 + +#### 4.5.1 日志字段 +- `objectApi`:对象API名称 +- `operationType`:操作类型("全量同步" 或 "增量同步") +- `operationStatus`:操作状态("成功"、"失败" 或 "部分失败") +- `errorMessage`:错误信息(失败时) +- `executionTime`:执行时间(秒) +- `deptId`:部门ID(从当前登录用户获取) + +#### 4.5.2 日志记录时机 +- 成功时:记录成功日志 +- 失败时:记录失败日志,包含错误信息 +- 异常时:记录异常日志,包含异常信息 + +## 5. 错误处理 + +### 5.1 常见错误码 + +| 错误码 | 说明 | 处理建议 | +|--------|------|----------| +| 401 | 未授权 | 检查token是否有效 | +| 403 | 权限不足 | 确认用户是否具有 `integration:object:syncData` 权限 | +| 500 | 服务器内部错误 | 查看服务器日志,检查具体错误信息 | + +### 5.2 常见错误场景 + +#### 5.2.1 对象ID为空 +``` +错误信息:对象ID不能为空 +原因:请求路径中的id参数未提供或为null +解决方案:确保请求路径中包含有效的对象ID +``` + +#### 5.2.2 对象不存在 +``` +错误信息:对象不存在 +原因:根据提供的对象ID未查询到对应的对象配置 +解决方案:检查对象ID是否正确,确认对象是否已创建 +``` + +#### 5.2.3 对象API为空 +``` +错误信息:对象API不能为空 +原因:对象的API字段为空或未配置 +解决方案:检查对象配置,确保API字段已正确配置 +``` + +#### 5.2.4 对象未启用同步 +``` +错误信息:对象未启用同步 +原因:对象的isWork字段为false +解决方案:先调用更新对象启用同步状态接口,将isWork设置为true +``` + +#### 5.2.5 没有找到全量同步批次 +``` +错误信息:没有找到全量同步批次 +原因:对象没有配置全量同步批次 +解决方案:先为对象创建全量同步批次 +``` + +#### 5.2.6 对象未启用增量更新 +``` +错误信息:对象未启用增量更新 +原因:对象的isIncremental字段为false +解决方案:先调用更新对象增量更新状态接口,将isIncremental设置为true +``` + +#### 5.2.7 无法获取Salesforce连接 +``` +错误信息:无法获取Salesforce连接 +原因:SOAP连接配置错误或网络问题 +解决方案:检查SOAPConnectionFactory配置和网络连接 +``` + +#### 5.2.8 无法确定批次字段 +``` +错误信息:无法确定批次字段 +原因:对象未配置批次字段,且无法获取默认日期字段 +解决方案:为对象配置批次字段,或确保对象包含日期字段 +``` + +#### 5.2.9 同步异常 +``` +错误信息:同步对象数据时发生异常: {具体异常信息} +原因:同步过程中发生异常,如网络超时、数据格式错误等 +解决方案:查看服务器日志,根据具体异常信息进行排查 +``` + +### 5.3 错误处理机制 + +#### 5.3.1 参数校验 +- 在方法开始时进行参数校验 +- 校验失败立即返回错误信息 + +#### 5.3.2 异常捕获 +- 使用try-catch捕获所有异常 +- 异常发生时记录错误日志 +- 更新对象的错误状态和错误信息 +- 记录同步日志(操作状态为"失败") + +#### 5.3.3 状态回滚 +- 异常发生时,将对象的syncStatus设置为false +- 记录错误信息到对象的errorMessage字段 + +## 6. 使用注意事项 + +### 6.1 性能考虑 + +#### 6.1.1 全量同步 +- 全量同步会拉取所有批次的数据,数据量可能很大 +- 使用多线程并行拉取,可以提高同步效率 +- 建议在业务低峰期执行全量同步 +- 监控同步耗时,避免超时 + +#### 6.1.2 增量同步 +- 增量同步只拉取增量数据,数据量相对较小 +- 增量同步的耗时取决于增量数据量 +- 建议定期执行增量同步,保持数据最新 + +### 6.2 数据一致性 + +#### 6.2.1 批次同步 +- 每个批次的同步是独立的 +- 部分批次失败不影响其他批次 +- 全量同步时,即使部分批次失败,成功的批次数据也会保存 +- 建议检查failedBatchCount,必要时重新同步失败的批次 + +#### 6.2.2 对象状态 +- 同步成功后,对象的同步状态(syncStatus)会更新 +- 部分批次失败时,syncStatus会设置为false +- 增量同步失败时,不会影响已同步的数据 + +### 6.3 并发控制 + +#### 6.3.1 同一对象 +- 不建议对同一对象同时发起多次同步请求 +- 可能导致数据重复或状态不一致 +- 建议等待前一次同步完成后再发起下一次同步 + +#### 6.3.2 不同对象 +- 可以对不同对象同时发起同步请求 +- 使用线程池控制并发数 +- 监控线程池状态,避免资源耗尽 + +### 6.4 监控建议 + +#### 6.4.1 同步状态监控 +- 监控对象的syncStatus字段 +- 监控对象的errorMessage字段 +- 定期检查同步日志 + +#### 6.4.2 性能监控 +- 监控同步耗时(duration) +- 监控同步记录数(totalCount) +- 监控批次成功率(successBatchCount / totalBatchCount) + +#### 6.4.3 异常监控 +- 监控同步异常日志 +- 监控失败批次数量 +- 及时处理异常情况 + +### 6.5 最佳实践 + +#### 6.5.1 首次同步 +- 首次同步时,确保对象已配置全量同步批次 +- 确保对象已启用同步(isWork = true) +- 在业务低峰期执行首次同步 + +#### 6.5.2 定期同步 +- 建议定期执行增量同步 +- 根据业务需求确定同步频率 +- 监控增量数据量,调整同步策略 + +#### 6.5.3 错误处理 +- 同步失败时,检查错误信息 +- 根据错误类型采取相应措施 +- 必要时重新执行同步 + +#### 6.5.4 日志管理 +- 定期检查同步日志 +- 分析同步趋势和异常情况 +- 优化同步策略 + +## 7. 相关接口 + +### 7.1 查询对象同步控制列表 +``` +GET /integration/object/list +``` + +### 7.2 获取对象同步控制详细信息 +``` +GET /integration/object/{id} +``` + +### 7.3 获取对象同步统计信息 +``` +GET /integration/object/{id}/statistics +``` + +### 7.4 变更对象启用同步状态 +``` +PUT /integration/object/{id}/workStatus?isWork=true +``` + +### 7.5 变更对象增量更新状态 +``` +PUT /integration/object/{id}/incrementalStatus?isIncremental=true +``` + +### 7.6 创建对象表结构 +``` +POST /integration/object/{id}/createStructure +``` + +### 7.7 获取对象依赖关系 +``` +GET /integration/object/{id}/dependencies +``` + +### 7.8 同步批次数据 +``` +POST /integration/batch/{id}/syncData +``` + +## 8. 附录 + +### 8.1 对象同步控制表结构(DataiIntegrationObject) + +| 字段名 | 类型 | 说明 | +|--------|------|------| +| id | Integer | 主键ID | +| api | String | 对象API名称 | +| label | String | 对象标签名称 | +| isWork | Boolean | 是否启用同步 | +| isIncremental | Boolean | 是否启用增量更新 | +| batchField | String | 批次字段 | +| lastFullSyncDate | LocalDateTime | 最后全量同步时间 | +| lastSyncDate | LocalDateTime | 最后同步时间 | +| totalRows | Integer | 总记录数 | +| syncStatus | Boolean | 同步状态 | +| errorMessage | String | 错误信息 | +| createTime | Date | 创建时间 | +| updateTime | Date | 更新时间 | +| createBy | String | 创建人 | +| updateBy | String | 更新人 | + +### 8.2 批次同步表结构(DataiIntegrationBatch) + +| 字段名 | 类型 | 说明 | +|--------|------|------| +| id | Integer | 主键ID | +| api | String | 对象API名称 | +| label | String | 对象标签名称 | +| batchField | String | 批次字段 | +| syncType | String | 同步类型(FULL/INCREMENTAL) | +| syncStatus | Boolean | 同步状态 | +| syncStartDate | LocalDateTime | 同步开始时间 | +| syncEndDate | LocalDateTime | 同步结束时间 | +| createTime | Date | 创建时间 | +| updateTime | Date | 更新时间 | + +### 8.3 同步日志表结构(DataiIntegrationSyncLog) + +| 字段名 | 类型 | 说明 | +|--------|------|------| +| id | Long | 主键ID | +| objectApi | String | 对象API名称 | +| operationType | String | 操作类型 | +| operationStatus | String | 操作状态 | +| errorMessage | String | 错误信息 | +| executionTime | BigDecimal | 执行时间(秒) | +| deptId | Long | 部门ID | +| createTime | Date | 创建时间 | + +### 8.4 SalesforceExecutor 线程池配置 + +| 配置项 | 默认值 | 说明 | +|--------|--------|------| +| corePoolSize | CPU核心数 | 核心线程数 | +| maxPoolSize | CPU核心数 * 2 | 最大线程数 | +| keepAliveTime | 60秒 | 线程存活时间 | +| queueCapacity | 100 | 队列容量 | +| allowCoreThreadTimeout | false | 是否允许核心线程超时 | + +### 8.5 版本历史 + +| 版本 | 日期 | 说明 | +|------|------|------| +| 1.0 | 2025-12-27 | 初始版本 | + +### 8.6 联系方式 +如有问题,请联系技术支持团队。 + +--- + +**文档生成时间**:2025-12-27 +**最后更新时间**:2025-12-27 +**文档版本**:1.0 diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/DataiIntegrationObjectMapper.xml b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/DataiIntegrationObjectMapper.xml index 1aae082f..14b7153f 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/DataiIntegrationObjectMapper.xml +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/resources/mapper/integration/DataiIntegrationObjectMapper.xml @@ -106,6 +106,11 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" where dio.id = #{id} + + insert into datai_integration_object diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/ComparableFutureTask.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/ComparableFutureTask.java index 99b03dc3..5cf8ab12 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/ComparableFutureTask.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/ComparableFutureTask.java @@ -1,14 +1,17 @@ package com.datai.setting.future; import lombok.Getter; +import org.springframework.security.core.context.SecurityContext; +import org.springframework.security.core.context.SecurityContextHolder; import java.util.concurrent.FutureTask; /** - * 自定义FutureTask,支持优先级排序 + * 自定义FutureTask,支持优先级排序和安全上下文传递 *

* 该类扩展了FutureTask,实现了Comparable接口,可以根据优先级对任务进行排序。 * 优先级规则:index值越大优先级越高;index相等时,batch值越小优先级越高。 + * 同时支持将父线程的Spring Security上下文传递到子线程中。 *

*/ @Getter @@ -34,6 +37,11 @@ public class ComparableFutureTask extends FutureTask implements Comparab */ private final Integer batch; + /** + * 父线程的Spring Security上下文 + */ + private final SecurityContext parentSecurityContext; + /** * 构造函数 * @@ -45,6 +53,21 @@ public class ComparableFutureTask extends FutureTask implements Comparab super(runnable, null); this.index = index; this.batch = batch; + this.parentSecurityContext = SecurityContextHolder.getContext(); + } + + @Override + public void run() { + SecurityContext originalContext = null; + try { + originalContext = SecurityContextHolder.getContext(); + if (parentSecurityContext != null) { + SecurityContextHolder.setContext(parentSecurityContext); + } + super.run(); + } finally { + SecurityContextHolder.setContext(originalContext); + } } /** diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/SalesforceExecutor.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/SalesforceExecutor.java index 0321c438..7fe423d1 100644 --- a/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/SalesforceExecutor.java +++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-setting/src/main/java/com/datai/setting/future/SalesforceExecutor.java @@ -17,7 +17,13 @@ import java.util.concurrent.*; * 该类提供了Salesforce数据处理任务的线程池管理功能,支持任务优先级排序和批量处理。 * 使用PriorityBlockingQueue作为任务队列,确保高优先级任务优先执行。 *

- * + *

+ * 特性: + * 1. 支持任务优先级排序(基于ComparableFutureTask) + * 2. 自动将父线程的Spring Security上下文传递到子线程 + * 3. 可配置的线程池参数(核心线程数、最大线程数、队列容量等) + * 4. 支持优雅关闭和强制关闭 + *

*/ @Service @Slf4j @@ -180,6 +186,10 @@ public class SalesforceExecutor { *

* 将任务提交到线程池执行,支持优先级设置。 *

+ *

+ * 注意:该方法会自动将当前线程的Spring Security上下文传递到子线程中, + * 确保子线程可以访问父线程的安全信息(如登录用户信息)。 + *

* * @param runnable 要执行的任务 * @param batch 批次号(优先级相同时,值越小优先级越高)