【feat】 对象数据同步本地数据库

This commit is contained in:
Kris 2026-01-06 15:14:02 +08:00
parent cc3a366e01
commit a78da492c7
11 changed files with 1735 additions and 20 deletions

View File

@ -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)
与本项目同为一个作者开发,兼容性最好,学习成本最低。

View File

@ -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");

View File

@ -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<String, Object> 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());
}
}
}

View File

@ -20,6 +20,14 @@ public interface DataiIntegrationObjectMapper
*/
public DataiIntegrationObject selectDataiIntegrationObjectById(Integer id);
/**
* 根据API查询对象同步控制
*
* @param api 对象API
* @return 对象同步控制
*/
public DataiIntegrationObject selectDataiIntegrationObjectByApi(String api);
/**
* 查询对象同步控制列表
*

View File

@ -108,4 +108,12 @@ public interface IDataiIntegrationObjectService
* @return 统计信息
*/
public Map<String, Object> getObjectStatistics();
/**
* 同步单对象数据到本地数据库
*
* @param id 对象ID
* @return 同步结果
*/
public Map<String, Object> syncSingleObjectData(Integer id);
}

View File

@ -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);
}
}
/**
* 重试操作
*

View File

@ -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<String, Object> 同步结果包含successmessageobjectIdobjectApitotalCountdurationsyncTypelastFullSyncDate等字段
* @throws ConnectionException 当获取Salesforce连接失败时抛出
* @throws RuntimeException 当同步过程中发生其他异常时抛出
*/
@Override
public Map<String, Object> syncSingleObjectData(Integer id)
{
Map<String, Object> 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<String, Object> 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<String, Object> syncFullData(DataiIntegrationObject object, long startTime) {
Map<String, Object> result = new HashMap<>();
String objectApi = object.getApi().trim();
try {
DataiIntegrationBatch queryBatch = new DataiIntegrationBatch();
queryBatch.setApi(objectApi);
queryBatch.setSyncType("FULL");
List<DataiIntegrationBatch> 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<java.util.concurrent.Future<?>> futures = new ArrayList<>();
List<Map<String, Object>> 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<String, Object> 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<String, Object> 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<String, Object> 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<String, Object> syncIncrementalData(DataiIntegrationObject object, long startTime) {
Map<String, Object> 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<String> 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<String, Object> 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<String> getSalesforceObjectFields(PartnerConnection connection, String objectApi) throws ConnectionException
{
List<String> 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<String> 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<String> 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<String> 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<String, List<Map<String, Object>>> partitionedData = new HashMap<>();
// 创建普通数据列表用于存储非分区表的数据
List<Map<String, Object>> normalData = new ArrayList<>();
// 遍历所有记录进行数据转换和分区分配
for (SObject record : records) {
// 跳过空记录避免空指针异常
if (record == null) {
continue;
}
try {
// 将SObject对象转换为Map格式便于后续处理
Map<String, Object> 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<String, List<Map<String, Object>>> entry : partitionedData.entrySet()) {
String partitionName = entry.getKey();
List<Map<String, Object>> 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<String, Object> convertSObjectToMap(SObject record, List<String> fieldList)
{
// 创建Map对象用于存储转换后的记录数据
Map<String, Object> 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<Map<String, Object>> dataList)
{
// 检查数据列表是否为空为空则直接返回
if (dataList == null || dataList.isEmpty()) {
return;
}
try {
// 创建字段名列表用于存储表的所有列名
List<String> keys = new ArrayList<>();
// 创建字段值列表用于存储所有记录的字段值
List<Collection<Object>> values = new ArrayList<>();
// 遍历所有数据记录
for (Map<String, Object> 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<Map<String, Object>> dataList)
{
// 检查数据列表是否为空为空则直接返回
if (dataList == null || dataList.isEmpty()) {
return;
}
try {
// 创建字段名列表用于存储表的所有列名
List<String> keys = new ArrayList<>();
// 创建字段值列表用于存储所有记录的字段值
List<Collection<Object>> values = new ArrayList<>();
// 遍历所有数据记录
for (Map<String, Object> 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;
}
}
}

View File

@ -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

View File

@ -106,6 +106,11 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
<include refid="selectDataiIntegrationObjectVo"/>
where dio.id = #{id}
</select>
<select id="selectDataiIntegrationObjectByApi" parameterType="String" resultMap="DataiIntegrationObjectResult">
<include refid="selectDataiIntegrationObjectVo"/>
where dio.api = #{api}
</select>
<insert id="insertDataiIntegrationObject" parameterType="DataiIntegrationObject" useGeneratedKeys="true" keyProperty="id">
insert into datai_integration_object

View File

@ -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支持优先级排序和安全上下文传递
* <p>
* 该类扩展了FutureTask实现了Comparable接口可以根据优先级对任务进行排序
* 优先级规则index值越大优先级越高index相等时batch值越小优先级越高
* 同时支持将父线程的Spring Security上下文传递到子线程中
* </p>
*/
@Getter
@ -34,6 +37,11 @@ public class ComparableFutureTask extends FutureTask<Object> implements Comparab
*/
private final Integer batch;
/**
* 父线程的Spring Security上下文
*/
private final SecurityContext parentSecurityContext;
/**
* 构造函数
*
@ -45,6 +53,21 @@ public class ComparableFutureTask extends FutureTask<Object> 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);
}
}
/**

View File

@ -17,7 +17,13 @@ import java.util.concurrent.*;
* 该类提供了Salesforce数据处理任务的线程池管理功能支持任务优先级排序和批量处理
* 使用PriorityBlockingQueue作为任务队列确保高优先级任务优先执行
* </p>
*
* <p>
* 特性
* 1. 支持任务优先级排序基于ComparableFutureTask
* 2. 自动将父线程的Spring Security上下文传递到子线程
* 3. 可配置的线程池参数核心线程数最大线程数队列容量等
* 4. 支持优雅关闭和强制关闭
* </p>
*/
@Service
@Slf4j
@ -180,6 +186,10 @@ public class SalesforceExecutor {
* <p>
* 将任务提交到线程池执行支持优先级设置
* </p>
* <p>
* 注意该方法会自动将当前线程的Spring Security上下文传递到子线程中
* 确保子线程可以访问父线程的安全信息如登录用户信息
* </p>
*
* @param runnable 要执行的任务
* @param batch 批次号优先级相同时值越小优先级越高