分布式任务调度在自动化购物系统中的应用:从定时任务到弹性调度的演进

简介: 本文介绍了一种基于Redis ZSET的分布式任务调度方案,用于高效支撑数千个煤炉/雅虎代拍监控任务。相比原始多线程轮询和APScheduler方案,该架构通过集中式任务队列、Worker竞争消费与弹性扩缩容,将延迟压至200ms内,显著提升资源利用率与系统稳定性。(239字)

一、业务场景
在煤炉自动下单和雅虎代拍场景中,系统需要同时监控成千上万个用户设置的关键词和商品。每个监控任务都需要定期执行——搜索新上架的商品、检查价格变化、触发自动下单。
Bidfins最早实现的是一个简单的while True循环 + time.sleep()方案:
python
import timedef monitor_task(keyword, max_price): while True: items = search_mercari(keyword, max_price) for item in items: if match_condition(item): place_order(item) time.sleep(5) # 每5秒检查一次
当监控任务从10个增长到1000个的时候,这个方案彻底崩溃了——1000个线程同时运行,每个线程5秒一次请求,每秒就是200次请求,服务器CPU直接飙到100%。
二、APScheduler方案
我首先用APScheduler重构了任务调度:
python
from apscheduler.schedulers.background import BackgroundSchedulerfrom apscheduler.triggers.interval import IntervalTriggerscheduler = BackgroundScheduler()def add_monitor_task(keyword, max_price, user_id): scheduler.add_job( func=monitor_once, trigger=IntervalTrigger(seconds=5), args=[keyword, max_price, userid], id=f"monitor{userid}{keyword}", replace_existing=True )
APScheduler比手动循环好一些,但本质上还是每个任务独立调度。当任务数量达到几千个时,调度器本身的开销就很大了。
三、基于Redis ZSET的集中式任务调度
最终方案是把所有监控任务放在Redis的有序集合(ZSET)中,由一组Worker竞争消费:
python
import redisimport timeimport jsonclass TaskScheduler: def init(self): self.redis = redis.Redis(decode_responses=True) self.task_queue = "scheduler:tasks" self.running_key = "scheduler:running" def add_task(self, task_id, task_data, interval=5): """添加一个周期性任务""" # 计算下一次执行时间 next_run = time.time() + interval # 存储任务详情 detail_key = f"task:detail:{task_id}" self.redis.hset(detail_key, mapping={ 'data': json.dumps(task_data), 'interval': interval, 'last_run': '0' }) # 加入ZSET,score为下一次执行时间 self.redis.zadd(self.task_queue, {task_id: next_run}) def acquire_task(self, worker_id): """Worker获取一个待执行的任务""" # 获取score最小的任务(最早需要执行的) tasks = self.redis.zrange(self.task_queue, 0, 0, withscores=True) if not tasks: return None task_id, score = tasks[0] now = time.time() if score > now: # 还没到执行时间 return None # 尝试锁定这个任务(使用Lua脚本保证原子性) lock_script = """ local task_id = ARGV[1] local worker_id = ARGV[2] local now = tonumber(ARGV[3]) -- 从ZSET中移除 local removed = redis.call('ZREM', KEYS[1], task_id) if removed == 0 then return nil end -- 记录正在运行 redis.call('HSET', KEYS[2], task_id, worker_id) return 'ok' """ result = self.redis.eval( lock_script, 2, self.task_queue, self.running_key, task_id, worker_id, now ) if result is None: return None # 获取任务详情 detail_key = f"task:detail:{task_id}" task_data = self.redis.hgetall(detail_key) return { 'task_id': task_id, 'data': json.loads(task_data['data']), 'interval': int(task_data['interval']) } def complete_task(self, task_id): """任务执行完成,重新调度""" detail_key = f"task:detail:{task_id}" interval = int(self.redis.hget(detail_key, 'interval')) next_run = time.time() + interval # 更新最后执行时间 self.redis.hset(detail_key, 'last_run', str(time.time())) # 重新加入ZSET self.redis.zadd(self.task_queue, {task_id: next_run}) # 从运行中移除 self.redis.hdel(self.running_key, task_id)
Worker的实现:
python
import threadingimport timeclass Worker: def init(self, worker_id, scheduler): self.worker_id = worker_id self.scheduler = scheduler self.running = True def run(self): while self.running: task = self.scheduler.acquire_task(self.worker_id) if task: self._execute(task) else: time.sleep(0.1) # 没有任务时短暂休眠 def _execute(self, task): try: # 执行监控逻辑 items = search_mercari( task['data']['keyword'], task['data']['max_price'] ) for item in items: if match_condition(item, task['data']): place_order(item, task['data']['user_id']) except Exception as e: # 记录错误,但不影响任务重新调度 log_error(task['task_id'], str(e)) finally: self.scheduler.complete_task(task['task_id'])
四、弹性扩缩容
这套架构支持动态增减Worker——流量大的时候启动更多Worker,流量小的时候减少Worker:
python
class WorkerManager: def init(self, scheduler): self.scheduler = scheduler self.workers = [] self.min_workers = 2 self.max_workers = 20 def adjust_workers(self): # 根据队列长度调整Worker数量 queue_length = self.scheduler.redis.zcard(self.scheduler.task_queue) target = max(self.min_workers, min(self.max_workers, queuelength // 50 + 1)) while len(self.workers) < target: worker = Worker(f"worker{len(self.workers)}", self.scheduler) threading.Thread(target=worker.run, daemon=True).start() self.workers.append(worker)
五、总结
这套基于Redis ZSET的分布式任务调度方案,支撑了数千个煤炉监控任务的同时运行,任务延迟控制在200毫秒以内。核心经验是:不要把调度逻辑分散在各个任务中,而是集中管理、统一调度。

目录
相关文章
|
3月前
|
SQL 关系型数据库 Java
跨境电商独立站多租户架构设计:从零搭建SaaS平台
本文从架构师视角对比多租户三大隔离方案(独立库/Schema/共享表),结合Taoify跨境电商实践,详解基于Spring Boot + MyBatis Plus的租户上下文传递、SQL自动注入与动态数据源切换实现,并分享阿里云RDS部署最佳实践。(239字)
381 1
|
3月前
|
人工智能 Rust 监控
这 3 个开源小工具,帮你让 Coding Agent 少吃点 Token
今天我们就来分享 3 个有用的开源项目,专门帮你的 Coding Agent 整理“上下文”:让它少翻无关代码,少吞冗长日志,把 token 留给更关键的信息。
585 0
这 3 个开源小工具,帮你让 Coding Agent 少吃点 Token
|
负载均衡 网络虚拟化
网络技术基础(17)——以太网链路聚合
【3月更文挑战第4天】网络基础笔记(加班了几天,中途耽搁了,预计推迟6天)
|
存储 缓存 Android开发
android分区概述
android分区概述
1409 0
|
2月前
|
人工智能 NoSQL 测试技术
测试开发必备的 AI 技能库:推荐6 个让接口自动化测试效率翻倍的 Skills
本文详解接口自动化测试“执行与报告”阶段的6款AI Agent Skill:打标、执行、诊断、清理、报告、编排,覆盖脚本筛选→智能运行→失败自愈→数据净化→决策报告→全链路调度,实现真正智能化、流水线化的测试闭环。
424 1
测试开发必备的 AI 技能库:推荐6 个让接口自动化测试效率翻倍的 Skills
|
2月前
|
存储 人工智能 运维
让 Agent 越用越准、成本越来越低:AgentLoop 的 Agent 经验自进化闭环
企业 Agent 上线后,如何在不重新训练模型的情况下提升准确率、稳定性并优化单位成功成本?本文介绍 AgentLoop 如何基于真实运行轨迹自动生成可复用经验,通过 Skill 和 CLI 按需召回,并给出控制台接入步骤及与 Memory、RAG、微调、RL 的差异。
562 1
|
2月前
|
存储 前端开发 NoSQL
高并发场景下的订单幂等性设计:从重复支付到精准扣款的实战
跨境电商订单系统需解决重复提交导致的重复下单问题。本文详解幂等性原理与实践:通过唯一订单号+数据库索引、Redis幂等令牌(含Lua原子校验)、订单号复用及多服务协同等方案,兼顾可靠性与性能,核心在于唯一标识、状态记录与原子操作三要素。(239字)
223 0
|
2月前
|
云安全 人工智能 安全
AI Agent 同时拿着浏览器、文件系统和命令行——三层防线为什么一起失效了
本文揭示AI编程代理三大致命漏洞:私有仓库因Issue评论被掏空、大模型绕过安全策略100%生成恶意代码、Agent自主窃取凭证摧毁生产环境。根源在于权限与判断脱节,上下文即攻击面。沙箱与规则难防“合法但有害”的决策,亟需上下文感知的动态策略网关。
297 0
|
2月前
|
前端开发 Java API
[049][Crypto模块]前后端混合加密API实战:基于Spring Boot的AES+RSA安全传输方案
本文详解Spring Boot中AES+RSA混合加密实战:前端用RSA公钥加密随机AES密钥并传输,后端通过`@Crypto`注解自动解密请求体。涵盖公钥分发、Hex/Base64编码统一、ECB模式适配及Caffeine缓存优化,提供开箱即用的端到端安全传输方案。(239字)
303 0
|
2月前
|
Java 应用服务中间件 数据库连接
Spring Boot这5个配置项默认值就是错的,高并发第1天就502
Spring Boot上线前5项关键配置必查:Tomcat线程数(建议≥500)、HikariCP连接池(单实例≤30)、Actuator端点最小暴露(禁用env/heapdump)、生产日志分级(业务INFO+第三方WARN)、强制激活prod profile。缺一即可能导致502、超时或安全风险。

热门文章

最新文章