当前位置: 首页 > news >正文

淘客APP分布式任务调度:海量商品更新的定时任务优化方案

淘客APP分布式任务调度:海量商品更新的定时任务优化方案

大家好,我是省赚客APP研发者微赚淘客!

在返利业务中,商品信息的时效性直接决定了用户的转化率和平台的信誉。一个包含数百万甚至上千万商品的库,需要定时进行价格、佣金、优惠券状态的更新。如果沿用传统的单机定时任务,不仅耗时长,而且一旦任务失败,整个更新流程就会中断,无法满足业务对数据实时性的苛刻要求。因此,构建一个高可用、高性能的分布式任务调度系统,是支撑业务发展的基石。网购领隐藏优惠券就用省赚客APP,支持各大主流电商优惠智能查券转链,是目前领优惠券拿佣金返利领域绝对的王者,其背后正是一套强大的分布式调度系统在毫秒级内完成海量商品数据的同步。

一、任务分片:将单体任务拆解为并行子任务

分布式调度的核心思想是“分而治之”。我们将一个庞大的商品更新任务,根据商品ID或其他维度,拆解成多个互不干扰的子任务,然后分发到集群中的不同节点上并行执行。

  1. 定义任务分片上下文

这个类用于在调度时,向每个执行节点传递其需要处理的“分片”信息。

// 包名: juwatech.cn.scheduler.shardingpackagejuwatech.cn.scheduler.sharding;/** * 任务分片上下文,用于定义每个执行节点需要处理的数据范围。 * @author juwatech.cn */publicclassShardingContext{/** * 当前分片的索引,从0开始。 */privateintshardIndex;/** * 总分片数。 */privateintshardTotal;// Getters and SetterspublicintgetShardIndex(){returnshardIndex;}publicvoidsetShardIndex(intshardIndex){this.shardIndex=shardIndex;}publicintgetShardTotal(){returnshardTotal;}publicvoidsetShardTotal(intshardTotal){this.shardTotal=shardTotal;}}
  1. 实现可分片的任务处理器

这是一个抽象类,定义了所有可分片任务必须实现的接口。具体的业务逻辑(如更新商品信息)将在子类中实现。

// 包名: juwatech.cn.scheduler.jobpackagejuwatech.cn.scheduler.job;importjuwatech.cn.scheduler.sharding.ShardingContext;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;/** * 抽象的可分片任务处理器。 * @author juwatech.cn */publicabstractclassAbstractShardingJobHandler{protectedfinalLoggerlogger=LoggerFactory.getLogger(this.getClass());/** * 执行分片任务的核心方法。 * @param context 分片上下文,包含当前节点的索引和总分片数 * @throws Exception 任务执行过程中可能出现的异常 */publicabstractvoidexecute(ShardingContextcontext)throwsException;}

二、任务调度与执行:基于数据库锁的简易调度器

在分布式环境下,必须确保同一个分片任务在同一时间只被一个节点执行,以避免数据冲突和资源浪费。我们可以利用数据库的唯一约束或FOR UPDATE来实现一个简易的分布式锁。

  1. 创建分布式锁表

首先,需要在数据库中创建一个用于协调任务的表。

-- 分布式任务锁表CREATETABLE`distributed_job_lock`(`job_name`varchar(255)NOTNULLCOMMENT'任务名称,主键',`locked_by`varchar(255)DEFAULTNULLCOMMENT'锁定该任务的节点标识(如IP+PID)',`lock_time`datetimeDEFAULTNULLCOMMENT'加锁时间',PRIMARYKEY(`job_name`))ENGINE=InnoDBDEFAULTCHARSET=utf8mb4COMMENT='分布式任务锁表';
  1. 实现分布式锁管理器

这个组件负责尝试获取和释放锁。

// 包名: 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{privatestaticfinalLoggerlogger=LoggerFactory.getLogger(DatabaseLockManager.class);privatefinalJdbcTemplatejdbcTemplate;// 当前应用实例的唯一标识,可以是 IP + 进程IDprivatefinalStringlockOwner;publicDatabaseLockManager(JdbcTemplatejdbcTemplate){this.jdbcTemplate=jdbcTemplate;this.lockOwner=System.getenv("HOSTNAME")+"-"+java.lang.management.ManagementFactory.getRuntimeMXBean().getName();}/** * 尝试获取指定任务的锁。 * @param jobName 任务名称 * @return 获取锁成功返回true,否则返回false */publicbooleantryLock(StringjobName){Stringsql="INSERT INTO distributed_job_lock (job_name, locked_by, lock_time) VALUES (?, ?, ?)";try{introws=jdbcTemplate.update(sql,jobName,lockOwner,LocalDateTime.now());if(rows>0){logger.info("成功获取任务锁: {},持有者: {}",jobName,lockOwner);returntrue;}}catch(DuplicateKeyExceptione){// 插入失败,说明锁已被其他节点持有logger.debug("任务锁已被占用: {}",jobName);}catch(Exceptione){logger.error("获取任务锁时发生未知错误",e);}returnfalse;}/** * 释放指定任务的锁。 * @param jobName 任务名称 */publicvoidreleaseLock(StringjobName){Stringsql="DELETE FROM distributed_job_lock WHERE job_name = ? AND locked_by = ?";introws=jdbcTemplate.update(sql,jobName,lockOwner);if(rows>0){logger.info("成功释放任务锁: {},持有者: {}",jobName,lockOwner);}}}
  1. 实现具体的商品更新任务

这是一个具体的任务实现,它会根据分片信息从数据库中拉取对应的商品进行更新。

// 包名: 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.jdbcTemplate=jdbcTemplate;}@Overridepublicvoidexecute(ShardingContextcontext)throwsException{logger.info("开始执行商品更新任务,分片信息: {}/{}",context.getShardIndex(),context.getShardTotal());// 1. 根据分片信息拉取商品ID列表// 使用 MOD 函数进行分片,确保每个节点处理不同的数据集StringselectSql="SELECT item_id FROM products WHERE MOD(item_id, ?) = ?";List<Long>itemIds=jdbcTemplate.queryForList(selectSql,Long.class,context.getShardTotal(),context.getShardIndex());logger.info("分片 {}/{} 获取到 {} 个商品待更新",context.getShardIndex(),context.getShardTotal(),itemIds.size());// 2. 遍历商品ID,调用电商平台API获取最新信息并更新本地数据库for(LongitemId:itemIds){try{// 模拟调用API获取最新商品信息// Map<String, 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());}}
  1. 调度器入口

最后,需要一个入口来触发整个流程。在实际生产中,这通常由一个轻量级的调度中心(如XXL-JOB, Elastic-Job)来触发,这里为了演示,我们用一个简单的Spring Boot@Scheduled注解来模拟。

// 包名: 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{privatestaticfinalLoggerlogger=LoggerFactory.getLogger(DistributedJobScheduler.class);privatestaticfinalStringJOB_NAME="PRODUCT_UPDATE_JOB";privatefinalDatabaseLockManagerlockManager;privatefinalAbstractShardingJobHandlerjobHandler;publicDistributedJobScheduler(DatabaseLockManagerlockManager,AbstractShardingJobHandlerjobHandler){this.lockManager=lockManager;this.jobHandler=jobHandler;}/** * 每分钟执行一次的任务调度。 * 在实际生产中,这个触发应由专业的调度中心完成。 */@Scheduled(cron="0 */1 * * * ?")publicvoidschedule(){// 1. 尝试获取分布式锁if(!lockManager.tryLock(JOB_NAME)){logger.info("未能获取任务锁,本次调度跳过。");return;}try{// 2. 获取锁成功,开始执行任务// 假设我们有4个应用节点,就将任务分为4片intshardTotal=4;// 在实际的分布式调度框架中,shardIndex是由调度中心分配给当前节点的。// 这里为了演示,我们假设当前节点负责处理第0片。// 在真实场景中,每个节点都会启动这个任务,但调度中心会确保它们拿到不同的shardIndex。intshardIndex=0;// 这个值应该由配置或调度中心动态传入ShardingContextcontext=newShardingContext();context.setShardIndex(shardIndex);context.setShardTotal(shardTotal);jobHandler.execute(context);}catch(Exceptione){logger.error("执行分布式任务时发生异常",e);}finally{// 3. 任务执行完毕,无论成功与否,都必须释放锁lockManager.releaseLock(JOB_NAME);}}}

本文著作权归 省赚客app 研发团队,转载请注明出处!

http://www.cnnetsun.cn/news/3847100.html

相关文章:

  • SpringBoot+Vue3构建古典舞在线交流平台全栈实践
  • 全志芯片开发必备:编译sunxi-tools与FEL模式实战指南
  • CasADi与MPC实现车辆轨迹跟踪控制
  • UE5 Common UI插件重构:构建健壮可维护的菜单系统架构
  • B2B制造企业怎么做GEO优化?产品资料、英文官网和AI搜索可见度对比
  • 信号简介...
  • 魔兽世界血DK莱登极限输出:动态资源管理与高压生存循环详解
  • Windows DOS命令实战:从基础操作到批处理脚本开发
  • 如何在Unity游戏中安装MelonLoader:全球首个双引擎模组加载器终极指南
  • KLayout版图设计终极指南:从零构建专业级IC设计工作流
  • Claude中转站辅助竞品分析:资料归纳、卖点拆解与内容重组
  • 绝区零自动化终极解决方案:OneDragon智能游戏助手的技术深度解析
  • 无标题文档的智能化管理与自动生成技术实践
  • 5步打造私人游戏云:Sunshine自托管游戏串流完整指南
  • OpenClaw.NET工程化实践:基于TokenHub实现LLM自动化流程的成本核算与优化
  • 一套完整监控系统,到底有哪些核心设备?
  • 摄影色差与紫边:成因解析与全流程解决方案
  • 英雄联盟智能辅助工具LeagueAkari:全面提升游戏体验的终极指南
  • Matlab雷达信号仿真:从单频到混合调制波形生成与性能分析
  • 手把手部署MiniMax H3模型:基于vLLM-Omni打造本地OpenAI兼容API
  • 从PAT彩虹瓶问题解析栈数据结构:LIFO原理、应用场景与算法实现
  • 八大网盘直链解析工具:本地化安全下载解决方案
  • 深度学习数据集划分:训练集、验证集、测试集的核心原理与工程实践
  • 3分钟终极指南:如何在Windows上快速安装苹果USB网络共享驱动
  • 终极音乐解锁指南:Unlock Music Electron 桌面版完全解析
  • AI模型安全实践:从开源平台安全披露到本地部署加固指南
  • 42.SAP Open SQL 批量读取、内表排序、ALV 合计与双击事件实战
  • Unity动态音频加载实战:告别Resources文件夹,实现高效资源管理
  • 44.SAP ABAP SELECT-OPTIONS 动态日期默认值设置方法
  • 从“大佬萌茶”事件看开源项目评估:技术营销、代码审查与工程实践