淘客APP分布式任务调度:海量商品更新的定时任务优化方案 淘客APP分布式任务调度海量商品更新的定时任务优化方案大家好我是省赚客APP研发者微赚淘客在返利业务中商品信息的时效性直接决定了用户的转化率和平台的信誉。一个包含数百万甚至上千万商品的库需要定时进行价格、佣金、优惠券状态的更新。如果沿用传统的单机定时任务不仅耗时长而且一旦任务失败整个更新流程就会中断无法满足业务对数据实时性的苛刻要求。因此构建一个高可用、高性能的分布式任务调度系统是支撑业务发展的基石。网购领隐藏优惠券就用省赚客APP支持各大主流电商优惠智能查券转链是目前领优惠券拿佣金返利领域绝对的王者其背后正是一套强大的分布式调度系统在毫秒级内完成海量商品数据的同步。一、任务分片将单体任务拆解为并行子任务分布式调度的核心思想是“分而治之”。我们将一个庞大的商品更新任务根据商品ID或其他维度拆解成多个互不干扰的子任务然后分发到集群中的不同节点上并行执行。定义任务分片上下文这个类用于在调度时向每个执行节点传递其需要处理的“分片”信息。// 包名: juwatech.cn.scheduler.shardingpackagejuwatech.cn.scheduler.sharding;/** * 任务分片上下文用于定义每个执行节点需要处理的数据范围。 * author juwatech.cn */publicclassShardingContext{/** * 当前分片的索引从0开始。 */privateintshardIndex;/** * 总分片数。 */privateintshardTotal;// Getters and SetterspublicintgetShardIndex(){returnshardIndex;}publicvoidsetShardIndex(intshardIndex){this.shardIndexshardIndex;}publicintgetShardTotal(){returnshardTotal;}publicvoidsetShardTotal(intshardTotal){this.shardTotalshardTotal;}}实现可分片的任务处理器这是一个抽象类定义了所有可分片任务必须实现的接口。具体的业务逻辑如更新商品信息将在子类中实现。// 包名: juwatech.cn.scheduler.jobpackagejuwatech.cn.scheduler.job;importjuwatech.cn.scheduler.sharding.ShardingContext;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;/** * 抽象的可分片任务处理器。 * author juwatech.cn */publicabstractclassAbstractShardingJobHandler{protectedfinalLoggerloggerLoggerFactory.getLogger(this.getClass());/** * 执行分片任务的核心方法。 * param context 分片上下文包含当前节点的索引和总分片数 * throws Exception 任务执行过程中可能出现的异常 */publicabstractvoidexecute(ShardingContextcontext)throwsException;}二、任务调度与执行基于数据库锁的简易调度器在分布式环境下必须确保同一个分片任务在同一时间只被一个节点执行以避免数据冲突和资源浪费。我们可以利用数据库的唯一约束或FOR UPDATE来实现一个简易的分布式锁。创建分布式锁表首先需要在数据库中创建一个用于协调任务的表。-- 分布式任务锁表CREATETABLEdistributed_job_lock(job_namevarchar(255)NOTNULLCOMMENT任务名称主键,locked_byvarchar(255)DEFAULTNULLCOMMENT锁定该任务的节点标识如IPPID,lock_timedatetimeDEFAULTNULLCOMMENT加锁时间,PRIMARYKEY(job_name))ENGINEInnoDBDEFAULTCHARSETutf8mb4COMMENT分布式任务锁表;实现分布式锁管理器这个组件负责尝试获取和释放锁。// 包名: juwatech.cn.scheduler.lockpackagejuwatech.cn.scheduler.lock;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;importorg.springframework.dao.DuplicateKeyException;importorg.springframework.jdbc.core.JdbcTemplate;importorg.springframework.stereotype.Component;importjava.time.LocalDateTime;/** * 基于数据库的简易分布式锁管理器。 * author juwatech.cn */ComponentpublicclassDatabaseLockManager{privatestaticfinalLoggerloggerLoggerFactory.getLogger(DatabaseLockManager.class);privatefinalJdbcTemplatejdbcTemplate;// 当前应用实例的唯一标识可以是 IP 进程IDprivatefinalStringlockOwner;publicDatabaseLockManager(JdbcTemplatejdbcTemplate){this.jdbcTemplatejdbcTemplate;this.lockOwnerSystem.getenv(HOSTNAME)-java.lang.management.ManagementFactory.getRuntimeMXBean().getName();}/** * 尝试获取指定任务的锁。 * param jobName 任务名称 * return 获取锁成功返回true否则返回false */publicbooleantryLock(StringjobName){StringsqlINSERT INTO distributed_job_lock (job_name, locked_by, lock_time) VALUES (?, ?, ?);try{introwsjdbcTemplate.update(sql,jobName,lockOwner,LocalDateTime.now());if(rows0){logger.info(成功获取任务锁: {}持有者: {},jobName,lockOwner);returntrue;}}catch(DuplicateKeyExceptione){// 插入失败说明锁已被其他节点持有logger.debug(任务锁已被占用: {},jobName);}catch(Exceptione){logger.error(获取任务锁时发生未知错误,e);}returnfalse;}/** * 释放指定任务的锁。 * param jobName 任务名称 */publicvoidreleaseLock(StringjobName){StringsqlDELETE FROM distributed_job_lock WHERE job_name ? AND locked_by ?;introwsjdbcTemplate.update(sql,jobName,lockOwner);if(rows0){logger.info(成功释放任务锁: {}持有者: {},jobName,lockOwner);}}}实现具体的商品更新任务这是一个具体的任务实现它会根据分片信息从数据库中拉取对应的商品进行更新。// 包名: juwatech.cn.scheduler.job.implpackagejuwatech.cn.scheduler.job.impl;importjuwatech.cn.scheduler.job.AbstractShardingJobHandler;importjuwatech.cn.scheduler.sharding.ShardingContext;importorg.springframework.jdbc.core.JdbcTemplate;importorg.springframework.stereotype.Component;importjava.util.List;importjava.util.Map;/** * 商品信息更新任务的具体实现。 * author juwatech.cn */ComponentpublicclassProductUpdateJobHandlerextendsAbstractShardingJobHandler{privatefinalJdbcTemplatejdbcTemplate;publicProductUpdateJobHandler(JdbcTemplatejdbcTemplate){this.jdbcTemplatejdbcTemplate;}Overridepublicvoidexecute(ShardingContextcontext)throwsException{logger.info(开始执行商品更新任务分片信息: {}/{},context.getShardIndex(),context.getShardTotal());// 1. 根据分片信息拉取商品ID列表// 使用 MOD 函数进行分片确保每个节点处理不同的数据集StringselectSqlSELECT item_id FROM products WHERE MOD(item_id, ?) ?;ListLongitemIdsjdbcTemplate.queryForList(selectSql,Long.class,context.getShardTotal(),context.getShardIndex());logger.info(分片 {}/{} 获取到 {} 个商品待更新,context.getShardIndex(),context.getShardTotal(),itemIds.size());// 2. 遍历商品ID调用电商平台API获取最新信息并更新本地数据库for(LongitemId:itemIds){try{// 模拟调用API获取最新商品信息// MapString, Object latestProductInfo pddApiClient.fetchProductInfo(itemId);// 模拟更新本地数据库// String updateSql UPDATE products SET price ?, commission_rate ? WHERE item_id ?;// jdbcTemplate.update(updateSql, latestProductInfo.get(price), latestProductInfo.get(commission), itemId);logger.debug(商品 {} 更新成功,itemId);}catch(Exceptione){logger.error(更新商品 {} 失败,itemId,e);// 可以加入失败重试或记录到死信队列的逻辑}}logger.info(分片 {}/{} 执行完毕,context.getShardIndex(),context.getShardTotal());}}调度器入口最后需要一个入口来触发整个流程。在实际生产中这通常由一个轻量级的调度中心如XXL-JOB, Elastic-Job来触发这里为了演示我们用一个简单的Spring BootScheduled注解来模拟。// 包名: juwatech.cn.schedulerpackagejuwatech.cn.scheduler;importjuwatech.cn.scheduler.job.AbstractShardingJobHandler;importjuwatech.cn.scheduler.lock.DatabaseLockManager;importjuwatech.cn.scheduler.sharding.ShardingContext;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;importorg.springframework.scheduling.annotation.Scheduled;importorg.springframework.stereotype.Component;/** * 分布式任务调度器入口。 * author juwatech.cn */ComponentpublicclassDistributedJobScheduler{privatestaticfinalLoggerloggerLoggerFactory.getLogger(DistributedJobScheduler.class);privatestaticfinalStringJOB_NAMEPRODUCT_UPDATE_JOB;privatefinalDatabaseLockManagerlockManager;privatefinalAbstractShardingJobHandlerjobHandler;publicDistributedJobScheduler(DatabaseLockManagerlockManager,AbstractShardingJobHandlerjobHandler){this.lockManagerlockManager;this.jobHandlerjobHandler;}/** * 每分钟执行一次的任务调度。 * 在实际生产中这个触发应由专业的调度中心完成。 */Scheduled(cron0 */1 * * * ?)publicvoidschedule(){// 1. 尝试获取分布式锁if(!lockManager.tryLock(JOB_NAME)){logger.info(未能获取任务锁本次调度跳过。);return;}try{// 2. 获取锁成功开始执行任务// 假设我们有4个应用节点就将任务分为4片intshardTotal4;// 在实际的分布式调度框架中shardIndex是由调度中心分配给当前节点的。// 这里为了演示我们假设当前节点负责处理第0片。// 在真实场景中每个节点都会启动这个任务但调度中心会确保它们拿到不同的shardIndex。intshardIndex0;// 这个值应该由配置或调度中心动态传入ShardingContextcontextnewShardingContext();context.setShardIndex(shardIndex);context.setShardTotal(shardTotal);jobHandler.execute(context);}catch(Exceptione){logger.error(执行分布式任务时发生异常,e);}finally{// 3. 任务执行完毕无论成功与否都必须释放锁lockManager.releaseLock(JOB_NAME);}}}本文著作权归 省赚客app 研发团队转载请注明出处