SpringBoot對接第三方系統(tǒng)的實現(xiàn)
根據(jù)實際場景需求去選擇需要的解決方案。
HTTP客戶端選擇方案:RestTemplate、Feign、WebClient。
同步方案:全量同步、增量同步、實時同步 三種核心方案。
一、HTTP客戶端方案
Spring Boot 對接第三方接口有多種常用方案,適配不同場景,比如簡單場景用RestTemplate,微服務架構(gòu)用Feign,高并發(fā)場景用響應式的WebClient。以下是每種方案的詳細教程,包含依賴配置、代碼實現(xiàn)和核心說明。
Spring Boot 官方在文檔中推薦使用 RestTemplate(傳統(tǒng)項目)或 WebClient(響應式項目),而 Feign 作為 Spring Cloud 的一部分,也是微服務場景的首選。
方案一:RestTemplate(同步基礎款,適合簡單場景)
RestTemplate是 Spring 框架提供的同步 HTTP 客戶端,適配大多數(shù)簡單的第三方接口調(diào)用場景,Spring Boot 2.x 中可直接集成使用。
步驟1:添加依賴Spring Boot 2.x 的spring-boot-starter-web已內(nèi)置RestTemplate,在pom.xml中添加 web 依賴即可:
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>
步驟2:配置 RestTemplate Bean 創(chuàng)建配置類,將RestTemplate注入 Spring 容器,可配置超時時間等參數(shù):
@Configuration
public class RestTemplateConfig {
@Bean
public RestTemplate restTemplate() {
HttpComponentsClientHttpRequestFactory factory = new HttpComponentsClientHttpRequestFactory();
factory.setConnectTimeout(5000); // 連接超時5秒
factory.setReadTimeout(5000); // 讀取超時5秒
return new RestTemplate(factory);
}
}
步驟3: 調(diào)用第三方接口在 Service 層注入RestTemplate,分別實現(xiàn) GET 和 POST 請求調(diào)用。這里以調(diào)用模擬的用戶接口為例:
@Service
public class ThirdPartyService {
@Resource
private RestTemplate restTemplate;
// GET請求:根據(jù)ID查詢用戶
public UserDTO getUserById(Long userId) {
String url = "https://api.example.com/users/{id}";
// 占位符替換,返回結(jié)果自動轉(zhuǎn)為UserDTO
return restTemplate.getForObject(url, UserDTO.class, userId);
}
// POST請求:創(chuàng)建用戶
public UserDTO createUser(UserRequest request) {
String url = "https://api.example.com/users";
// 發(fā)送POST請求,攜帶JSON請求體,返回UserDTO
return restTemplate.postForObject(url, request, UserDTO.class);
}
}
步驟4: 定義實體類創(chuàng)建與接口請求 / 響應對應的實體類UserRequest和UserDTO:
// 請求實體
public class UserRequest {
private String username;
private String email;
// getter和setter
}
// 響應實體
public class UserDTO {
private Long id;
private String username;
private String email;
// getter和setter
}
方案二:Feign(聲明式調(diào)用,適配微服務)
Feign 是聲明式 HTTP 客戶端,通過注解簡化請求代碼,且能與 Spring Cloud 集成實現(xiàn)負載均衡,適合微服務架構(gòu)下的第三方接口調(diào)用。
步驟1:添加依賴在pom.xml中添加 OpenFeign 依賴:
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-openfeign</artifactId>
<version>3.1.3</version>
</dependency>
</dependencies>
步驟2:啟用 Feign 客戶端在 Spring Boot 啟動類添加@EnableFeignClients注解:
@SpringBootApplication
@EnableFeignClients // 啟用Feign客戶端
public class DemoApplication {
public static void main(String[] args) {
SpringApplication.run(DemoApplication.class, args);
}
}
步驟3:定義 Feign 接口創(chuàng)建 Feign 接口,通過注解聲明第三方接口的請求規(guī)則:
// name為客戶端名稱,url為第三方接口基地址
@FeignClient(name = "user-api", url = "https://api.example.com")
public interface UserFeignClient {
@GetMapping("/users/{id}")
UserDTO getUserById(@PathVariable("id") Long userId);
@PostMapping("/users")
UserDTO createUser(@RequestBody UserRequest request);
}
步驟4:調(diào)用 Feign 接口在 Service 層注入 Feign 接口直接調(diào)用,無需手動構(gòu)建請求:
@Service
public class UserService {
@Resource
private UserFeignClient userFeignClient;
public UserDTO getUser(Long userId) {
return userFeignClient.getUserById(userId);
}
public UserDTO addUser(UserRequest request) {
return userFeignClient.createUser(request);
}
}
方案三:WebClient(響應式非阻塞,適配高并發(fā))
WebClient是 Spring WebFlux 提供的響應式 HTTP 客戶端,非阻塞 IO,適合高并發(fā)場景,Spring Boot 2.x 及以上版本支持。
步驟1:添加依賴在pom.xml中添加 WebFlux 依賴(內(nèi)置 WebClient):
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
</dependencies>
步驟2:配置 WebClient Bean創(chuàng)建配置類,統(tǒng)一配置基礎 URL 和請求頭:
@Configuration
public class WebClientConfig {
@Bean
public WebClient webClient() {
return WebClient.builder()
.baseUrl("https://api.example.com") // 第三方接口基地址
.defaultHeader("Content-Type", "application/json")
.build();
}
}
步驟3:調(diào)用第三方接口WebClient返回Mono(單結(jié)果)或Flux(多結(jié)果),通過響應式編程處理結(jié)果:
@Service
public class ReactiveThirdPartyService {
@Resource
private WebClient webClient;
// GET請求:查詢用戶
public Mono<UserDTO> getUserById(Long userId) {
return webClient.get()
.uri("/users/{id}", userId)
.retrieve() // 發(fā)送請求并接收響應
.bodyToMono(UserDTO.class); // 響應體轉(zhuǎn)為UserDTO的Mono對象
}
// POST請求:創(chuàng)建用戶
public Mono<UserDTO> createUser(UserRequest request) {
return webClient.post()
.uri("/users")
.bodyValue(request) // 設置請求體
.retrieve()
.bodyToMono(UserDTO.class);
}
}
步驟4:控制器層調(diào)用響應式接口需返回Mono或Flux對象:
@RestController
@RequestMapping("/api/users")
public class UserController {
@Resource
private ReactiveThirdPartyService service;
@GetMapping("/{id}")
public Mono<UserDTO> getUser(@PathVariable Long id) {
return service.getUserById(id);
}
@PostMapping
public Mono<UserDTO> addUser(@RequestBody UserRequest request) {
return service.createUser(request);
}
}
二、數(shù)據(jù)同步方案
方案一:定時全量同步
適用于數(shù)據(jù)量小、對實時性要求不高的場景。
實現(xiàn)思路:
- 每天凌晨2點執(zhí)行一次全量拉取
- 刪除舊數(shù)據(jù),插入新數(shù)據(jù)(或軟刪除 + 更新)
- 使用事務保證一致性
1、全量刪除 + 批量插入
@Slf4j
@Component
@RequiredArgsConstructor
public class FullSyncScheduler {
private final RestTemplate restTemplate = new RestTemplate();
private final DepartmentService departmentService;
private final UserService userService;
@Value("${third-party.api-base-url}")
private String apiBaseUrl;
// 每天凌晨2點執(zhí)行
@Scheduled(cron = "0 0 2 * * *")
public void performFullSync() {
log.info("--- 開始執(zhí)行全量同步 ---");
Instant startTime = Instant.now();
try {
// 步驟 1: 刪除本地所有數(shù)據(jù)
departmentService.remove(new QueryWrapper<>());
UserService.remove(new QueryWrapper<>());
// 步驟 2: 從第三方拉取全量數(shù)據(jù)
syncDepartments(); // 1. 同步部門
syncUsers(); // 2. 同步用戶
log.info("--- 全量同步成功完成,總耗時: {} ms ---",
Duration.between(startTime, Instant.now()).toMillis());
} catch (Exception e) {
log.error("全量同步失敗", e);
}
}
// 同步部門邏輯
private void syncDepartments() {
log.info("同步部門數(shù)據(jù)...");
Instant depStartTime = Instant.now();
// 通過第三方接口獲取數(shù)據(jù)
String url = apiBaseUrl + "/api/departments";
Department[] remoteDepartments = restTemplate.getForObject(url, Department[].class);
if (remoteDepartments == null || remoteDepartments.length == 0) {
log.warn("從第三方API獲取部門數(shù)據(jù)為空");
return;
}
List<Department> deptList = Arrays.asList(remoteDepartments);
// 批量插入到本地數(shù)據(jù)庫
departmentService.saveBatch(deptList);
log.info("部門同步完成,共 {} 個部門,耗時:{}",
remoteDepartments.length,Duration.between(depStartTime, Instant.now()).toMillis());
}
// 同步用戶邏輯
private void syncUsers() {
log.info("同步用戶數(shù)據(jù)...");
// 通過第三方接口獲取數(shù)據(jù)
String url = apiBaseUrl + "/api/users";
User[] remoteUsers = restTemplate.getForObject(url, User[].class);
if (remoteUsers == null || remoteUsers.length == 0) {
log.warn("從第三方API獲取用戶數(shù)據(jù)為空");
return;
}
List<User> userList = Arrays.asList(remoteUsers);
// 批量插入到本地數(shù)據(jù)庫
userService.saveBatch(userList);
log.info("用戶同步完成,共 {} 個用戶。", remoteUsers.length);
}
}
2、UPSERT + 刪除多余 (SaveOrUpdateBatch + Delete Not In)
@Slf4j
@Component
@RequiredArgsConstructor
public class FullSyncScheduler {
private final RestTemplate restTemplate = new RestTemplate();
private final DepartmentService departmentService;
private final UserService userService;
@Value("${third-party.api-base-url}")
private String apiBaseUrl;
// 每天凌晨2點執(zhí)行
@Scheduled(cron = "0 0 2 * * *")
public void performFullSync() {
log.info("--- 開始執(zhí)行全量同步 ---");
Instant startTime = Instant.now();
try {
// 1. 同步部門
syncDepartments();
// 2. 同步用戶
syncUsers();
log.info("--- 全量同步成功完成,總耗時: {} ms ---",
Duration.between(startTime, Instant.now()).toMillis());
} catch (Exception e) {
log.error("全量同步失敗", e);
}
}
// 同步部門邏輯
private void syncDepartments() {
log.info("同步部門數(shù)據(jù)...");
Instant depStartTime = Instant.now();
// 步驟 1: 從第三方拉取全量數(shù)據(jù)
String url = apiBaseUrl + "/api/departments";
Department[] remoteDepartments = restTemplate.getForObject(url, Department[].class);
if (remoteDepartments == null || remoteDepartments.length == 0) {
log.warn("從第三方API獲取部門數(shù)據(jù)為空");
return;
}
List<Department> deptList = Arrays.asList(remoteDepartments);
List<String> remoteIds = deptList.stream()
.map(Department::getExternalId)
.collect(Collectors.toList());
// 步驟 2: 執(zhí)行 UPSERT (更新或插入)
departmentService.saveOrUpdateBatch(deptList);
// 收集 externalId
List<String> remoteIds = deptList.stream()
.map(Department::getExternalId)
.collect(Collectors.toList());
// 步驟3:找出并刪除本地存在但遠程不存在的數(shù)據(jù)
departmentService.removeByExternalIdNotIn(remoteIds);
log.info("部門同步完成,共 {} 個部門,耗時:{}",
remoteDepartments.length,Duration.between(depStartTime, Instant.now()).toMillis());
}
// 同步用戶邏輯
private void syncUsers() {
log.info("同步用戶數(shù)據(jù)...");
// 步驟 1: 從第三方拉取全量數(shù)據(jù)
String url = apiBaseUrl + "/api/users";
User[] remoteUsers = restTemplate.getForObject(url, User[].class);
if (remoteUsers == null || remoteUsers.length == 0) {
log.warn("從第三方API獲取用戶數(shù)據(jù)為空");
return;
}
List<User> userList = Arrays.asList(remoteUsers);
// 步驟 2: 執(zhí)行 UPSERT (更新或插入)
userService.saveOrUpdateBatch(userList);
// 收集 externalId
List<String> remoteIds = userList.stream()
.map(User::getExternalId)
.collect(Collectors.toList());
// 步驟3:找出并刪除本地存在但遠程不存在的數(shù)據(jù)
userService.removeByExternalIdNotIn(remoteIds);
log.info("用戶同步完成,共 {} 個用戶。", remoteUsers.length);
}
}
幾乎在所有其他情況下,方案二都是更優(yōu)、更安全的選擇。
它能最大限度地保證數(shù)據(jù)的一致性和業(yè)務的連續(xù)性,雖然在性能上可能比方案一略遜一籌,但在絕大多數(shù)企業(yè)級應用中,數(shù)據(jù)一致性和系統(tǒng)穩(wěn)定性遠比同步快幾秒更為重要。
因此,在組織架構(gòu)同步場景中,強烈推薦使用方案二(saveOrUpdateBatch + delete not in) 。它能確保在同步過程中,業(yè)務系統(tǒng)總能查詢到有效的部門和用戶信息,避免了因同步失敗或數(shù)據(jù)真空期導致的業(yè)務異常。
方案二:定時增量同步
1、基于時間戳的增量同步(最常用)
記錄上次同步的時間戳(如 last_sync_time),每次同步時只拉取第三方系統(tǒng)中 update_time > last_sync_time 的數(shù)據(jù)。
步驟1:記錄同步時間戳
在本地數(shù)據(jù)庫中維護一張同步記錄表(如 sync_checkpoint),存儲每個同步任務的上次成功時間戳。
CREATE TABLE sync_checkpoint (
id INT PRIMARY KEY AUTO_INCREMENT,
task_name VARCHAR(50) NOT NULL COMMENT '任務名稱(如部門同步、用戶同步)',
last_sync_time DATETIME NOT NULL COMMENT '上次同步時間戳',
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);
步驟2:拉取增量數(shù)據(jù) 每次同步時,從 sync_checkpoint 讀取 last_sync_time,調(diào)用第三方 API 時傳入該時間戳,只獲取更新時間晚于該值的數(shù)據(jù)。
GET /api/departments?since=2024-05-20T10:00:00Z GET /api/users?since=2024-05-20T10:00:00Z
步驟3:更新時間戳
同步成功后,將 last_sync_time 更新為當前時間(或第三方返回的最新數(shù)據(jù)時間戳)。
優(yōu)點
- 實現(xiàn)簡單,第三方 API 通常自帶
since/update_time篩選參數(shù)。 - 資源消耗低,只處理變化的數(shù)據(jù)。
缺點
- 依賴第三方系統(tǒng)的
update_time字段準確性(若第三方未正確更新該字段,會導致數(shù)據(jù)漏同步)。 - 若同步失敗,需手動處理時間戳回滾,否則會丟失中間數(shù)據(jù)。
適用場景
- 第三方 API 支持按時間戳篩選(如釘釘、企業(yè)微信的增量接口)。
- 數(shù)據(jù)變更頻率適中,對漏同步可通過后續(xù)全量同步兜底。
@Component
@Slf4j
@RequiredArgsConstructor
public class IncrementalSyncScheduler {
private final RestTemplate restTemplate = new RestTemplate();
private final DepartmentService departmentService;
private final UserService userService;
private final SyncCheckpointService checkpointService;
@Value("${third-party.api-base-url}")
private String apiBaseUrl;
private static final String TASK_NAME = "DEPT_USER_SYNC_MP";
private static final DateTimeFormatter dtf = DateTimeFormatter.ISO_LOCAL_DATE_TIME;
// 每10分鐘執(zhí)行一次
@Scheduled(cron = "0 */10 * * * *")
public void performIncrementalSync() {
log.info("--- 開始執(zhí)行增量同步 ---");
try {
// 1. 獲取上次同步時間
LocalDateTime lastSyncTime = getLastSyncTime();
// 2. 拉取增量數(shù)據(jù)
String url = apiBaseUrl + "/api/changes?since=" + dtf.format(lastSyncTime);
ChangeEventWrapper changes = restTemplate.getForObject(url, ChangeEventWrapper.class);
if (changes == null || (changes.getDepartments().isEmpty() && changes.getUsers().isEmpty())) {
log.info("沒有增量數(shù)據(jù)。");
updateCheckpoint();
return;
}
// 3. 應用變更
applyChanges(changes);
// 4. 更新檢查點
updateCheckpoint();
log.info("--- [MyBatis-Plus] 增量同步成功完成 ---");
} catch (Exception e) {
log.error("[MyBatis-Plus] 增量同步失敗", e);
}
}
private void applyChanges(ChangeEventWrapper changes) {
if (!changes.getDepartments().isEmpty()) {
log.info("處理 {} 條部門變更...", changes.getDepartments().size());
for (ChangeEvent<Department> event : changes.getDepartments()) {
Department data = event.getData();
switch (event.getType()) {
case CREATE:
case UPDATE:
departmentService.saveOrUpdate(data);
break;
case DELETE:
departmentService.removeById(data.getId());
break;
}
}
}
if (!changes.getUsers().isEmpty()) {
log.info("處理 {} 條用戶變更...", changes.getUsers().size());
for (ChangeEvent<User> event : changes.getUsers()) {
User data = event.getData();
switch (event.getType()) {
case CREATE:
case UPDATE:
userService.saveOrUpdate(data);
break;
case DELETE:
userService.removeById(data.getId());
break;
}
}
}
}
private LocalDateTime getLastSyncTime() {
SyncCheckpoint checkpoint = checkpointService.getByTaskName(TASK_NAME);
return checkpoint != null ? checkpoint.getLastSyncTimestamp() : LocalDateTime.of(2000, 1, 1, 0, 0);
}
private void updateCheckpoint() {
SyncCheckpoint checkpoint = checkpointService.getByTaskName(TASK_NAME);
if (checkpoint == null) {
checkpoint = new SyncCheckpoint();
checkpoint.setTaskName(TASK_NAME);
}
checkpoint.setLastSyncTimestamp(LocalDateTime.now());
checkpointService.saveOrUpdate(checkpoint);
}
// 輔助類
public static class ChangeEventWrapper {
private java.util.List<ChangeEvent<Department>> departments;
private java.util.List<ChangeEvent<User>> users;
// getters and setters
}
public static class ChangeEvent<T> {
private String type;
private T data;
// getters and setters
}
}
2、方案二:基于變更 ID 的增量同步(高可靠性)
第三方系統(tǒng)為每條數(shù)據(jù)分配唯一的變更 ID(如 change_id),每次同步時只拉取 change_id > last_change_id 的數(shù)據(jù)。
變更 ID 通常按時間遞增生成。
步驟1:記錄上次變更 ID
在 sync_checkpoint 表中增加 last_change_id 字段,存儲上次同步的最大變更 ID。
ALTER TABLE sync_checkpoint ADD COLUMN last_change_id BIGINT DEFAULT 0 COMMENT '上次同步的最大變更ID';
步驟2:拉取增量數(shù)據(jù)
調(diào)用第三方 API 時傳入 last_change_id,只獲取變更 ID 更大的數(shù)據(jù)。示例 API 請求:
GET /api/changes?last_change_id=12345
步驟3:更新變更 ID
同步成功后,將 last_change_id 更新為本次同步到的最大變更 ID。
總結(jié):
- 可靠性高,變更 ID 唯一且遞增,不會漏同步或重復同步。
- 無需依賴時間戳,避免因時間偏差導致的問題。
- 第三方系統(tǒng)需支持變更 ID 篩選(并非所有 API 都提供)。
對數(shù)據(jù)一致性要求極高的場景(如金融、支付數(shù)據(jù)同步)。其實這種方案實現(xiàn)思路和時間戳類似,只是手動維護了一個自增的變更ID,用來規(guī)避時間戳未設值之類的情況。
方案三:實時同步 (Webhook)
實時同步的核心目標是 “數(shù)據(jù)變更后立即同步” ,實現(xiàn) “準實時” 或 “實時” 的數(shù)據(jù)一致性。與定時同步不同,實時同步無需依賴定時任務觸發(fā),而是由 “事件驅(qū)動” (數(shù)據(jù)變更事件觸發(fā)同步)。
以下是幾種常見的實時同步實現(xiàn)方案,從簡單到復雜,覆蓋不同技術(shù)棧和場景:
1、Webhook 回調(diào)(最常用)
第三方系統(tǒng)(如釘釘、企業(yè)微信、CRM 系統(tǒng))在數(shù)據(jù)發(fā)生變更時(如新增用戶、修改部門),主動調(diào)用你的系統(tǒng)提供的 回調(diào)接口(Webhook Endpoint) ,將變更數(shù)據(jù)推送到你的系統(tǒng),你的系統(tǒng)接收并處理這些數(shù)據(jù)。
步驟1:提供 Webhook 接口
- 在你的系統(tǒng)中開發(fā)一個公開的接口(如
/api/webhook/sync),用于接收第三方推送的變更事件。 - 接口需支持
POST請求,通常接收 JSON 格式的事件數(shù)據(jù)。
步驟2:配置第三方 Webhook
- 在第三方系統(tǒng)的管理后臺(如釘釘開放平臺),配置你的 Webhook 接口地址,并選擇需要監(jiān)聽的事件類型(如用戶新增、部門刪除)。
步驟3:接收并處理事件
- 你的系統(tǒng)接收事件數(shù)據(jù)后,解析數(shù)據(jù)內(nèi)容(如變更類型、變更數(shù)據(jù)、時間戳),并執(zhí)行同步操作(插入、更新、刪除本地數(shù)據(jù)庫)。
- 關(guān)鍵注意點:
- 簽名驗證:第三方會在請求頭中攜帶簽名(如
X-Signature),你需要驗證簽名的合法性,防止惡意請求。 - 冪等性處理:由于網(wǎng)絡重試等原因,可能會收到重復事件,需確保同步邏輯冪等(如通過
event_id去重)。 - 異步處理:接收到事件后,應立即返回響應(如 HTTP 200),再通過線程池或消息隊列異步處理同步邏輯,避免阻塞第三方的回調(diào)請求。
- 簽名驗證:第三方會在請求頭中攜帶簽名(如
@RestController
@RequestMapping("/api/webhook")
@Slf4j
public class WebhookController {
@Autowired
private SyncService syncService;
@Autowired
private WebhookSignatureService signatureService;
@PostMapping("/dingtalk")
public ResponseEntity<?> handleDingTalkWebhook(
@RequestBody String requestBody,
@RequestHeader("X-Signature") String signature,
@RequestHeader("X-Timestamp") String timestamp) {
// 1. 驗證簽名
if (!signatureService.validateSignature(requestBody, timestamp, signature)) {
log.warn("Webhook簽名驗證失敗");
return ResponseEntity.badRequest().body("Invalid signature");
}
// 2. 解析事件數(shù)據(jù)
DingTalkWebhookEvent event = JsonUtils.parseObject(requestBody, DingTalkWebhookEvent.class);
log.info("收到釘釘Webhook事件:{}", event.getEventType());
// 3. 異步處理同步邏輯(避免阻塞)
syncService.asyncProcessEvent(event);
// 4. 立即返回響應
return ResponseEntity.ok().body("{"errcode":0,"errmsg":"success"}");
}
}
總結(jié):
- 第三方系統(tǒng)支持 Webhook(如釘釘、企業(yè)微信、GitHub、PayPal 等)。
- 對實時性要求中等(秒級延遲可接受),且不希望引入復雜中間件的場景。
2、消息隊列(MQ)異步同步(高可靠)
通過 消息隊列(如 RabbitMQ、Kafka、RocketMQ)解耦數(shù)據(jù)變更源和同步目標:
- 數(shù)據(jù)變更源(如業(yè)務系統(tǒng)、第三方 API)將變更事件寫入消息隊列。
- 你的系統(tǒng)作為消費者,監(jiān)聽消息隊列,讀取事件并執(zhí)行同步操作。
步驟1:選擇并部署消息隊列
- 根據(jù)場景選擇 MQ(如 RabbitMQ 適合可靠性優(yōu)先,Kafka 適合高吞吐)。
步驟2:生產(chǎn)端寫入消息
- 數(shù)據(jù)變更時(如用戶更新),生產(chǎn)端(如業(yè)務系統(tǒng)的服務)將變更事件(如用戶 ID、變更字段、操作類型)序列化為消息,發(fā)送到 MQ 的指定主題 / 隊列。
- 示例(Java + RabbitMQ):
@Service
public class EventProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendUserUpdateEvent(User user) {
UserUpdateEvent event = new UserUpdateEvent(user.getId(), user.getName(), LocalDateTime.now());
rabbitTemplate.convertAndSend("user.sync.exchange", "user.update", event);
log.info("發(fā)送用戶更新事件:{}", user.getId());
}
}
步驟3:消費端處理消息
- 你的系統(tǒng)作為消費者,訂閱 MQ 的主題 / 隊列,接收消息后解析并同步到本地數(shù)據(jù)庫。
- 示例(Java + RabbitMQ):
@Service
public class EventConsumer {
@Autowired
private UserRepository userRepository;
@RabbitListener(queues = "user.update.queue")
public void handleUserUpdateEvent(UserUpdateEvent event) {
log.info("接收用戶更新事件:{}", event.getUserId());
// 執(zhí)行同步操作
User user = userRepository.findByExternalId(event.getUserId())
.orElseThrow(() -> new RuntimeException("用戶不存在"));
user.setName(event.getUserName());
userRepository.save(user);
}
}
步驟4:保障可靠性
- 消息持久化:將消息和隊列設置為持久化,避免 MQ 重啟后消息丟失。
- 消費者確認(Ack) :消費者處理完消息后手動發(fā)送 Ack,確保消息被成功處理。
- 死信隊列(DLQ) :處理失敗的消息(如數(shù)據(jù)庫異常)轉(zhuǎn)入死信隊列,避免阻塞正常消息,后續(xù)可人工重試。
總結(jié):
- 對數(shù)據(jù)可靠性要求高(如金融交易、訂單同步),不允許消息丟失。
- 高并發(fā)場景(如每秒數(shù)千條數(shù)據(jù)變更),需要 MQ 削峰填谷。
- 多系統(tǒng)間數(shù)據(jù)同步(如業(yè)務系統(tǒng) → 數(shù)據(jù)倉庫 → 報表系統(tǒng))。
到此這篇關(guān)于SpringBoot對接第三方系統(tǒng)的實現(xiàn)的文章就介紹到這了,更多相關(guān)SpringBoot對接第三方系統(tǒng)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java數(shù)據(jù)結(jié)構(gòu)之雙端鏈表原理與實現(xiàn)方法
這篇文章主要介紹了Java數(shù)據(jù)結(jié)構(gòu)之雙端鏈表原理與實現(xiàn)方法,簡單描述了雙端鏈表的概念、原理并結(jié)合實例形式分析了java實現(xiàn)雙端鏈表的相關(guān)操作技巧,需要的朋友可以參考下2017-10-10
Java使用TarsosDSP庫實現(xiàn)音頻的處理和格式轉(zhuǎn)換
在音頻處理領域,Java雖然有原生的音頻處理類庫,但其功能相對基礎,而TarsosDSP是一個強大的開源音頻處理庫,提供了豐富的功能,本文將介紹如何在Java中結(jié)合使用TarsosDSP庫,來實現(xiàn)音頻的處理和格式轉(zhuǎn)換,需要的朋友可以參考下2025-04-04
Springboot @Transactional使用時需注意的幾個問題記錄
本文詳細介紹了Spring Boot中使用`@Transactional`注解進行事務管理的多個方面,包括事務的隔離級別(如REPEATABLE_READ)和傳播行為(如REQUIRES_NEW),并指出了在同一個類中調(diào)用事務方法時可能遇到的問題以及解決方案,感興趣的朋友跟隨小編一起看看吧2025-01-01

