feat(payment): 交易订单定时同步任务(分层窗口扫描+ShedLock分布式锁)

- TradeSyncJob 分层窗口扫描: 支付PROCESSING(4窗口)/CLOSE纠正/退款PROGRESS(4窗口), 越新查越勤, 超7天淘汰
- TradeSyncService 跨租户引导读(runAs装载mchNo)后委托租户内同步, 共享 payment:trade:{id} 锁与关单互斥
- PayTradeManager.findSyncTrades / RefundOrderManager.findProgressRefunds+findByRefundNoNotTenant 跨租户扫描(单次上限500)
- PaySyncService 增 syncByContainer 容器视角同步入口(admin/merchant复用)
- 引入 ShedLock(6.3.1) + ShedLockConfiguration(Redis), @SchedulerLock 防多节点重复执行
- application-prod.yml 增 trade-sync-enabled 开关(默认true)
This commit is contained in:
daxpay
2026-08-02 10:44:39 +08:00
parent e27748840f
commit 647565e7e2
9 changed files with 332 additions and 0 deletions

View File

@@ -83,6 +83,20 @@ public class PayTradeManager extends BaseManager<PayTradeMapper, PayTrade> {
return this.page(mpPage, wrapper);
}
/// 按状态 + 创建时间窗口扫描(定时同步任务用)
///
/// 跨租户扫描(定时任务无 HTTP 上下文), 单次上限 500 防积压爆量。
/// 支持 processing(常规同步)和 close(CLOSE→SUCCESS 纠正)两种状态扫描。
/// 命中索引 idx_pay_trade_status_create_time。
@IgnoreTenant
public List<PayTrade> findSyncTrades(String status, OffsetDateTime start, OffsetDateTime end) {
return listLimit(500, q -> q
.eq(PayTrade::getStatus, status)
.ge(PayTrade::getCreateTime, start)
.le(PayTrade::getCreateTime, end)
.orderByAsc(PayTrade::getCreateTime));
}
/// 查询普通支付已超时但仍处理中的资金交易(兜底定时任务用)
///
/// 条件: tradeType=NORMAL 且 status=PROCESSING 且容器 expiredTime < now。

View File

@@ -3,6 +3,7 @@ package cn.daxpay.open.payment.trade.order.dao;
import cn.daxpay.open.platform.common.mybatisplus.impl.BaseManager;
import cn.daxpay.open.platform.common.mybatisplus.query.generator.QueryGenerator;
import cn.daxpay.open.platform.common.mybatisplus.util.MpUtil;
import cn.daxpay.open.platform.core.annotation.IgnoreTenant;
import cn.daxpay.open.platform.core.rest.param.PageParam;
import cn.daxpay.open.payment.trade.order.entity.RefundOrder;
import cn.daxpay.open.payment.trade.order.param.RefundOrderQuery;
@@ -10,6 +11,8 @@ import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
import org.springframework.stereotype.Repository;
import java.time.OffsetDateTime;
import java.util.List;
import java.util.Optional;
import java.util.Set;
@@ -38,6 +41,25 @@ public class RefundOrderManager extends BaseManager<RefundOrderMapper, RefundOrd
.oneOpt();
}
/// 根据退款号查询(忽略租户, 定时任务引导读用)
@IgnoreTenant
public Optional<RefundOrder> findByRefundNoNotTenant(String refundNo) {
return findByField(RefundOrder::getRefundNo, refundNo);
}
/// 按创建时间窗口扫描退款中订单(定时同步任务用)
///
/// 跨租户扫描(定时任务无 HTTP 上下文), 单次上限 500 防积压爆量。
/// 固定 status=progress, 命中索引 idx_refund_order_status_create_time。
@IgnoreTenant
public List<RefundOrder> findProgressRefunds(OffsetDateTime start, OffsetDateTime end) {
return listLimit(500, q -> q
.eq(RefundOrder::getStatus, "progress")
.ge(RefundOrder::getCreateTime, start)
.le(RefundOrder::getCreateTime, end)
.orderByAsc(RefundOrder::getCreateTime));
}
/// 根据实际上送串查询(回调容错: 特殊通道仅回传变形号)
public Optional<RefundOrder> findByRelationOrderNo(String relationOrderNo) {
return findByField(RefundOrder::getRelationOrderNo, relationOrderNo);

View File

@@ -0,0 +1,158 @@
package cn.daxpay.open.payment.trade.runtime.job;
import cn.daxpay.open.payment.trade.enums.PayFundStatusEnum;
import cn.daxpay.open.payment.trade.order.dao.PayTradeManager;
import cn.daxpay.open.payment.trade.order.dao.RefundOrderManager;
import cn.daxpay.open.payment.trade.order.entity.PayTrade;
import cn.daxpay.open.payment.trade.order.entity.RefundOrder;
import cn.daxpay.open.payment.trade.runtime.service.sync.TradeSyncService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import net.javacrumbs.shedlock.spring.annotation.SchedulerLock;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.time.OffsetDateTime;
import java.time.ZoneOffset;
import java.util.List;
/// # 交易订单定时同步任务
///
/// 分层窗口扫描支付/退款中间态订单, 调通道查单纠正本地状态。
///
/// ## 设计要点
/// - **支付 PROCESSING 段**(4 窗口): 回调丢失/超时关单失败后, 查通道真实状态纠正为 SUCCESS/CLOSE/FAIL
/// - **支付 CLOSE 纠正段**(1 窗口): 超时关单后通道实际已付款, 触发 CLOSE→SUCCESS 纠正
/// - **退款 PROGRESS 段**(4 窗口): 退款回调丢失后, 查通道真实退款状态
///
/// 订单越"新"查得越勤, 越久越稀疏, 超 7 天自然淘汰(不显式置 FAIL)。
/// 与 [NormalPayTimeoutJob] / [GatewayTimeoutJob] 共享 `payment:trade:{id}` Redis 锁, 天然互斥。
/// 单笔同步异常不阻断整批, 锁冲突(RepetitiveOperationException)跳过等下一轮。
///
/// 全局开关: `daxpay.platform.config.trade-sync-enabled`(默认 true)。
@Slf4j
@Component
@RequiredArgsConstructor
@ConditionalOnProperty(prefix = "daxpay.platform.config", name = "trade-sync-enabled", havingValue = "true", matchIfMissing = true)
public class TradeSyncJob {
private final PayTradeManager payTradeManager;
private final RefundOrderManager refundOrderManager;
private final TradeSyncService tradeSyncService;
// ==================== 支付 PROCESSING 同步(分层窗口) ====================
/// 最新窗口: 创建 1~10 分钟的 PROCESSING 订单, 每分钟同步一次
@Scheduled(cron = "0 */1 * * * ?")
@SchedulerLock(name = "lock:tradeSync:payProc10M", lockAtMostFor = "50s", lockAtLeastFor = "5s")
public void syncPayProcessing10M() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncPayBatch(PayFundStatusEnum.PROCESSING.getCode(), now.minusMinutes(10), now.minusMinutes(1));
}
/// 次新窗口: 创建 10 分钟~1 小时的 PROCESSING 订单, 每 10 分钟同步一次
@Scheduled(cron = "0 */10 * * * ?")
@SchedulerLock(name = "lock:tradeSync:payProc1H", lockAtMostFor = "8m", lockAtLeastFor = "30s")
public void syncPayProcessing1H() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncPayBatch(PayFundStatusEnum.PROCESSING.getCode(), now.minusHours(1), now.minusMinutes(10));
}
/// 陈旧窗口: 创建 1~24 小时的 PROCESSING 订单, 每小时同步一次
@Scheduled(cron = "0 0 */1 * * ?")
@SchedulerLock(name = "lock:tradeSync:payProc1D", lockAtMostFor = "50m", lockAtLeastFor = "1m")
public void syncPayProcessing1D() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncPayBatch(PayFundStatusEnum.PROCESSING.getCode(), now.minusHours(24), now.minusHours(1));
}
/// 死单兜底: 创建 1~7 天的 PROCESSING 订单, 每天凌晨同步一次
@Scheduled(cron = "0 10 1 * * ?")
@SchedulerLock(name = "lock:tradeSync:payProc7D", lockAtMostFor = "30m", lockAtLeastFor = "5m")
public void syncPayProcessing7D() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncPayBatch(PayFundStatusEnum.PROCESSING.getCode(), now.minusDays(7), now.minusHours(24));
}
// ==================== 支付 CLOSE 纠正(CLOSE→SUCCESS) ====================
/// 超时关单后通道实际已付款的纠正窗口: 创建 1~30 分钟的 CLOSE 订单, 每 5 分钟同步一次
///
/// 超过 30 分钟的 CLOSE 基本确定是真关闭(通道侧也确认未付), 不再浪费通道查询。
@Scheduled(cron = "0 */5 * * * ?")
@SchedulerLock(name = "lock:tradeSync:payCloseFix", lockAtMostFor = "4m", lockAtLeastFor = "30s")
public void syncPayCloseFix() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncPayBatch(PayFundStatusEnum.CLOSE.getCode(), now.minusMinutes(30), now.minusMinutes(1));
}
// ==================== 退款 PROGRESS 同步(分层窗口) ====================
/// 最新窗口: 创建 1~10 分钟的 PROGRESS 退款, 每分钟同步一次
@Scheduled(cron = "0 */1 * * * ?")
@SchedulerLock(name = "lock:tradeSync:refund10M", lockAtMostFor = "50s", lockAtLeastFor = "5s")
public void syncRefund10M() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncRefundBatch(now.minusMinutes(10), now.minusMinutes(1));
}
/// 次新窗口: 创建 10 分钟~1 小时的 PROGRESS 退款, 每 10 分钟同步一次
@Scheduled(cron = "0 */10 * * * ?")
@SchedulerLock(name = "lock:tradeSync:refund1H", lockAtMostFor = "8m", lockAtLeastFor = "30s")
public void syncRefund1H() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncRefundBatch(now.minusHours(1), now.minusMinutes(10));
}
/// 陈旧窗口: 创建 1~24 小时的 PROGRESS 退款, 每小时同步一次
@Scheduled(cron = "0 0 */1 * * ?")
@SchedulerLock(name = "lock:tradeSync:refund1D", lockAtMostFor = "50m", lockAtLeastFor = "1m")
public void syncRefund1D() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncRefundBatch(now.minusHours(24), now.minusHours(1));
}
/// 死单兜底: 创建 1~7 天的 PROGRESS 退款, 每天凌晨同步一次
@Scheduled(cron = "0 20 1 * * ?")
@SchedulerLock(name = "lock:tradeSync:refund7D", lockAtMostFor = "30m", lockAtLeastFor = "5m")
public void syncRefund7D() {
OffsetDateTime now = OffsetDateTime.now(ZoneOffset.UTC);
syncRefundBatch(now.minusDays(7), now.minusHours(24));
}
// ==================== 批量处理 ====================
/// 支付同步批量处理(逐笔容错, 单笔锁冲突/异常不阻断整批)
private void syncPayBatch(String status, OffsetDateTime start, OffsetDateTime end) {
List<PayTrade> trades = payTradeManager.findSyncTrades(status, start, end);
if (trades.isEmpty()) {
return;
}
log.info("定时同步扫描支付 status={} 命中 {} 笔, 开始处理", status, trades.size());
for (PayTrade trade : trades) {
try {
tradeSyncService.syncPayTrade(trade.getTradeNo());
} catch (Exception e) {
// 单笔失败不阻断整批(锁冲突/通道异常等), 记录后继续
log.warn("支付同步跳过 tradeNo={}", trade.getTradeNo(), e);
}
}
}
/// 退款同步批量处理(逐笔容错)
private void syncRefundBatch(OffsetDateTime start, OffsetDateTime end) {
List<RefundOrder> refunds = refundOrderManager.findProgressRefunds(start, end);
if (refunds.isEmpty()) {
return;
}
log.info("定时同步扫描退款命中 {} 笔, 开始处理", refunds.size());
for (RefundOrder refund : refunds) {
try {
tradeSyncService.syncRefundOrder(refund.getRefundNo());
} catch (Exception e) {
log.warn("退款同步跳过 refundNo={}", refund.getRefundNo(), e);
}
}
}
}

View File

@@ -51,6 +51,15 @@ public class PaySyncService {
private final PaySyncRecordService paySyncRecordService;
private final LockExecutor lockExecutor;
/// 按容器ID同步支付状态
///
/// 供容器视角的对外 Service(Admin/Merchant)调用, 内部反查资金凭证后委托 [syncPayOrder]。
public NormalPaySyncResult syncByContainer(Long containerId, String tradeType) {
PayTrade trade = payTradeManager.findByContainerId(containerId, tradeType)
.orElseThrow(() -> new BizInfoException(CommonErrorCode.VALIDATE_PARAMETERS_ERROR, "pay.error.payOrderNotExist"));
return this.syncPayOrder(trade);
}
/// 支付同步
@Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class)
public NormalPaySyncResult sync(NormalPaySyncParam param) {

View File

@@ -0,0 +1,91 @@
package cn.daxpay.open.payment.trade.runtime.service.sync;
import cn.daxpay.open.payment.common.context.PaymentContext;
import cn.daxpay.open.payment.trade.enums.PayFundStatusEnum;
import cn.daxpay.open.payment.trade.enums.RefundOrderStatusEnum;
import cn.daxpay.open.payment.trade.order.dao.PayTradeManager;
import cn.daxpay.open.payment.trade.order.dao.RefundOrderManager;
import cn.daxpay.open.payment.trade.order.entity.PayTrade;
import cn.daxpay.open.payment.trade.order.entity.RefundOrder;
import cn.daxpay.open.payment.trade.runtime.service.refund.RefundSyncService;
import cn.hutool.core.util.StrUtil;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.Objects;
import java.util.Set;
/// # 交易定时同步服务
///
/// 定时任务无商户登录上下文, 本服务负责:
/// ① 引导读(NotTenant)定位订单获取 mchNo;
/// ② 通过 [PaymentContext#runAs] + `setMchNo` 装载租户身份(仅 setMchNo, 不校验商户启用);
/// ③ 委托 [PaySyncService#syncPayOrder] / [RefundSyncService#syncById] 执行租户内同步。
///
/// 与 [PayCloseService#closeForTimeout] 的上下文装载范式完全对称,
/// 共享 `payment:trade:{id}` Redis 锁, 同步与关单天然互斥。
@Slf4j
@Service
@RequiredArgsConstructor
public class TradeSyncService {
/// 需要同步的支付资金态: 处理中(回调丢失纠正) + 已关闭(CLOSE→SUCCESS 纠正)
private static final Set<String> SYNC_PAY_STATUSES = Set.of(
PayFundStatusEnum.PROCESSING.getCode(),
PayFundStatusEnum.CLOSE.getCode());
private final PayTradeManager payTradeManager;
private final RefundOrderManager refundOrderManager;
private final PaySyncService paySyncService;
private final RefundSyncService refundSyncService;
private final PaymentContext paymentContext;
/// 定时同步支付订单(幂等)
///
/// 仅处理 PROCESSING 和 CLOSE 状态:
/// - PROCESSING: 回调丢失/超时关单失败后, 查通道真实状态纠正
/// - CLOSE: 超时关单后通道实际已付款的场景, 触发 CLOSE→SUCCESS 纠正
public void syncPayTrade(String tradeNo) {
// 引导读: 跨租户定位订单
PayTrade boot = payTradeManager.findByTradeNoNotTenant(tradeNo).orElse(null);
if (Objects.isNull(boot)) {
return;
}
if (!SYNC_PAY_STATUSES.contains(boot.getStatus())) {
return;
}
if (StrUtil.isBlank(boot.getMchNo())) {
log.error("定时同步交易缺少 mchNo, tradeNo={}", tradeNo);
return;
}
// 装载租户身份后执行同步(syncPayOrder 自带 Redis 锁 + REQUIRES_NEW 事务)
paymentContext.runAs(() -> {
paymentContext.setMchNo(boot.getMchNo());
paySyncService.syncPayOrder(boot);
});
}
/// 定时同步退款订单(幂等)
///
/// 仅处理 PROGRESS 状态(退款中), 查通道真实退款状态后回写结算。
public void syncRefundOrder(String refundNo) {
// 引导读: 跨租户定位退款单
RefundOrder boot = refundOrderManager.findByRefundNoNotTenant(refundNo).orElse(null);
if (Objects.isNull(boot)) {
return;
}
if (!Objects.equals(RefundOrderStatusEnum.PROGRESS.getCode(), boot.getStatus())) {
return;
}
if (StrUtil.isBlank(boot.getMchNo())) {
log.error("定时同步退款缺少 mchNo, refundNo={}", refundNo);
return;
}
// 装载租户身份后执行同步
paymentContext.runAs(() -> {
paymentContext.setMchNo(boot.getMchNo());
refundSyncService.syncById(boot.getId());
});
}
}

View File

@@ -35,6 +35,17 @@
<artifactId>lock4j-redis-template-spring-boot-starter</artifactId>
<version>${lock4j.version}</version>
</dependency>
<!-- 定时任务分布式锁(@SchedulerLock防止多节点重复执行 -->
<dependency>
<groupId>net.javacrumbs.shedlock</groupId>
<artifactId>shedlock-spring</artifactId>
<version>${shedlock.version}</version>
</dependency>
<dependency>
<groupId>net.javacrumbs.shedlock</groupId>
<artifactId>shedlock-provider-redis-spring</artifactId>
<version>${shedlock.version}</version>
</dependency>
</dependencies>
</project>

View File

@@ -0,0 +1,24 @@
package cn.daxpay.open.platform.common.redis.configuration;
import net.javacrumbs.shedlock.provider.redis.spring.RedisLockProvider;
import net.javacrumbs.shedlock.spring.annotation.EnableSchedulerLock;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
/// # 定时任务分布式锁配置(ShedLock)
///
/// 基于 Redis 实现, 与 [SpringAutoConfiguration] 上的 `@EnableScheduling` 配合,
/// 通过 `@SchedulerLock` 注解保证多节点部署时同一任务仅单节点执行。
@Configuration
@EnableSchedulerLock(defaultLockAtMostFor = "5m")
public class ShedLockConfiguration {
/// 锁 key 前缀环境名, 用于多应用共用 Redis 时隔离
private static final String ENVIRONMENT = "daxpay";
@Bean
public RedisLockProvider lockProvider(RedisConnectionFactory connectionFactory) {
return new RedisLockProvider(connectionFactory, ENVIRONMENT);
}
}

View File

@@ -129,6 +129,8 @@ daxpay:
config:
# 沙箱环境(必须关闭; 不写时 Java 默认 true 会被 DeploymentModeEnforcer 拦截)
sandbox-enabled: false
# 交易定时同步开关(生产环境开启)
trade-sync-enabled: true
# RSA 密钥(PEM 文本通过环境变量注入)
key-config:
private-key: ${RSA_PRIVATE_KEY:?missing RSA_PRIVATE_KEY}

View File

@@ -50,6 +50,7 @@
<mybatis-plus.version>3.5.16</mybatis-plus.version>
<mybatis-plus-join.version>1.5.7</mybatis-plus-join.version>
<lock4j.version>2.2.7</lock4j.version>
<shedlock.version>6.3.1</shedlock.version>
<aws-sdk.version>2.46.15</aws-sdk.version>
<sa-token.version>1.45.0</sa-token.version>
<caffeine.version>3.2.4</caffeine.version>