资讯动态

大规模请求任务的工程化架构:调度、去重、存储、重试怎么搭

发布时间:2026/10/8 17:28:08 来源:尧图企业网站定制
大规模请求任务的工程化架构调度、去重、存储、重试怎么搭几十个请求一个 for 循环就够了。但任务量上到几万、几十万你很快会发现脚本越跑越慢、内存越吃越多、重跑一遍还得从头来、失败的任务悄无声息地丢了。这时候要的已经不是更快的循环而是一套最小的工程化架构。这篇把它拆成四块调度、去重、存储、重试——一块一块讲清楚该怎么做、用什么。一、先看一个 for 循环会怎么死假设你要处理 10 万个待办任务importrequestsforurlinurls:# urls 有 10 万个requests.get(url)跑起来之后你会依次撞上这些墙太慢串行一个接一个大部分时间在等网络重跑等于白跑中途挂了下次从头开始前面的全白干做完的又做一遍源数据里有重复同一个任务被反复处理内存涨到爆把结果全塞进一个 list几十万条就撑不住了失败无声无息超时、报错的任务直接跳过没人知道丢了多少。这四个问题刚好对应要补的四块调度、去重、存储、重试。二、调度让任务排队 并发调度的核心只有两件事并发度和节奏。并发度决定同时有多少任务在跑节奏决定单位时间内发出去多少。前者用线程池 / 协程池解决后者用令牌桶或漏桶限速。一个最小骨架importqueue,threading qqueue.Queue()fortaskintasks:q.put(task)# 先把任务全部入队defworker():whileTrue:try:taskq.get_nowait()exceptqueue.Empty:returnhandle(task)# 处理单个任务q.task_done()threads[threading.Thread(targetworker)for_inrange(20)]fortinthreads:t.start()fortinthreads:t.join()要点任务队列和消费解耦。生产任务的速度和处理速度可以不一样队列起到缓冲作用。并发度要压测出来不是拍脑袋。这一点在讲并发的文章里说过逐步加压看错误率。加个限速。很多服务端有速率限制你打太快只会换来一堆 429。用一个简单的 sleep 或令牌桶控制 QPS。顺带说一句如果每个任务都新建一次连接光 TLS 握手就够呛。用requests.Session()或连接池把 TCP 连接复用起来往往是单个请求提速最立竿见影的一招——握手和慢启动都省了。这个点在讲 Session 复用的文章里展开过这里不重复。三、去重别让同一个任务跑第二遍去重看着简单坑却最多。内存 set最直接但任务量大就吃内存进程一重启全丢布隆过滤器极省内存代价是有极小的误判率会把没做的判成做过了Redis / 数据库可持久化、多进程共享适合要断点续跑的场合。选哪个取决于你的规模和对漏判的容忍度。完整对比我放在同期的配菜里这里先记住一句话小批量用 set大批量上布隆要断点续跑用 Redis。四、存储落盘而且要幂等处理完的结果必须及时落盘而不是攒在内存里最后一起写。原因很简单中途挂了内存里的全没了。三种常见落点落点适合注意CSV / JSONL 文件中小批量、要人工看追加写别整份重写SQLite / 本地数据库单机、要查询开 WAL 模式对象存储 / 服务端数据库大规模、要共享批量写别一条一提交还有个小细节别一条一条地写库。攒够一批比如 500 条再批量提交写入速度能差出好几倍。落文件也一样——用追加模式一行一行写比最后把整个列表 dump 一次更抗中断。一个必须做到的点是幂等同一个任务重复写入结果应该一致而不是产生两条。最常用的做法是给每条记录一个唯一键比如任务参数的 hash写入时用INSERT OR REPLACE或者去重约束兜底。五、重试让失败的任务有第二次机会没有重试的批量任务等于成功靠运气。重试要讲策略不能死循环重试指数退避第 1 次等 1 秒第 2 次 2 秒第 3 次 4 秒……给服务端喘息的空间重试上限一般 3 到 5 次超过就放弃区分错误超时、429 值得重试404、参数错误重试多少次都没用死信队列重试到上限仍失败的任务单独记下来别静默丢弃。importtimedefwith_retry(fn,tries4):foriinrange(tries):try:returnfn()exceptException:ifitries-1:raisetime.sleep(2**i)# 1, 2, 4 秒还有一点容易忽略重试和去重是一对搭档。如果去重记录的是“处理成功”那失败重试的任务就不能被误标成已处理反过来如果去重记录的是“已经排过队”重试又会被它挡住。工程上通常把“已入队”和“已成功”分成两种状态记录别混用一个集合。六、把四块拼起来源数据 → [去重] → 任务队列 → [调度/限速] → Worker → [存储(幂等)] ↑ ↓ └──── [重试/死信] ←─────┘顺序上有个关键点去重放在入队之前。任务进队列前先过滤掉做过的能省掉大量无谓的调度和网络开销。放反了去重就只是事后补救。另外整套流程最好能断点续跑。把这四块的状态哪些做完了、哪些还在队列里尽量外置到磁盘或 Redis脚本被 kill 之后重启才能接着上次的位置继续而不是从零开始。七、别忘了一件事可观测架构搭好之后你还需要知道它到底跑得怎么样。至少要盯这几个数队列长度一直涨说明消费跟不上生产成功率 / 失败数失败率突然升高往往是服务端限流或网络出了问题QPS便于和服务端的限制保持一致耗时分布只看平均值会骗人看 P95 / P99。这几个指标不复杂但有没有差别很大——出问题时前者是我看一眼就知道后者是我只能猜。八、小结把一句话总结成表模块解决什么常用做法调度慢、并发失控队列 线程/协程池 限速去重重复处理set / 布隆 / Redis存储结果丢失追加落盘 幂等键重试静默失败指数退避 上限 死信不用一上来就上 Kafka、上分布式——先把这四块补齐一个单机脚本就能稳稳跑完几十万任务。架构不是越复杂越好而是该有的地方不能缺。互动时间你的批量任务脚本踩过最大的坑是哪个是重复跑、内存爆还是失败没重试评论区聊聊我挑典型的写成下一篇 声明本文为原创技术分享涉及的采集思路仅供学习与合法用途参考。请遵守目标网站的 robots 协议及相关法律法规勿将技术用于任何违规场景。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑