diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/pom.xml b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/pom.xml
index fa9f4861..e86e90b6 100644
--- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/pom.xml
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/pom.xml
@@ -31,10 +31,21 @@
${project.basedir}/src/main/resources/lib/partner.jar
+
+ com.salesforce.eventbus
+ salesforce-pubsub-api
+ 1.0.0
+
+
org.springframework.boot
spring-boot-starter-aop
+
+
+ org.springframework.boot
+ spring-boot-starter-webflux
+
\ No newline at end of file
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationBatchController.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationBatchController.java
index 5ec76cb4..4075fc6c 100644
--- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationBatchController.java
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationBatchController.java
@@ -164,4 +164,20 @@ public class DataiIntegrationBatchController extends BaseController
return error((String) result.get("message"));
}
}
+
+ /**
+ * 获取所有批次统计信息
+ */
+ @Operation(summary = "获取所有批次统计信息")
+ @PreAuthorize("@ss.hasPermi('integration:batch:statistics')")
+ @GetMapping("/statistics")
+ public AjaxResult getAllBatchStatistics()
+ {
+ Map result = dataiIntegrationBatchService.getAllBatchStatistics();
+ if ((Boolean) result.get("success")) {
+ return success(result);
+ } else {
+ return error((String) result.get("message"));
+ }
+ }
}
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 f6db3177..9f56f0cc 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
@@ -202,6 +202,23 @@ public class DataiIntegrationObjectController extends BaseController
}
}
+ /**
+ * 变更对象实时同步状态
+ */
+ @Operation(summary = "变更对象实时同步状态")
+ @PreAuthorize("@ss.hasPermi('integration:object:updateRealtimeSyncStatus')")
+ @Log(title = "对象同步控制", businessType = BusinessType.UPDATE)
+ @PutMapping("/{id}/realtimeSyncStatus")
+ public AjaxResult updateRealtimeSyncStatus(@PathVariable("id") Integer id, @org.springframework.web.bind.annotation.RequestParam("isRealtimeSync") Boolean isRealtimeSync)
+ {
+ Map result = dataiIntegrationObjectService.updateRealtimeSyncStatus(id, isRealtimeSync);
+ if ((Boolean) result.get("success")) {
+ return success(result);
+ } else {
+ return error((String) result.get("message"));
+ }
+ }
+
/**
* 获取对象整体统计信息
*/
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationRealtimeSyncController.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationRealtimeSyncController.java
new file mode 100644
index 00000000..c6a3e224
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationRealtimeSyncController.java
@@ -0,0 +1,172 @@
+package com.datai.integration.controller;
+
+import com.datai.integration.realtime.RealtimeSyncService;
+import com.datai.integration.realtime.impl.ObjectRegistryImpl;
+import com.datai.integration.model.domain.DataiIntegrationObject;
+import io.swagger.v3.oas.annotations.Operation;
+import io.swagger.v3.oas.annotations.tags.Tag;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.security.access.prepost.PreAuthorize;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.List;
+import java.util.Map;
+import java.util.HashMap;
+
+/**
+ * 实时同步控制器
+ * 用于管理实时同步服务
+ */
+@RestController
+@RequestMapping("/integration/realtime")
+@Tag(name = "【实时同步管理】管理")
+@Slf4j
+public class DataiIntegrationRealtimeSyncController {
+
+ @Autowired
+ private RealtimeSyncService realtimeSyncService;
+
+ @Autowired
+ private ObjectRegistryImpl objectRegistry;
+
+ /**
+ * 获取实时同步服务状态
+ */
+ @Operation(summary = "获取实时同步服务状态")
+ @PreAuthorize("@ss.hasPermi('integration:realtime:status')")
+ @GetMapping("/status")
+ public Map getStatus() {
+ log.info("获取实时同步服务状态");
+
+ Map result = new HashMap<>();
+
+ try {
+ // 获取启用实时同步的对象列表
+ List realtimeSyncObjects = objectRegistry.getRealtimeSyncObjects();
+
+ result.put("success", true);
+ result.put("message", "获取实时同步服务状态成功");
+ result.put("realtimeSyncObjects", realtimeSyncObjects);
+ result.put("objectCount", realtimeSyncObjects.size());
+
+ log.info("获取实时同步服务状态成功,共 {} 个启用实时同步的对象", realtimeSyncObjects.size());
+ } catch (Exception e) {
+ log.error("获取实时同步服务状态时发生异常: {}", e.getMessage(), e);
+ result.put("success", false);
+ result.put("message", "获取实时同步服务状态失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+
+ /**
+ * 启动实时同步服务
+ */
+ @Operation(summary = "启动实时同步服务")
+ @PreAuthorize("@ss.hasPermi('integration:realtime:start')")
+ @PostMapping("/start")
+ public Map start() {
+ log.info("启动实时同步服务");
+
+ Map result = new HashMap<>();
+
+ try {
+ realtimeSyncService.start();
+ result.put("success", true);
+ result.put("message", "实时同步服务启动成功");
+ log.info("实时同步服务启动成功");
+ } catch (Exception e) {
+ log.error("启动实时同步服务时发生异常: {}", e.getMessage(), e);
+ result.put("success", false);
+ result.put("message", "实时同步服务启动失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+
+ /**
+ * 停止实时同步服务
+ */
+ @Operation(summary = "停止实时同步服务")
+ @PreAuthorize("@ss.hasPermi('integration:realtime:stop')")
+ @PostMapping("/stop")
+ public Map stop() {
+ log.info("停止实时同步服务");
+
+ Map result = new HashMap<>();
+
+ try {
+ realtimeSyncService.stop();
+ result.put("success", true);
+ result.put("message", "实时同步服务停止成功");
+ log.info("实时同步服务停止成功");
+ } catch (Exception e) {
+ log.error("停止实时同步服务时发生异常: {}", e.getMessage(), e);
+ result.put("success", false);
+ result.put("message", "实时同步服务停止失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+
+ /**
+ * 重启实时同步服务
+ */
+ @Operation(summary = "重启实时同步服务")
+ @PreAuthorize("@ss.hasPermi('integration:realtime:restart')")
+ @PostMapping("/restart")
+ public Map restart() {
+ log.info("重启实时同步服务");
+
+ Map result = new HashMap<>();
+
+ try {
+ realtimeSyncService.restart();
+ result.put("success", true);
+ result.put("message", "实时同步服务重启成功");
+ log.info("实时同步服务重启成功");
+ } catch (Exception e) {
+ log.error("重启实时同步服务时发生异常: {}", e.getMessage(), e);
+ result.put("success", false);
+ result.put("message", "实时同步服务重启失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+
+ /**
+ * 刷新对象注册表
+ */
+ @Operation(summary = "刷新对象注册表")
+ @PreAuthorize("@ss.hasPermi('integration:realtime:refresh')")
+ @PostMapping("/refresh")
+ public Map refresh() {
+ log.info("刷新对象注册表");
+
+ Map result = new HashMap<>();
+
+ try {
+ realtimeSyncService.refreshObjectRegistry();
+
+ // 获取刷新后的对象列表
+ List realtimeSyncObjects = objectRegistry.getRealtimeSyncObjects();
+
+ result.put("success", true);
+ result.put("message", "对象注册表刷新成功");
+ result.put("realtimeSyncObjects", realtimeSyncObjects);
+ result.put("objectCount", realtimeSyncObjects.size());
+
+ log.info("对象注册表刷新成功,共 {} 个启用实时同步的对象", realtimeSyncObjects.size());
+ } catch (Exception e) {
+ log.error("刷新对象注册表时发生异常: {}", e.getMessage(), e);
+ result.put("success", false);
+ result.put("message", "对象注册表刷新失败: " + e.getMessage());
+ }
+
+ return result;
+ }
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationRealtimeSyncLogController.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationRealtimeSyncLogController.java
new file mode 100644
index 00000000..65459657
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/controller/DataiIntegrationRealtimeSyncLogController.java
@@ -0,0 +1,121 @@
+package com.datai.integration.controller;
+
+import java.util.List;
+import java.util.stream.Collectors;
+
+import com.datai.common.utils.PageUtils;
+import com.datai.integration.model.domain.DataiIntegrationRealtimeSyncLog;
+import com.datai.integration.model.dto.DataiIntegrationRealtimeSyncLogDto;
+import com.datai.integration.model.vo.DataiIntegrationRealtimeSyncLogVo;
+import jakarta.servlet.http.HttpServletResponse;
+import org.springframework.security.access.prepost.PreAuthorize;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.PutMapping;
+import org.springframework.web.bind.annotation.DeleteMapping;
+import org.springframework.web.bind.annotation.PathVariable;
+import org.springframework.web.bind.annotation.RequestBody;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+import com.datai.common.annotation.Log;
+import com.datai.common.core.controller.BaseController;
+import com.datai.common.core.domain.AjaxResult;
+import com.datai.common.enums.BusinessType;
+
+import com.datai.integration.service.IDataiIntegrationRealtimeSyncLogService;
+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;
+
+/**
+ * 实时同步日志Controller
+ *
+ * @author datai
+ * @date 2026-01-09
+ */
+@RestController
+@RequestMapping("/integration/realtimelog")
+@Tag(name = "【实时同步日志】管理")
+public class DataiIntegrationRealtimeSyncLogController extends BaseController
+{
+ @Autowired
+ private IDataiIntegrationRealtimeSyncLogService dataiIntegrationRealtimeSyncLogService;
+
+ /**
+ * 查询实时同步日志列表
+ */
+ @Operation(summary = "查询实时同步日志列表")
+ @PreAuthorize("@ss.hasPermi('integration:realtimelog:list')")
+ @GetMapping("/list")
+ public TableDataInfo list(DataiIntegrationRealtimeSyncLogDto dataiIntegrationRealtimeSyncLogDto)
+ {
+ startPage();
+ List list = dataiIntegrationRealtimeSyncLogService.selectDataiIntegrationRealtimeSyncLogList(DataiIntegrationRealtimeSyncLogDto.toObj(dataiIntegrationRealtimeSyncLogDto));
+ List voList = list.stream().map(DataiIntegrationRealtimeSyncLogVo::objToVo).collect(Collectors.toList());
+ return getDataTableByPage(voList, PageUtils.getTotal(list));
+ }
+
+ /**
+ * 导出实时同步日志列表
+ */
+ @Operation(summary = "导出实时同步日志列表")
+ @PreAuthorize("@ss.hasPermi('integration:realtimelog:export')")
+ @Log(title = "实时同步日志", businessType = BusinessType.EXPORT)
+ @PostMapping("/export")
+ public void export(HttpServletResponse response, DataiIntegrationRealtimeSyncLogDto dataiIntegrationRealtimeSyncLogDto)
+ {
+ List list = dataiIntegrationRealtimeSyncLogService.selectDataiIntegrationRealtimeSyncLogList(DataiIntegrationRealtimeSyncLogDto.toObj(dataiIntegrationRealtimeSyncLogDto));
+ ExcelUtil util = new ExcelUtil(DataiIntegrationRealtimeSyncLog.class);
+ util.exportExcel(response, list, "实时同步日志数据");
+ }
+
+ /**
+ * 获取实时同步日志详细信息
+ */
+ @Operation(summary = "获取实时同步日志详细信息")
+ @PreAuthorize("@ss.hasPermi('integration:realtimelog:query')")
+ @GetMapping(value = "/{id}")
+ public AjaxResult getInfo(@PathVariable("id") Long id)
+ {
+ DataiIntegrationRealtimeSyncLog dataiIntegrationRealtimeSyncLog = dataiIntegrationRealtimeSyncLogService.selectDataiIntegrationRealtimeSyncLogById(id);
+ return success(DataiIntegrationRealtimeSyncLogVo.objToVo(dataiIntegrationRealtimeSyncLog));
+ }
+
+ /**
+ * 新增实时同步日志
+ */
+ @Operation(summary = "新增实时同步日志")
+ @PreAuthorize("@ss.hasPermi('integration:realtimelog:add')")
+ @Log(title = "实时同步日志", businessType = BusinessType.INSERT)
+ @PostMapping
+ public AjaxResult add(@RequestBody DataiIntegrationRealtimeSyncLogDto dataiIntegrationRealtimeSyncLogDto)
+ {
+ return toAjax(dataiIntegrationRealtimeSyncLogService.insertDataiIntegrationRealtimeSyncLog(DataiIntegrationRealtimeSyncLogDto.toObj(dataiIntegrationRealtimeSyncLogDto)));
+ }
+
+ /**
+ * 修改实时同步日志
+ */
+ @Operation(summary = "修改实时同步日志")
+ @PreAuthorize("@ss.hasPermi('integration:realtimelog:edit')")
+ @Log(title = "实时同步日志", businessType = BusinessType.UPDATE)
+ @PutMapping
+ public AjaxResult edit(@RequestBody DataiIntegrationRealtimeSyncLogDto dataiIntegrationRealtimeSyncLogDto)
+ {
+ return toAjax(dataiIntegrationRealtimeSyncLogService.updateDataiIntegrationRealtimeSyncLog(DataiIntegrationRealtimeSyncLogDto.toObj(dataiIntegrationRealtimeSyncLogDto)));
+ }
+
+ /**
+ * 删除实时同步日志
+ */
+ @Operation(summary = "删除实时同步日志")
+ @PreAuthorize("@ss.hasPermi('integration:realtimelog:remove')")
+ @Log(title = "实时同步日志", businessType = BusinessType.DELETE)
+ @DeleteMapping("/{ids}")
+ public AjaxResult remove(@PathVariable( name = "ids" ) Long[] ids)
+ {
+ return toAjax(dataiIntegrationRealtimeSyncLogService.deleteDataiIntegrationRealtimeSyncLogByIds(ids));
+ }
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/core/PubSubClient.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/core/PubSubClient.java
new file mode 100644
index 00000000..b6814454
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/core/PubSubClient.java
@@ -0,0 +1,162 @@
+package com.datai.integration.core;
+
+import com.salesforce.eventbus.EventBusClient;
+import com.salesforce.eventbus.EventBusClientFactory;
+import com.salesforce.eventbus.Subscription;
+import com.salesforce.eventbus.SubscriptionListener;
+import com.salesforce.eventbus.protobuf.EventBatch;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.stereotype.Component;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Salesforce Pub/Sub API 客户端
+ * 用于管理与 Salesforce Pub/Sub API 的连接和订阅
+ */
+@Slf4j
+@Component
+public class PubSubClient {
+
+ private EventBusClient eventBusClient;
+ private Subscription subscription;
+ private boolean connected = false;
+
+ @Value("${salesforce.pubsub.endpoint:https://api.salesforce.com/eventbus/v1}")
+ private String pubSubEndpoint;
+
+ @Value("${salesforce.pubsub.replayId:-1}")
+ private long replayId;
+
+ @Value("${salesforce.pubsub.timeout:30}")
+ private int timeout;
+
+ /**
+ * 初始化 Pub/Sub API 客户端
+ * @param accessToken Salesforce 访问令牌
+ * @param instanceUrl Salesforce 实例 URL
+ * @return 是否初始化成功
+ */
+ public boolean initialize(String accessToken, String instanceUrl) {
+ try {
+ log.info("开始初始化 Salesforce Pub/Sub API 客户端");
+
+ // 创建 EventBusClientFactory
+ EventBusClientFactory factory = EventBusClientFactory.builder()
+ .withAuthProvider(() -> accessToken)
+ .withEndpoint(pubSubEndpoint)
+ .withInstanceUrl(instanceUrl)
+ .build();
+
+ // 创建 EventBusClient
+ eventBusClient = factory.createClient();
+
+ // 连接到 Salesforce Event Bus
+ CountDownLatch latch = new CountDownLatch(1);
+ eventBusClient.connect(connection -> {
+ if (connection.isSuccess()) {
+ log.info("成功连接到 Salesforce Event Bus");
+ connected = true;
+ } else {
+ log.error("连接到 Salesforce Event Bus 失败: {}", connection.getError());
+ connected = false;
+ }
+ latch.countDown();
+ });
+
+ // 等待连接完成
+ if (!latch.await(timeout, TimeUnit.SECONDS)) {
+ log.error("连接到 Salesforce Event Bus 超时");
+ return false;
+ }
+
+ return connected;
+ } catch (Exception e) {
+ log.error("初始化 Salesforce Pub/Sub API 客户端失败: {}", e.getMessage(), e);
+ return false;
+ }
+ }
+
+ /**
+ * 订阅事件通道
+ * @param topic 事件通道名称
+ * @param listener 事件监听器
+ * @return 是否订阅成功
+ */
+ public boolean subscribe(String topic, SubscriptionListener listener) {
+ try {
+ if (!connected || eventBusClient == null) {
+ log.error("Pub/Sub API 客户端未连接,无法订阅事件通道");
+ return false;
+ }
+
+ log.info("开始订阅事件通道: {}", topic);
+
+ // 创建订阅
+ subscription = eventBusClient.subscribe(topic, replayId, listener);
+
+ log.info("成功订阅事件通道: {}", topic);
+ return true;
+ } catch (Exception e) {
+ log.error("订阅事件通道 {} 失败: {}", topic, e.getMessage(), e);
+ return false;
+ }
+ }
+
+ /**
+ * 取消订阅
+ */
+ public void unsubscribe() {
+ try {
+ if (subscription != null) {
+ log.info("开始取消订阅事件通道");
+ subscription.cancel();
+ log.info("成功取消订阅事件通道");
+ }
+ } catch (Exception e) {
+ log.error("取消订阅事件通道失败: {}", e.getMessage(), e);
+ }
+ }
+
+ /**
+ * 断开连接
+ */
+ public void disconnect() {
+ try {
+ if (eventBusClient != null) {
+ log.info("开始断开与 Salesforce Event Bus 的连接");
+ eventBusClient.disconnect();
+ connected = false;
+ log.info("成功断开与 Salesforce Event Bus 的连接");
+ }
+ } catch (Exception e) {
+ log.error("断开与 Salesforce Event Bus 的连接失败: {}", e.getMessage(), e);
+ }
+ }
+
+ /**
+ * 检查客户端是否已连接
+ * @return 是否已连接
+ */
+ public boolean isConnected() {
+ return connected;
+ }
+
+ /**
+ * 获取 EventBusClient
+ * @return EventBusClient
+ */
+ public EventBusClient getEventBusClient() {
+ return eventBusClient;
+ }
+
+ /**
+ * 获取当前订阅
+ * @return Subscription
+ */
+ public Subscription getSubscription() {
+ return subscription;
+ }
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/factory/impl/PubSubConnectionFactory.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/factory/impl/PubSubConnectionFactory.java
new file mode 100644
index 00000000..09842128
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/factory/impl/PubSubConnectionFactory.java
@@ -0,0 +1,76 @@
+package com.datai.integration.factory.impl;
+
+import com.datai.integration.core.PubSubClient;
+import com.datai.integration.factory.AbstractConnectionFactory;
+import com.datai.salesforce.auth.service.ISalesforceAuthService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+/**
+ * Salesforce Pub/Sub API 连接工厂
+ * 用于创建和管理 PubSubClient 实例
+ */
+@Slf4j
+@Component
+public class PubSubConnectionFactory extends AbstractConnectionFactory {
+
+ @Autowired
+ private ISalesforceAuthService salesforceAuthService;
+
+ @Autowired
+ private PubSubClient pubSubClient;
+
+ /**
+ * 获取 Salesforce Pub/Sub API 客户端
+ * @return PubSubClient 实例
+ */
+ @Override
+ public PubSubClient getConnection() {
+ try {
+ // 检查是否已连接
+ if (pubSubClient.isConnected()) {
+ log.info("Pub/Sub API 客户端已连接,直接返回实例");
+ return pubSubClient;
+ }
+
+ // 获取 Salesforce 访问令牌和实例 URL
+ String accessToken = salesforceAuthService.getAccessToken();
+ String instanceUrl = salesforceAuthService.getInstanceUrl();
+
+ if (accessToken == null || instanceUrl == null) {
+ log.error("无法获取 Salesforce 访问令牌或实例 URL");
+ return null;
+ }
+
+ // 初始化 Pub/Sub API 客户端
+ boolean initialized = pubSubClient.initialize(accessToken, instanceUrl);
+ if (initialized) {
+ log.info("成功获取 Pub/Sub API 客户端实例");
+ return pubSubClient;
+ } else {
+ log.error("初始化 Pub/Sub API 客户端失败");
+ return null;
+ }
+ } catch (Exception e) {
+ log.error("获取 Pub/Sub API 客户端实例失败: {}", e.getMessage(), e);
+ return null;
+ }
+ }
+
+ /**
+ * 关闭 Salesforce Pub/Sub API 连接
+ * @param connection PubSubClient 实例
+ */
+ @Override
+ public void closeConnection(PubSubClient connection) {
+ try {
+ if (connection != null) {
+ connection.disconnect();
+ log.info("成功关闭 Pub/Sub API 连接");
+ }
+ } catch (Exception e) {
+ log.error("关闭 Pub/Sub API 连接失败: {}", e.getMessage(), e);
+ }
+ }
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationRealtimeSyncLogMapper.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationRealtimeSyncLogMapper.java
new file mode 100644
index 00000000..48fcccbc
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/mapper/DataiIntegrationRealtimeSyncLogMapper.java
@@ -0,0 +1,62 @@
+package com.datai.integration.mapper;
+
+import com.datai.integration.model.domain.DataiIntegrationRealtimeSyncLog;
+
+import java.util.List;
+
+/**
+ * 实时同步日志Mapper接口
+ *
+ * @author datai
+ * @date 2026-01-09
+ */
+public interface DataiIntegrationRealtimeSyncLogMapper
+{
+ /**
+ * 查询实时同步日志
+ *
+ * @param id 实时同步日志主键
+ * @return 实时同步日志
+ */
+ public DataiIntegrationRealtimeSyncLog selectDataiIntegrationRealtimeSyncLogById(Long id);
+
+ /**
+ * 查询实时同步日志列表
+ *
+ * @param dataiIntegrationRealtimeSyncLog 实时同步日志
+ * @return 实时同步日志集合
+ */
+ public List selectDataiIntegrationRealtimeSyncLogList(DataiIntegrationRealtimeSyncLog dataiIntegrationRealtimeSyncLog);
+
+ /**
+ * 新增实时同步日志
+ *
+ * @param dataiIntegrationRealtimeSyncLog 实时同步日志
+ * @return 结果
+ */
+ public int insertDataiIntegrationRealtimeSyncLog(DataiIntegrationRealtimeSyncLog dataiIntegrationRealtimeSyncLog);
+
+ /**
+ * 修改实时同步日志
+ *
+ * @param dataiIntegrationRealtimeSyncLog 实时同步日志
+ * @return 结果
+ */
+ public int updateDataiIntegrationRealtimeSyncLog(DataiIntegrationRealtimeSyncLog dataiIntegrationRealtimeSyncLog);
+
+ /**
+ * 删除实时同步日志
+ *
+ * @param id 实时同步日志主键
+ * @return 结果
+ */
+ public int deleteDataiIntegrationRealtimeSyncLogById(Long id);
+
+ /**
+ * 批量删除实时同步日志
+ *
+ * @param ids 需要删除的数据主键集合
+ * @return 结果
+ */
+ public int deleteDataiIntegrationRealtimeSyncLogByIds(Long[] ids);
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationObject.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationObject.java
index 5e80ccba..db392b00 100644
--- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationObject.java
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationObject.java
@@ -147,6 +147,11 @@ public class DataiIntegrationObject extends BaseEntity
@Schema(title = "失败原因")
@Excel(name = "失败原因")
private String errorMessage;
+
+ /** 实时同步 */
+ @Schema(title = "实时同步")
+ @Excel(name = "实时同步")
+ private Boolean isRealtimeSync;
public void setId(Integer id)
{
this.id = id;
@@ -432,6 +437,16 @@ public class DataiIntegrationObject extends BaseEntity
return errorMessage;
}
+ public void setIsRealtimeSync(Boolean isRealtimeSync)
+ {
+ this.isRealtimeSync = isRealtimeSync;
+ }
+
+ public Boolean getIsRealtimeSync()
+ {
+ return isRealtimeSync;
+ }
+
@Override
@@ -463,6 +478,7 @@ public class DataiIntegrationObject extends BaseEntity
.append("lastBatchDate", getLastBatchDate())
.append("syncStatus", getSyncStatus())
.append("errorMessage", getErrorMessage())
+ .append("isRealtimeSync", getIsRealtimeSync())
.append("remark", getRemark())
.append("createBy", getCreateBy())
.append("createTime", getCreateTime())
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationRealtimeSyncLog.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationRealtimeSyncLog.java
new file mode 100644
index 00000000..b1cc4e5e
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/domain/DataiIntegrationRealtimeSyncLog.java
@@ -0,0 +1,199 @@
+package com.datai.integration.model.domain;
+
+import java.time.LocalDateTime;
+import com.fasterxml.jackson.annotation.JsonFormat;
+import io.swagger.v3.oas.annotations.media.Schema;
+import org.apache.commons.lang3.builder.ToStringBuilder;
+import org.apache.commons.lang3.builder.ToStringStyle;
+import com.datai.common.annotation.Excel;
+import com.datai.common.core.domain.BaseEntity;
+
+/**
+ * 实时同步日志对象 datai_integration_realtime_sync_log
+ *
+ * @author datai
+ * @date 2026-01-09
+ */
+@Schema(description = "实时同步日志对象")
+public class DataiIntegrationRealtimeSyncLog extends BaseEntity
+{
+ private static final long serialVersionUID = 1L;
+
+
+ /** 主键ID */
+ @Schema(title = "主键ID")
+ private Long id;
+
+ /** 对象名称 */
+ @Schema(title = "对象名称")
+ @Excel(name = "对象名称")
+ private String objectName;
+
+ /** 记录ID */
+ @Schema(title = "记录ID")
+ @Excel(name = "记录ID")
+ private String recordId;
+
+ /** 操作类型 */
+ @Schema(title = "操作类型 ")
+ @Excel(name = "操作类型 ")
+ private String operationType;
+
+ /** 变更数据 */
+ @Schema(title = "变更数据")
+ @Excel(name = "变更数据")
+ private String changeData;
+
+ /** 同步状态 */
+ @Schema(title = "同步状态")
+ @Excel(name = "同步状态")
+ private String syncStatus;
+
+ /** 错误信息 */
+ @Schema(title = "错误信息")
+ @Excel(name = "错误信息")
+ private String errorMessage;
+
+ /** 重试次数 */
+ @Schema(title = "重试次数")
+ @Excel(name = "重试次数")
+ private Integer retryCount;
+
+ /** Salesforce时间戳 */
+ @Schema(title = "Salesforce时间戳")
+ @Excel(name = "Salesforce时间戳")
+ private LocalDateTime salesforceTimestamp;
+
+ /** 同步时间戳 */
+ @Schema(title = "同步时间戳")
+ @Excel(name = "同步时间戳")
+ private LocalDateTime syncTimestamp;
+ public void setId(Long id)
+ {
+ this.id = id;
+ }
+
+ public Long getId()
+ {
+ return id;
+ }
+
+
+ public void setObjectName(String objectName)
+ {
+ this.objectName = objectName;
+ }
+
+ public String getObjectName()
+ {
+ return objectName;
+ }
+
+
+ public void setRecordId(String recordId)
+ {
+ this.recordId = recordId;
+ }
+
+ public String getRecordId()
+ {
+ return recordId;
+ }
+
+
+ public void setOperationType(String operationType)
+ {
+ this.operationType = operationType;
+ }
+
+ public String getOperationType()
+ {
+ return operationType;
+ }
+
+
+ public void setChangeData(String changeData)
+ {
+ this.changeData = changeData;
+ }
+
+ public String getChangeData()
+ {
+ return changeData;
+ }
+
+
+ public void setSyncStatus(String syncStatus)
+ {
+ this.syncStatus = syncStatus;
+ }
+
+ public String getSyncStatus()
+ {
+ return syncStatus;
+ }
+
+
+ public void setErrorMessage(String errorMessage)
+ {
+ this.errorMessage = errorMessage;
+ }
+
+ public String getErrorMessage()
+ {
+ return errorMessage;
+ }
+
+
+ public void setRetryCount(Integer retryCount)
+ {
+ this.retryCount = retryCount;
+ }
+
+ public Integer getRetryCount()
+ {
+ return retryCount;
+ }
+
+
+ public void setSalesforceTimestamp(LocalDateTime salesforceTimestamp)
+ {
+ this.salesforceTimestamp = salesforceTimestamp;
+ }
+
+ public LocalDateTime getSalesforceTimestamp()
+ {
+ return salesforceTimestamp;
+ }
+
+
+ public void setSyncTimestamp(LocalDateTime syncTimestamp)
+ {
+ this.syncTimestamp = syncTimestamp;
+ }
+
+ public LocalDateTime getSyncTimestamp()
+ {
+ return syncTimestamp;
+ }
+
+
+
+ @Override
+ public String toString() {
+ return new ToStringBuilder(this,ToStringStyle.MULTI_LINE_STYLE)
+ .append("id", getId())
+ .append("objectName", getObjectName())
+ .append("recordId", getRecordId())
+ .append("operationType", getOperationType())
+ .append("changeData", getChangeData())
+ .append("syncStatus", getSyncStatus())
+ .append("errorMessage", getErrorMessage())
+ .append("retryCount", getRetryCount())
+ .append("salesforceTimestamp", getSalesforceTimestamp())
+ .append("syncTimestamp", getSyncTimestamp())
+ .append("createTime", getCreateTime())
+ .append("updateTime", getUpdateTime())
+ .toString();
+ }
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationObjectDto.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationObjectDto.java
index b8465fec..9472b0e5 100644
--- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationObjectDto.java
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationObjectDto.java
@@ -102,6 +102,9 @@ public class DataiIntegrationObjectDto implements Serializable
/** 失败原因 */
private String errorMessage;
+ /** 实时同步 */
+ private Boolean isRealtimeSync;
+
/** 备注 */
private String remark;
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationRealtimeSyncLogDto.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationRealtimeSyncLogDto.java
new file mode 100644
index 00000000..cc05c2e8
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/dto/DataiIntegrationRealtimeSyncLogDto.java
@@ -0,0 +1,90 @@
+package com.datai.integration.model.dto;
+
+import java.io.Serializable;
+import java.util.Map;
+import java.util.Date;
+import java.util.List;
+ import java.time.LocalDateTime;
+ import com.fasterxml.jackson.annotation.JsonFormat;
+import lombok.Data;
+import org.springframework.beans.BeanUtils;
+import com.fasterxml.jackson.annotation.JsonFormat;
+import com.fasterxml.jackson.annotation.JsonInclude;
+import com.datai.integration.model.domain.DataiIntegrationRealtimeSyncLog;
+
+/**
+ * 实时同步日志通用业务传输对象 (Dto)
+ * 整合了查询、新增、修改的所有字段
+ *
+ * @author datai
+ * @date 2026-01-09
+ */
+@Data
+public class DataiIntegrationRealtimeSyncLogDto implements Serializable
+{
+ private static final long serialVersionUID = 1L;
+
+ /** 主键ID */
+ private Long id;
+
+ /** 对象名称 */
+ private String objectName;
+
+ /** 记录ID */
+ private String recordId;
+
+ /** 操作类型 */
+ private String operationType;
+
+ /** 变更数据 */
+ private String changeData;
+
+ /** 同步状态 */
+ private String syncStatus;
+
+ /** 错误信息 */
+ private String errorMessage;
+
+ /** 重试次数 */
+ private Integer retryCount;
+
+ /** Salesforce时间戳 */
+ private LocalDateTime salesforceTimestamp;
+
+ /** 同步时间戳 */
+ private LocalDateTime syncTimestamp;
+
+ /** 创建时间 */
+ private LocalDateTime createTime;
+
+ /** 更新时间 */
+ private LocalDateTime updateTime;
+
+ /** 请求参数(用于存放查询范围等临时数据) */
+ @JsonInclude(JsonInclude.Include.NON_EMPTY)
+ private Map params;
+
+ /**
+ * Dto 转 业务对象 (DataiIntegrationRealtimeSyncLog)
+ */
+ public static DataiIntegrationRealtimeSyncLog toObj(DataiIntegrationRealtimeSyncLogDto Dto) {
+ if (Dto == null) {
+ return null;
+ }
+ DataiIntegrationRealtimeSyncLog obj = new DataiIntegrationRealtimeSyncLog();
+ BeanUtils.copyProperties(Dto, obj);
+ return obj;
+ }
+
+ /**
+ * 业务对象 (DataiIntegrationRealtimeSyncLog) 转 Dto
+ */
+ public static DataiIntegrationRealtimeSyncLogDto fromObj(DataiIntegrationRealtimeSyncLog obj) {
+ if (obj == null) {
+ return null;
+ }
+ DataiIntegrationRealtimeSyncLogDto Dto = new DataiIntegrationRealtimeSyncLogDto();
+ BeanUtils.copyProperties(obj, Dto);
+ return Dto;
+ }
+}
\ No newline at end of file
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationObjectVo.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationObjectVo.java
index 6e7cf0f4..fa10459b 100644
--- a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationObjectVo.java
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationObjectVo.java
@@ -97,6 +97,9 @@ public class DataiIntegrationObjectVo implements Serializable {
/** 失败原因 */
private String errorMessage;
+ /** 实时同步 */
+ private Boolean isRealtimeSync;
+
/** 备注 */
private String remark;
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationRealtimeSyncLogVo.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationRealtimeSyncLogVo.java
new file mode 100644
index 00000000..a922d273
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/model/vo/DataiIntegrationRealtimeSyncLogVo.java
@@ -0,0 +1,73 @@
+package com.datai.integration.model.vo;
+
+import java.io.Serializable;
+import java.util.Date;
+ import java.time.LocalDateTime;
+ import com.fasterxml.jackson.annotation.JsonFormat;
+import lombok.Data;
+import com.datai.common.annotation.Excel;
+import org.springframework.beans.BeanUtils;
+import com.fasterxml.jackson.annotation.JsonFormat;
+import com.datai.integration.model.domain.DataiIntegrationRealtimeSyncLog;
+/**
+ * 实时同步日志Vo对象 datai_integration_realtime_sync_log
+ *
+ * @author datai
+ * @date 2026-01-09
+ */
+@Data
+public class DataiIntegrationRealtimeSyncLogVo implements Serializable {
+ private static final long serialVersionUID = 1L;
+
+ /** 主键ID */
+ private Long id;
+
+ /** 对象名称 */
+ private String objectName;
+
+ /** 记录ID */
+ private String recordId;
+
+ /** 操作类型 */
+ private String operationType;
+
+ /** 变更数据 */
+ private String changeData;
+
+ /** 同步状态 */
+ private String syncStatus;
+
+ /** 错误信息 */
+ private String errorMessage;
+
+ /** 重试次数 */
+ private Integer retryCount;
+
+ /** Salesforce时间戳 */
+ private LocalDateTime salesforceTimestamp;
+
+ /** 同步时间戳 */
+ private LocalDateTime syncTimestamp;
+
+ /** 创建时间 */
+ private LocalDateTime createTime;
+
+ /** 更新时间 */
+ private LocalDateTime updateTime;
+
+
+ /**
+ * 对象转封装类
+ *
+ * @param dataiIntegrationRealtimeSyncLog DataiIntegrationRealtimeSyncLog实体对象
+ * @return DataiIntegrationRealtimeSyncLogVo
+ */
+ public static DataiIntegrationRealtimeSyncLogVo objToVo(DataiIntegrationRealtimeSyncLog dataiIntegrationRealtimeSyncLog) {
+ if (dataiIntegrationRealtimeSyncLog == null) {
+ return null;
+ }
+ DataiIntegrationRealtimeSyncLogVo dataiIntegrationRealtimeSyncLogVo = new DataiIntegrationRealtimeSyncLogVo();
+ BeanUtils.copyProperties(dataiIntegrationRealtimeSyncLog, dataiIntegrationRealtimeSyncLogVo);
+ return dataiIntegrationRealtimeSyncLogVo;
+ }
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/DataSynchronizer.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/DataSynchronizer.java
new file mode 100644
index 00000000..8369575d
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/DataSynchronizer.java
@@ -0,0 +1,21 @@
+package com.datai.integration.realtime;
+
+import java.util.Date;
+import java.util.Map;
+
+/**
+ * 数据同步器接口
+ * 用于将Salesforce Change Events变更数据同步至本地数据库
+ */
+public interface DataSynchronizer {
+
+ /**
+ * 同步数据变更
+ * @param objectType 对象类型
+ * @param recordId 记录ID
+ * @param changeType 变更类型
+ * @param changeData 变更数据
+ * @param changeDate 变更时间
+ */
+ void synchronizeData(String objectType, String recordId, String changeType, Map changeData, Date changeDate);
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/EventProcessor.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/EventProcessor.java
new file mode 100644
index 00000000..a168441a
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/EventProcessor.java
@@ -0,0 +1,16 @@
+package com.datai.integration.realtime;
+
+import com.sforce.soap.partner.sobject.SObject;
+
+/**
+ * 事件处理器接口
+ * 用于处理捕获到的Salesforce Change Events
+ */
+public interface EventProcessor {
+
+ /**
+ * 处理捕获到的事件
+ * @param event Salesforce Change Event
+ */
+ void processEvent(SObject event);
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/EventSubscriber.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/EventSubscriber.java
new file mode 100644
index 00000000..673c23e2
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/EventSubscriber.java
@@ -0,0 +1,27 @@
+package com.datai.integration.realtime;
+
+import com.sforce.ws.ConnectionException;
+
+/**
+ * 事件订阅器接口
+ * 用于订阅Salesforce Change Events,实时捕获数据变更
+ */
+public interface EventSubscriber {
+
+ /**
+ * 启动事件订阅
+ * @throws ConnectionException 当获取Salesforce连接失败时抛出
+ */
+ void startSubscription() throws ConnectionException;
+
+ /**
+ * 停止事件订阅
+ */
+ void stopSubscription();
+
+ /**
+ * 检查订阅状态
+ * @return 是否正在订阅
+ */
+ boolean isSubscribed();
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/ObjectRegistry.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/ObjectRegistry.java
new file mode 100644
index 00000000..01bc45ac
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/ObjectRegistry.java
@@ -0,0 +1,41 @@
+package com.datai.integration.realtime;
+
+import com.datai.integration.model.domain.DataiIntegrationObject;
+import java.util.List;
+
+/**
+ * 对象注册表接口
+ * 用于管理所有启用实时同步的对象信息
+ */
+public interface ObjectRegistry {
+
+ /**
+ * 注册启用实时同步的对象
+ * @param object 对象信息
+ */
+ void registerObject(DataiIntegrationObject object);
+
+ /**
+ * 注销对象
+ * @param objectApi 对象API
+ */
+ void unregisterObject(String objectApi);
+
+ /**
+ * 获取所有启用实时同步的对象
+ * @return 启用实时同步的对象列表
+ */
+ List getRealtimeSyncObjects();
+
+ /**
+ * 检查对象是否已注册
+ * @param objectApi 对象API
+ * @return 是否已注册
+ */
+ boolean isObjectRegistered(String objectApi);
+
+ /**
+ * 刷新对象注册表
+ */
+ void refreshRegistry();
+}
diff --git a/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/RealtimeSyncService.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/RealtimeSyncService.java
new file mode 100644
index 00000000..8e108392
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/RealtimeSyncService.java
@@ -0,0 +1,99 @@
+package com.datai.integration.realtime;
+
+import com.datai.integration.realtime.impl.ObjectRegistryImpl;
+import com.datai.integration.realtime.impl.PubSubEventSubscriberImpl;
+import jakarta.annotation.PostConstruct;
+import jakarta.annotation.PreDestroy;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+/**
+ * 实时同步服务
+ * 用于初始化和启动实时同步核心组件
+ */
+@Service
+@Slf4j
+public class RealtimeSyncService {
+
+ @Autowired
+ private PubSubEventSubscriberImpl eventSubscriber;
+
+ @Autowired
+ private ObjectRegistryImpl objectRegistry;
+
+ private boolean started = false;
+
+ /**
+ * 启动实时同步服务
+ */
+ @PostConstruct
+ public void start() {
+ log.info("开始启动实时同步服务");
+
+ try {
+ // 刷新对象注册表
+ objectRegistry.refreshRegistry();
+
+ // 启动事件订阅
+ eventSubscriber.startSubscription();
+
+ started = true;
+ log.info("实时同步服务启动成功");
+ } catch (Exception e) {
+ log.error("启动实时同步服务时发生异常: {}", e.getMessage(), e);
+ }
+ }
+
+ /**
+ * 停止实时同步服务
+ */
+ @PreDestroy
+ public void stop() {
+ log.info("开始停止实时同步服务");
+
+ try {
+ if (started) {
+ // 停止事件订阅
+ eventSubscriber.stopSubscription();
+ log.info("实时同步服务停止成功");
+ } else {
+ log.info("实时同步服务未启动,跳过停止操作");
+ }
+ } catch (Exception e) {
+ log.error("停止实时同步服务时发生异常: {}", e.getMessage(), e);
+ }
+ }
+
+ /**
+ * 重启实时同步服务
+ */
+ public void restart() {
+ log.info("开始重启实时同步服务");
+
+ try {
+ // 停止服务
+ stop();
+
+ // 启动服务
+ start();
+
+ log.info("实时同步服务重启成功");
+ } catch (Exception e) {
+ log.error("重启实时同步服务时发生异常: {}", e.getMessage(), e);
+ }
+ }
+
+ /**
+ * 刷新对象注册表
+ */
+ public void refreshObjectRegistry() {
+ log.info("开始刷新对象注册表");
+
+ try {
+ objectRegistry.refreshRegistry();
+ log.info("对象注册表刷新成功");
+ } 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/realtime/impl/DataSynchronizerImpl.java b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/impl/DataSynchronizerImpl.java
new file mode 100644
index 00000000..579d1bc6
--- /dev/null
+++ b/datai-scenes/datai-scene-salesforce/datai-salesforce-integration/src/main/java/com/datai/integration/realtime/impl/DataSynchronizerImpl.java
@@ -0,0 +1,205 @@
+package com.datai.integration.realtime.impl;
+
+import com.datai.integration.mapper.CustomMapper;
+import com.datai.integration.model.domain.DataiIntegrationObject;
+import com.datai.integration.realtime.DataSynchronizer;
+import com.datai.integration.service.IDataiIntegrationObjectService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import java.util.Date;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 数据同步器实现类
+ * 用于将Salesforce Change Events变更数据同步至本地数据库
+ */
+@Component
+@Slf4j
+public class DataSynchronizerImpl implements DataSynchronizer {
+
+ @Autowired
+ private CustomMapper customMapper;
+
+ @Autowired
+ private IDataiIntegrationObjectService dataiIntegrationObjectService;
+
+ @Override
+ public void synchronizeData(String objectType, String recordId, String changeType, Map changeData, Date changeDate) {
+ log.info("开始同步数据变更: 对象={}, 记录ID={}, 变更类型={}", objectType, recordId, changeType);
+
+ try {
+ // 检查对象是否开启实时同步
+ if (!isObjectRealtimeSyncEnabled(objectType)) {
+ log.info("对象 {} 未开启实时同步,跳过数据同步", objectType);
+ return;
+ }
+
+ // 执行upsert操作
+ upsertData(objectType, recordId, changeData);
+
+ // 更新对象的最后同步时间
+ updateLastSyncTime(objectType, changeDate);
+
+ log.info("数据变更同步成功: 对象={}, 记录ID={}", objectType, recordId);
+ } catch (Exception e) {
+ log.error("同步数据变更时发生异常: {}", e.getMessage(), e);
+ }
+ }
+
+ /**
+ * 检查对象是否开启实时同步
+ * @param objectType 对象类型
+ * @return 是否开启实时同步
+ */
+ private boolean isObjectRealtimeSyncEnabled(String objectType) {
+ log.info("检查对象是否开启实时同步: {}", objectType);
+
+ try {
+ // 查询对象信息
+ DataiIntegrationObject queryObject = new DataiIntegrationObject();
+ queryObject.setApi(objectType);
+ List objects = dataiIntegrationObjectService.selectDataiIntegrationObjectList(queryObject);
+
+ if (objects != null && !objects.isEmpty()) {
+ DataiIntegrationObject object = objects.get(0);
+ boolean isRealtimeSync = Boolean.TRUE.equals(object.getIsRealtimeSync());
+ log.info("对象 {} 的实时同步状态: {}", objectType, isRealtimeSync);
+ return isRealtimeSync;
+ }
+
+ log.warn("对象 {} 不存在,默认不开启实时同步", objectType);
+ return false;
+ } catch (Exception e) {
+ log.error("检查对象实时同步状态时发生异常: {}", e.getMessage(), e);
+ return false;
+ }
+ }
+
+ /**
+ * 执行upsert操作
+ * @param objectType 对象类型
+ * @param recordId 记录ID
+ * @param changeData 变更数据
+ */
+ private void upsertData(String objectType, String recordId, Map changeData) {
+ log.info("执行upsert操作: 对象={}, 记录ID={}", objectType, recordId);
+
+ try {
+ // 检查记录是否存在
+ boolean exists = checkRecordExists(objectType, recordId);
+
+ if (exists) {
+ // 更新记录
+ updateRecord(objectType, recordId, changeData);
+ log.info("更新记录成功: 对象={}, 记录ID={}", objectType, recordId);
+ } else {
+ // 插入记录
+ insertRecord(objectType, changeData);
+ log.info("插入记录成功: 对象={}, 记录ID={}", objectType, recordId);
+ }
+ } catch (Exception e) {
+ log.error("执行upsert操作时发生异常: {}", e.getMessage(), e);
+ throw e;
+ }
+ }
+
+ /**
+ * 检查记录是否存在
+ * @param objectType 对象类型
+ * @param recordId 记录ID
+ * @return 记录是否存在
+ */
+ private boolean checkRecordExists(String objectType, String recordId) {
+ log.info("检查记录是否存在: 对象={}, 记录ID={}", objectType, recordId);
+
+ try {
+ // 构建查询条件
+ Map condition = new java.util.HashMap<>();
+ condition.put("Id", recordId);
+
+ // 执行查询
+ int count = customMapper.countBySQL(objectType, condition);
+ boolean exists = count > 0;
+ log.info("记录存在状态: 对象={}, 记录ID={}, 存在={}", objectType, recordId, exists);
+ return exists;
+ } catch (Exception e) {
+ log.error("检查记录是否存在时发生异常: {}", e.getMessage(), e);
+ return false;
+ }
+ }
+
+ /**
+ * 插入记录
+ * @param objectType 对象类型
+ * @param changeData 变更数据
+ */
+ private void insertRecord(String objectType, Map changeData) {
+ log.info("插入记录: 对象={}", objectType);
+
+ try {
+ // 构建字段名列表
+ List fields = new java.util.ArrayList<>(changeData.keySet());
+
+ // 构建字段值列表
+ List