分布式系统中接口时序不确定性处理

发布时间:2026/7/30 23:25:05

分布式系统中接口时序不确定性处理 分布式系统中接口时序不确定性处理一、核心问题在分布式系统中多个独立接口由外部系统分别调用无法保证到达顺序。当接口 B 的处理依赖接口 A 的结果时如果 B 先到、A 后到B 的处理必然失败。期望时序 实际可能 A(创建订单) → B(补充条码) B(补充条码) → A(创建订单) ↓ ↓ ↓ ↓ 订单存在 → 条码关联成功 订单不存在 → 条码关联失败 ❌产生原因外部系统多个模块独立回调彼此不感知时序网络延迟不均先发的请求可能后到外部系统内部有队列/异步处理不同类型消息的处理速度不同多个微服务分别推送无全局编排不处理的后果依赖方处理失败需等定时任务补偿延迟从秒级变为分钟级甚至小时级失败日志堆积告警噪音定时任务负载增大用户体验下降数据延迟可见注博客https://blog.csdn.net/badao_liumang_qizhi二、解决思路核心原则谁后到谁负责串联方案一后到方主动触发先到方 ┌─────────────────────────────────────────────┐ │ B 先到 → 标记待处理等待 A │ │ A 后到 → 完成自身逻辑 → 检查 B 是否在等 → 触发 B │ └─────────────────────────────────────────────┘ 方案二先到方轮询等待 ┌─────────────────────────────────────────────┐ │ B 先到 → 检查 A 是否完成 → 未完成 → 进入重试 │ │ A 后到 → 完成自身逻辑 │ │ 定时重试 → B 再次检查 → A 已完成 → B 处理成功 │ └─────────────────────────────────────────────┘ 方案三聚合后统一处理 ┌─────────────────────────────────────────────┐ │ A 到达 → 写入聚合表 │ │ B 到达 → 写入聚合表 → 检查聚合完整性 → 完整则处理 │ └─────────────────────────────────────────────┘三、模式分类模式 1后到方主动触发推荐特点时延最短秒级不依赖定时任务响应及时。A 后到时完成自身 → 反查 B 的状态 → B 待处理则投递 MQ 触发 B B 后到时正常执行A 已完成依赖条件满足适用场景两个接口有明确的依赖关系且后到方能够通过数据关联找到先到方的记录。模式 2失败重试 定时补偿特点实现简单但延迟较高。B 先到 → 尝试执行 → 依赖不满足 → 标记失败 定时任务 → 扫描失败记录 → 重新投递 → 此时 A 已完成 → B 成功适用场景时效性要求不高或改造成本较高时的兜底方案。模式 3聚合等待模式特点所有条件到齐后才处理不存在失败重试。A 到达 → 写片段到聚合表 B 到达 → 写片段到聚合表 → 检查是否全部到齐 → 到齐则执行适用场景多方数据需要全部到齐才能处理如多物流节点聚合、分片数据汇总。模式 4延迟消费特点通过延迟队列给依赖方留出到达时间。B 先到 → 放入延迟队列等待 N 秒→ N 秒后消费 → 此时 A 大概率已到适用场景两个接口的时间差通常很小几秒内用固定延迟覆盖绝大多数情况。四、通用示例代码4.1 状态表设计CREATETABLEevent_dependency_log(idBIGINTAUTO_INCREMENTPRIMARYKEY,event_typeVARCHAR(32)NOTNULLCOMMENT事件类型,biz_keyVARCHAR(128)NOTNULLCOMMENT业务唯一键,payloadTEXTNOTNULLCOMMENT原始数据(JSON),statusCHAR(1)NOTNULLDEFAULTOCOMMENTO-待处理 P-失败 Y-成功,depend_onVARCHAR(32)COMMENT依赖的事件类型,error_msgVARCHAR(512),retry_countINTNOTNULLDEFAULT0,create_timeDATETIMENOTNULLDEFAULTCURRENT_TIMESTAMP,update_timeDATETIMENOTNULLDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP,UNIQUEINDEXuk_event_biz(event_type,biz_key),INDEXidx_status(status,event_type))COMMENT事件依赖日志表;4.2 模式 1 完整实现后到方主动触发接口 A被依赖方如创建订单ServiceSlf4jpublicclassEventAService{ResourceprivateEventDependencyLogRepositorylogRepository;ResourceprivateEventAProcessorprocessor;ResourceprivateDependencyTriggerServicetriggerService;/** * 接收事件 A被依赖方. * 完成自身逻辑后主动检查并触发依赖自己的事件 B。 */Transactional(rollbackForException.class)publicvoidhandleEventA(StringbizKey,Objectpayload){// 1. 落日志LonglogIdsaveLog(EVENT_A,bizKey,payload);if(logIdnull)return;// 已成功幂等返回// 2. 执行核心业务逻辑如创建订单、生成出库单ProcessResultresultprocessor.process(payload);// 3. 标记自身成功markSuccess(logId);// 4. 【关键】主动触发依赖自己的事件 BtriggerService.triggerDependentEvent(bizKey,EVENT_B);}}接口 B依赖方如补充条码ServiceSlf4jpublicclassEventBService{ResourceprivateEventDependencyLogRepositorylogRepository;ResourceprivateEventBMqSendermqSender;/** * 接收事件 B依赖方. * 只落日志并投递 MQ不关心依赖是否满足由消费者判断。 */Transactional(rollbackForException.class)publicvoidhandleEventB(StringbizKey,Objectpayload){LonglogIdsaveLog(EVENT_B,bizKey,payload);if(logIdnull)return;// 事务提交后投递 MQfinalLongfinalLogIdlogId;TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronizationAdapter(){OverridepublicvoidafterCommit(){mqSender.send(finalLogId);}});}}事件 B 的消费者ComponentSlf4jpublicclassEventBConsumer{ResourceprivateEventDependencyLogRepositorylogRepository;ResourceprivateEventBProcessorprocessor;RabbitListener(queues${mq.queue.event-b})publicvoidconsume(LonglogId){EventDependencyLogeventLoglogRepository.findById(logId).orElse(null);if(eventLognull||Y.equals(eventLog.getStatus())){return;}try{// 执行业务逻辑内部会查询事件A的结果// 如果事件 A 尚未完成这里会抛异常processor.process(eventLog.getPayload());eventLog.setStatus(Y);logRepository.saveAndFlush(eventLog);}catch(DependencyNotReadyExceptione){// 依赖未就绪标记 P等待事件 A 完成后主动触发log.info(事件B依赖未就绪, bizKey{}, 等待事件A触发,eventLog.getBizKey());eventLog.setStatus(P);eventLog.setErrorMsg(e.getMessage());logRepository.saveAndFlush(eventLog);}catch(Exceptione){log.warn(事件B处理失败, logId{},logId,e);eventLog.setStatus(P);eventLog.setRetryCount(eventLog.getRetryCount()1);eventLog.setErrorMsg(e.getMessage());logRepository.saveAndFlush(eventLog);throwe;}}}主动触发服务核心组件ServiceSlf4jpublicclassDependencyTriggerService{ResourceprivateEventDependencyLogRepositorylogRepository;ResourceprivateEventBMqSendereventBMqSender;/** * 当事件 A 完成时检查是否有依赖它的事件 B 在等待. * 如果存在且未完成事务提交后投递 MQ 触发重新消费。 * * param bizKey 关联的业务键两个事件通过此键关联 * param dependentEventType 要检查的依赖事件类型 */publicvoidtriggerDependentEvent(StringbizKey,StringdependentEventType){// 通过业务键关联查找依赖事件StringdependentBizKeyresolveDependentBizKey(bizKey);if(dependentBizKeynull){return;// 无关联关系如非该业务场景}EventDependencyLogdependentLoglogRepository.findByBizKeyAndEventType(dependentBizKey,dependentEventType);// 不存在尚未到达或已成功无需触发if(dependentLognull||Y.equals(dependentLog.getStatus())){return;}LongtargetLogIddependentLog.getId();log.info(主动触发依赖事件, eventType{}, bizKey{}, logId{},dependentEventType,dependentBizKey,targetLogId);// 事务提交后投递 MQ保证当前事务数据对消费者可见TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronizationAdapter(){OverridepublicvoidafterCommit(){eventBMqSender.send(targetLogId);}});}/** * 从事件 A 的业务键推导出事件 B 的业务键. * 两个事件可能用不同的标识需要通过数据关联。 */privateStringresolveDependentBizKey(StringeventABizKey){// 示例事件A用订单号事件B用子单号后缀// 通过关联表查出子单号SubOrdersubOrdersubOrderRepository.findByOrderNo(eventABizKey);if(subOrdernull||subOrder.getSubOrderNo()null){returnnull;}returnsubOrder.getSubOrderNo()D;}}4.3 模式 3 实现聚合等待ServiceSlf4jpublicclassAggregationService{ResourceprivateEventFragmentRepositoryfragmentRepository;ResourceprivateAggregatedProcessoraggregatedProcessor;/** * 接收事件片段. * 每个片段到达时检查是否所有片段都已到齐到齐则触发处理。 */Transactional(rollbackForException.class)publicvoidreceiveFragment(StringaggregateKey,StringfragmentType,Objectpayload){// 1. 存储片段EventFragmentfragmentnewEventFragment();fragment.setAggregateKey(aggregateKey);fragment.setFragmentType(fragmentType);fragment.setPayload(JsonUtil.toJson(payload));fragment.setArrivedTime(newDate());fragmentRepository.saveAndFlush(fragment);// 2. 检查是否所有必需片段都已到齐SetStringrequiredTypesgetRequiredFragmentTypes(aggregateKey);SetStringarrivedTypesfragmentRepository.findArrivedTypes(aggregateKey);if(arrivedTypes.containsAll(requiredTypes)){log.info(所有片段已到齐, aggregateKey{},aggregateKey);// 3. 聚合处理ListEventFragmentallFragmentsfragmentRepository.findByAggregateKey(aggregateKey);aggregatedProcessor.processAll(aggregateKey,allFragments);}else{SetStringmissingnewHashSet(requiredTypes);missing.removeAll(arrivedTypes);log.info(等待片段到达, aggregateKey{}, missing{},aggregateKey,missing);}}privateSetStringgetRequiredFragmentTypes(StringaggregateKey){// 根据业务规则确定需要哪些片段returnSets.newHashSet(ORDER_INFO,BARCODE_INFO,LOGISTICS_INFO);}}4.4 模式 4 实现延迟消费ServiceSlf4jpublicclassDelayedEventService{ResourceprivateDelayedMqSenderdelayedMqSender;/** * 接收依赖方事件延迟 N 秒后再消费. * 给被依赖方留出到达时间。 */Transactional(rollbackForException.class)publicvoidhandleWithDelay(StringbizKey,Objectpayload,intdelaySeconds){LonglogIdsaveLog(DELAYED_EVENT,bizKey,payload);if(logIdnull)return;finalLongfinalLogIdlogId;TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronizationAdapter(){OverridepublicvoidafterCommit(){// 投递到延迟队列delayedMqSender.sendWithDelay(finalLogId,delaySeconds);}});}}ComponentSlf4jpublicclassDelayedMqSender{ResourceprivateRabbitTemplaterabbitTemplate;/** * 投递延迟消息. * RabbitMQ 通过 x-delayed-message 插件或 TTL 死信队列实现。 */publicvoidsendWithDelay(LonglogId,intdelaySeconds){rabbitTemplate.convertAndSend(delayed-exchange,event.delayed,logId,message-{message.getMessageProperties().setDelay(delaySeconds*1000);returnmessage;});}}五、业务键关联策略两个独立接口通常使用不同的业务标识需要通过某种方式建立关联5.1 直接关联两个接口入参中有共同字段// 接口 A 入参orderId ORDER_001// 接口 B 入参orderId ORDER_001// 直接用 orderId 关联5.2 间接关联通过中间表// 接口 A 入参deliveryCode SO.001// 接口 B 入参subOrderNo 8900001D// 关联方式// delivery_master.id → sub_table.delivery_id (取 sub_order_no)// sub_order_no D 接口 B 的 biz_keyprivateStringresolveBizKey(StringdeliveryCode){SubTablesubsubTableRepo.findByDeliveryCode(deliveryCode);if(subnull||sub.getSubOrderNo()null){returnnull;// 非该业务场景}returnsub.getSubOrderNo()D;}5.3 规则推导// 接口 A 入参orderNo 8900001// 接口 B 入参bizKey 8900001D固定后缀 D// 关联方式字符串拼接规则privateStringderiveBizKey(StringorderNo){returnorderNoD;}六、并发安全设计时序不确定还带来并发问题两个接口可能几乎同时到达。6.1 分布式锁隔离ComponentSlf4jpublicclassConcurrencySafeConsumer{ResourceprivateDistributedLockProviderlockProvider;publicvoidconsumeWithLock(StringbizKey,LonglogId){// 同一业务键的两个事件用同一把锁StringlockKeyevent:process:bizKey;DistributedLocklocklockProvider.getLock(lockKey,60,TimeUnit.SECONDS);if(!lock.tryLock(30,TimeUnit.SECONDS)){log.warn(获取锁失败, bizKey{},bizKey);return;}try{doProcess(logId);}finally{lock.unlock();}}}6.2 不同锁粒度避免互斥如果两个事件用不同的锁需要确保不会死锁// 事件 A 的锁按操作维度StringlockKeyAorder:create:memberId;// 事件 B 的锁按单据维度StringlockKeyBbarcode:process:subOrderNo;// 两把锁互不影响A 完成后触发 B 时不会被 B 的锁阻塞// 因为触发动作是投递 MQ而非直接调用 B 的逻辑6.3 事务提交后才发消息// A 的事务创建订单 → 提交 → 释放锁 → 发送触发 B 的 MQ// B 的消费收到 MQ → 获取 B 的锁 → 查到 A 的数据已提交可见→ 处理// 时序保证// 1. A 的数据已提交 → B 能查到// 2. A 的锁已释放 → B 获取自己的锁不会被阻塞// 3. MQ 在提交后才发 → 不存在数据未提交就消费的问题七、多级保障组合生产环境通常组合使用多种策略ServiceSlf4jpublicclassRobustEventHandler{/** * 完整的时序不确定性处理流程. * * 第一级B先到尝试处理 * → 依赖不满足标记P * * 第二级A后到主动触发B秒级 * → 立即投递MQB重新消费 * * 第三级定时任务兜底分钟级 * → 扫描P状态记录重新投递 * * 幂等保证无论哪一级触发多次执行结果一致 */// 事件 A 处理完成Transactional(rollbackForException.class)publicvoidonEventACompleted(StringbizKey){// 核心逻辑...doBusinessLogic(bizKey);// 【第二级】主动触发triggerDependentEvent(bizKey);}// 事件 B MQ 消费publicvoidonEventBConsumed(LonglogId){try{process(logId);// 可能因 A 未到而失败}catch(DependencyNotReadyExceptione){markFailed(logId,e);// 【第一级】标记失败等待触发}}// 【第三级】定时补偿Scheduled(cron0 */5 * * * ?)publicvoidcompensate(){findFailedLogs().forEach(log-mqSender.send(log.getId()));}}八、适用场景与选型建议场景特征推荐模式理由时效性要求高秒级后到方主动触发最快响应不等调度周期时效性一般分钟级可接受失败重试 定时补偿实现简单改造量小两接口时间差极小5秒延迟消费固定延迟覆盖大多数情况多方数据缺一不可聚合等待数据完整性优先外部系统不可控可能长时间不回调主动触发 定时 人工多级保障九、注意事项关联查询要轻量主动触发时的反查操作应走索引避免全表扫描触发失败不能影响主流程MQ 发送放在afterCommit中且 sender 内部 try-catch不回滚主事务避免循环触发A 触发 BB 不应反过来触发 A通过事件类型区分状态过滤只触发O/P状态的记录Y状态不重复触发监控延迟分布统计从B 先到标记 P到A 触发 B 成功变 Y的时间差验证优化效果降级方案主动触发机制异常时定时任务仍能兜底不会造成数据永久不一致

相关新闻