Java定时同步流程如何统一

wen java案例 30

本文目录导读:

Java定时同步流程如何统一

  1. 统一的核心设计原则
  2. 统一架构组件图
  3. 具体实现步骤(Java + Spring Boot 示例)
  4. 进阶优化(可选)
  5. 统一后的好处

针对Java定时同步流程的统一,核心目标是 解耦、可观测、易维护,避免在业务代码中到处散落 @ScheduledTimer 导致难以管理。

以下是一套统一的架构方案,包含设计原则、核心组件和代码示例。

统一的核心设计原则

  1. 配置外部化:触发时间、开关、参数不写死,放在配置文件或配置中心。
  2. 任务集中注册:所有同步任务在一个地方登记,方便查看和管理。
  3. 状态统一追踪:每个同步任务都有统一的执行记录(开始时间、结束时间、成功/失败、耗时、异常信息),方便监控和排错。
  4. 错误处理统一化:定义通用的重试、告警、降级逻辑。
  5. 资源隔离:长时间运行的任务使用独立的线程池,避免影响系统核心请求。

统一架构组件图

┌─────────────────────────────────────────────────────┐
│                     统一调度层                        │
│  (Spring Schedule / Quartz / 自研调度框架)           │
│  - 解析配置/配置中心中的任务定义                      │
│  - 管理任务的生命周期(启动、暂停、停止)             │
└─────────────┬───────────────────────────────────────┘
              │ 触发
              ▼
┌─────────────────────────────────────────────────────┐
│                  任务执行器层                        │
│   SyncTaskExecutor<T>                               │
│  - 统一的 execute(taskName, Supplier<Result>) 模板   │
│  - 执行前:检查分布式锁、写开始日志                  │
│  - 执行中:捕获异常、记录耗时                        │
│  - 执行后:写结束日志、清理上下文、触发告警          │
└─────────────┬───────────────────────────────────────┘
              │ 调用
              ▼
┌─────────────────────────────────────────────────────┐
│             具体业务同步逻辑层                       │
│  - UserSyncServiceImpl: 同步用户数据                 │
│  - OrderSyncServiceImpl: 同步订单数据                │
│  - 只关注业务逻辑本身                                │
└─────────────────────────────────────────────────────┘
              │ 写入
              ▼
┌─────────────────────────────────────────────────────┐
│                  基础设施层                          │
│  - 数据库 (任务执行记录表)                           │
│  - 消息队列 (可选,用于解耦触发)                     │
│  - 分布式锁 (Redis/ZK)                              │
│  - 告警系统 (飞书/钉钉/短信)                         │
└─────────────────────────────────────────────────────┘

具体实现步骤(Java + Spring Boot 示例)

定义统一的任务元数据接口

// 所有同步任务必须实现此接口
public interface SyncTask {
    // 任务唯一标识
    String getTaskName();
    // 执行同步的核心逻辑,返回执行结果
    SyncResult execute(SyncContext context);
    // 可选的降级逻辑
    default void fallback(SyncContext context, Exception e) {
        log.error("Task [{}] fallback triggered: {}", getTaskName(), e.getMessage());
    }
}

统一的任务执行器(核心模板方法)

@Component
@Slf4j
public class UnifiedSyncTaskExecutor {
    @Resource
    private LockService lockService; // 分布式锁服务
    @Resource
    private SyncRecordRepository recordRepo; // 执行记录持久化
    @Resource
    private AlertService alertService; // 告警服务
    @Resource(name = "syncThreadPool")
    private ThreadPoolExecutor syncPool; // 专用线程池
    /**
     * 统一的执行入口
     */
    public SyncResult execute(SyncTask task) {
        String taskName = task.getTaskName();
        SyncContext context = new SyncContext(taskName);
        SyncResult result = new SyncResult();
        try {
            // 1. 分布式锁检查(防止集群重复执行)
            String lockKey = "sync:lock:" + taskName;
            if (!lockService.tryLock(lockKey, 10, TimeUnit.MINUTES)) {
                log.warn("Task [{}] skipped due to lock conflict.", taskName);
                return result.skipped("Lock already held by another instance.");
            }
            // 2. 记录任务开始
            context.setStartTime(LocalDateTime.now());
            recordRepo.save(SyncRecord(running, taskName, context));
            // 3. 真正的业务逻辑执行
            result = task.execute(context);
            // 4. 记录成功结果
            context.setEndTime(LocalDateTime.now());
            context.setCostMs(Duration.between(context.getStartTime(), context.getEndTime()).toMillis());
            recordRepo.updateStatus(taskName, SyncStatus.SUCCESS, context);
        } catch (Exception e) {
            // 5. 统一异常处理
            log.error("Task [{}] failed.", taskName, e);
            context.setEndTime(LocalDateTime.now());
            context.setErrorMessage(e.getMessage());
            recordRepo.updateStatus(taskName, SyncStatus.FAILED, context);
            // 6. 触发降级
            task.fallback(context, e);
            // 7. 发送告警
            alertService.sendAlert("同步任务失败", taskName, e.getMessage());
            result.setSuccess(false);
            result.setErrorMessage(e.getMessage());
        } finally {
            // 8. 释放锁
            lockService.unlock(lockKey);
        }
        return result;
    }
    // 支持异步执行
    public CompletableFuture<SyncResult> executeAsync(SyncTask task) {
        return CompletableFuture.supplyAsync(() -> execute(task), syncPool);
    }
}

配置层:读取调度规则

# application.yml 或 Nacos配置中心
sync:
  tasks:
    user-sync:
      enabled: true
      cron: "0 0 1 * * ?"  # 每天凌晨1点
      description: "同步用户基础信息"
      retryCount: 3
    order-sync:
      enabled: true
      cron: "0 */30 * * * ?" # 每30分钟
      description: "同步未完成的订单"
      pageSize: 500

任务注册与调度(Spring Schedule + 动态配置)

@Component
@Slf4j
public class SyncTaskScheduler {
    @Resource
    private UnifiedSyncTaskExecutor executor;
    // 集中注册所有任务(Bean方式,也可用配置驱动)
    private Map<String, SyncTask> taskMap = new ConcurrentHashMap<>();
    @PostConstruct
    public void registerTasks() {
        // 方式1:代码手动注册
        taskMap.put("user-sync", new UserSyncTask());
        taskMap.put("order-sync", new OrderSyncTask());
        // 方式2:自动扫描 @Component 的 SyncTask
        // applicationContext.getBeansOfType(SyncTask.class).values()
        //      .forEach(task -> taskMap.put(task.getTaskName(), task));
    }
    // 统一调度入口,由Spring Cron触发
    @Scheduled(cron = "${sync.tasks.user-sync.cron}")
    public void triggerUserSync() {
        SyncTask task = taskMap.get("user-sync");
        if (task != null && isTaskEnabled("user-sync")) {
            executor.execute(task);
        }
    }
    @Scheduled(cron = "${sync.tasks.order-sync.cron}")
    public void triggerOrderSync() {
        SyncTask task = taskMap.get("order-sync");
        if (task != null && isTaskEnabled("order-sync")) {
            executor.execute(task);
        }
    }
    // 支持运行时启停
    private boolean isTaskEnabled(String taskName) {
        // 从配置中心实时读取
        return configService.getBoolean("sync.tasks." + taskName + ".enabled", true);
    }
}

具体业务实现:只关注逻辑

@Component
public class UserSyncTask implements SyncTask {
    @Resource
    private UserSyncService userSyncService;
    @Override
    public String getTaskName() {
        return "user-sync";
    }
    @Override
    public SyncResult execute(SyncContext context) {
        log.info("开始同步用户数据...");
        // 纯业务逻辑,无需关心锁、日志、异常处理
        int count = userSyncService.syncAllUsers();
        return SyncResult.success("同步用户 " + count + " 条");
    }
}

关键:数据库同步记录表(用于监控和排错)

CREATE TABLE sync_task_record (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    task_name VARCHAR(100) NOT NULL COMMENT '任务名称',
    status VARCHAR(20) NOT NULL COMMENT 'RUNNING/SUCCESS/FAILED/SKIPPED',
    start_time DATETIME,
    end_time DATETIME,
    cost_ms BIGINT COMMENT '执行耗时(毫秒)',
    error_message TEXT,
    result_message VARCHAR(500),
    instance_id VARCHAR(100) COMMENT '执行实例IP',
    created_time DATETIME DEFAULT CURRENT_TIMESTAMP,
    INDEX idx_task_name_time (task_name, start_time)
);

进阶优化(可选)

  1. 基于配置中心动态调度
    • 使用 Nacos / Apollo,当修改 cronenabled 时,通过监听刷新 ScheduledFuture,实现不停机修改定时周期。
  2. 任务分片
    • 对于百万级数据同步,可将任务拆分为多个分片,分别执行。
      public interface ShardingSyncTask extends SyncTask {
      int totalShards();  // 总分片数
      void executeShard(ShardContext shardContext); // 执行特定分片
      }
  3. DAG 任务依赖
    • 如果任务有先后顺序(如:订单同步完成后才能同步订单明细),可引入 DAG 调度器,定义节点间的依赖关系。
  4. 可视化监控
    • 配合 Prometheus + Grafana,暴露 sync_task_duration_secondssync_task_status 指标。

统一后的好处

维度 统一前(@Scheduled 散落各地) 统一后
代码 重复的锁、日志、异常处理 只关注业务逻辑
运维 靠人力查看日志,排查慢 直接查 sync_task_record 表,一目了然
监控 缺失,任务挂掉不知道 失败自动告警,成功记录耗时
扩展 新增任务要改调度器 只需实现 SyncTask 接口即可
集群 手动处理重复执行 分布式锁统一保障

这套方案将“定时同步”从业务代码中的副作用转变为独立、可管理的基础设施,是大型项目中推荐的标准化做法。

抱歉,评论功能暂时关闭!