资讯动态

Scrapy分布式爬虫调度器架构设计:队列、去重与锁的工程实践

发布时间:2026/9/15 4:07:41 来源:尧图企业网站定制
我做分布式爬虫这几年感触最深的一点是很多人把分布式爬虫理解成多台机器同时跑Scrapy结果代码一上线就翻车——重复抓取、队列错乱、任务丢失甚至Redis被塞爆直接OOM。问题基本都出在同一个地方调度器没设计好。Scrapy分布式爬虫调度器架构设计本质上是解决三件大事请求怎么排队、去重怎么做、多台机器的任务怎么协调。这篇文章我会把调度器的角色边界、任务模型、核心实现、生产环境的坑以及从伪分布式到真分布式的部署演进一次性说透。适合正在做爬虫平台化、准备从单机爬虫升级到分布式的工程师参考也适合想搞懂scrapy-redis原理、被分布式锁和任务去重折磨的同学。1. 先搞清楚一件事Scrapy原生调度器为什么撑不起分布式1.1 Scrapy单机调度器的记忆全在进程内Scrapy官方框架自带的调度器底层用的是Scheduler类它管理两个队列一个内存队列memory queue一个磁盘队列disk queue。去重用的RFPDupeFilter本质上就是一个Pythonset集合存放请求指纹。这套设计在单机场景下没有问题但它有一个致命特征所有状态都活在当前进程里。当你启动一个Scrapy爬虫请求队列、已爬取指纹、失败重试信息全都绑定在这一个PID上。一旦爬虫停止队列清空一旦换了机器指纹集合不复存在一旦同时启动两个爬虫进程它们各自维护各自的记忆彼此完全不知道对方抓过什么。这意味着什么呢我用一个例子说明。假设你要爬一个电商网站商品详情页URL是/item/{id}.html总共有10万个商品。你在服务器A上启动了爬虫抓到第5万个商品的时候因为内存占用太高崩了。你重新拉起一个爬虫它又会从第1个商品开始爬——因为之前的请求队列和去重指纹已经随着进程退出全部丢了。1.2 多开几个进程/几台机器问题更大有的人说那我多开几个进程不就能加快速度了吗没错Scrapy是支持的但你要面对新的问题多个进程同时从起始URL开始爬同一个URL会被两个进程同时请求settings.py里配置的DUPEFILTER_CLASS默认是RFPDupeFilter它基于内存set去重多进程之间完全不共享目标网站收到的请求量翻倍甚至翻三倍但有效抓取量没有同步提升反而更容易触发反爬。我见过一个真实案例某团队用scrapy crawl spider同时拉起8个进程结果日志里全是重复请求目标网站的反爬策略直接把他们的IP段封了。最后一查8个进程里7个都在重复抓取同一个队列头部的请求。所以结论很直接要搞分布式第一步就是要把调度器的三样核心状态——请求队列、去重指纹、任务状态——从进程内存里搬出去搬到一个所有Worker都能访问的地方。目前最通用、成本最低的方案就是Redis。1.3 调度器架构设计的本质从单机时序到多机协调单机调度器的工作模式是串行化的Scrapy引擎从调度器拿一个请求交给下载器下载完交给爬虫解析产生新请求再放回调度器。整个过程像一条流水线调度器只管一个工人。分布式调度器的工作模式则完全不同它面对的是多个Worker每个Worker是一套完整的Scrapy爬虫实例同时来取任务、同时来交任务。这时候调度器就变成了一个多进多出的协调中心。我习惯用外卖平台的派单系统来类比你把外卖订单URL请求统一放到一个任务池里骑手Worker接单、送完、再接单。调度器要解决的问题就是——这份订单不会被两个骑手同时抢到同一个订单不会被多次派发订单的优先级要合理爆单的时候先处理高优订单。这套逻辑搬到分布式爬虫上就是调度器架构设计的核心。2. 分布式调度器在整体架构里的位置以及它的职责边界2.1 一个典型的Scrapy分布式架构长什么样从架构分层来看一个可落地的Scrapy分布式爬虫系统通常包含以下几层接入层/任务入口接收运营配置的抓取任务比如爬取分类A下的所有商品详情页生成初始请求调度中心也就是本文的核心负责请求的入队、出队、去重、优先级管理、失败重试Worker集群运行Scrapy爬虫的机器/容器从调度中心领取请求执行下载和解析存储层Redis调度器状态 目标数据库MySQL / MongoDB / MinIO等存抓取结果和原始HTML监控告警Redis队列积压量、Worker心跳、抓取速率、异常率等指标的可视化。这里面最容易混淆的地方是很多人以为调度器就是Redis里的一个list。严格来说Redis只是调度器的存储载体调度器的核心是一个逻辑组件它由三部分组成队列调度策略、指纹去重策略、任务状态管理。你可以用scrapy-redis库实现也可以基于Redis自己实现一套后面我会专门讲自研方案。2.2 调度器该管什么不该管什么明确职责边界是架构设计的第一步。调度器管的事情是请求入队爬虫解析出来的新请求统一推到队列请求出队按策略把请求分配给空闲WorkerURL去重判断一个请求是否已经抓过或在队列中优先级调度高优先级请求先出队重试管理抓取失败的请求按规则重新入队流量控制控制全局限速避免被封。调度器不该管的事情这也是很多刚做分布式的人容易犯的错不该管页面解析逻辑那是Spider的职责不该管数据存储逻辑那是Pipeline的职责不该管代理IP的获取和轮换那是Downloader Middleware的职责不该管Worker的弹性伸缩那是部署运维层的职责。一旦你让调度器去操心业务逻辑它就会变得又大又重改一个爬虫业务就要重新部署调度器这在生产环境是不可接受的。我见过有项目直接在调度器里写死了某个网站的解析函数后来网站改版整个调度系统跟着遭殃——这就是职责边界不清的典型反面教材。2.3 同步调用还是异步拉取Worker和调度器的交互模式Worker从调度器拿任务有两种主流交互模式推模式Push调度器主动把请求发给空闲Worker。好处是任务分配及时、均衡坏处是实现复杂调度器需要维护每个Worker的心跳和当前状态还要处理Worker掉线后任务回收的问题。拉模式PullWorker主动向调度器请求任务说我空闲了给我一个请求。好处是实现简单、天然负载均衡谁空闲谁取坏处是极端情况下会出现任务倾斜个别Worker一直抢到高优任务。我在实际项目中几乎都采用拉模式原因很简单爬虫场景下Worker的数量不稳定扩缩容是常态拉模式天然适配Worker自己决定是否要任务的场景调度器只需要保证同一个请求不会被两个Worker同时取走即可而这件事用Redis的原子操作就能解决——这就是后面要讲的分布式锁的用武之地。3. 面向项目落地的调度模型设计任务、队列、去重、优先级3.1 请求怎么序列化才能跨进程传输Scrapy的Request对象是Python内存对象包含URL、Method、Headers、Body、Meta、Callback等信息。跨进程传输之前必须先序列化成字符串。我建议的序列化格式是JSON因为可读性好、调试方便而且Redis本身对字符串最友好。一个序列化后的请求长这样{ url: https://example.com/item/12345.html, method: GET, headers: { User-Agent: Mozilla/5.0 (compatible; MySpider/1.0), Referer: https://example.com/list/2.html }, meta: { category_id: 205, retry_times: 0 }, callback: parse_item, priority: 5 }这里有一个容易被忽略的点回调函数信息也要一起序列化。因为Worker进程里可能有多个回调函数请求取出来之后下载完页面要调用哪个回调必须靠这个字段告诉爬虫。还有meta字段建议只放和业务判断强相关的数据不要塞一些很大的对象进去。我见过有人把整个HTML塞到meta里再入队结果Redis里存的全是几MB的大字符串性能直线下降。正确的做法是meta里只放ID、类型、页码这类小字段解析后再从数据库或对象存储里补全数据。3.2 队列设计不止一个List那么简单很多人用scrapy-redis时以为分布式调度器就是往Redis的List里rpush请求Worker用blpop取请求。这在demo阶段没问题但到了生产环境你就需要考虑多种队列并存待爬队列Ready Queue存放等待被Worker领取的请求。根据优先级可以进一步拆分为高、中、低三个List延迟队列Delay Queue存放需要延时重试的请求比如目标网站返回了限流要求1分钟后重试。用Redis的zset按时间戳排序Worker定时扫描到期的请求把它们搬运回待爬队列正在处理队列Processing QueueWorker取走请求后在抓取完成之前请求状态是进行中。这个状态需要单独记录否则Worker中途崩溃这个请求就永远消失了。这里我特别强调一下正在处理队列的必要性。用blpop取出请求后如果Worker进程崩溃这个请求是拿不回来重新入队的。生产级方案是给每个请求增加一个租约机制Worker取走请求时记录一个时间戳请求在时间戳N秒内未完成就允许重新入队给其他Worker。这个N一般是单个请求平均耗时的3到5倍。3.3 URL指纹去重set够用但要注意内存膨胀去重是调度器的核心能力之一。Scrapy默认的去重指纹计算逻辑是对URL做sha1加密生成40位十六进制字符串。在Redis里去重集合用的就是SADD命令。import hashlib def url_fingerprint(url: str) - str: 生成URL指纹忽略URL中的常见追踪参数 parsed urlparse(url) query parse_qs(parsed.query) # 去掉常见统计参数 for key in [utm_source, utm_medium, utm_campaign, spm, from]: query.pop(key, None) normalized_query urlencode(query, doseqTrue) normalized_url f{parsed.scheme}://{parsed.netloc}{parsed.path} if normalized_query: normalized_url ? normalized_query return hashlib.sha1(normalized_url.encode()).hexdigest()使用Redisset去重的优点是精确、实现简单、支持SISMEMBER原子判断。但缺点在数据量大时非常明显1亿条URL指纹占用Redis内存大约5-6GB而且全是纯内存存储成本很高。如果去重数据量上亿建议改用Bloom Filter它用1%左右的误判率换来90%以上的内存节省。具体做法是引入pybloom_live或Redis的bf模块Redis 4.0以上支持RedisBloom插件把SISMEMBER换成BF.EXISTS和BF.ADD。注意Bloom Filter不能删除元素所以它适合只增不减的去重场景——爬虫去重恰好就是这种场景。还有一个风险要提醒不要把已抓取和待抓取混在一个集合里。严谨的做法是拆成两个集合一个是seen:done已经抓取成功的一个是seen:queue已经入队还没抓的。否则你的爬虫一旦中途调整了队列重新入队的历史请求可能全部被去重挡掉。3.4 优先级策略高优请求永远先被处理分布式爬虫场景下不同请求的紧急程度是不同的。比如你有一个监控爬虫要每5分钟检查一次某个商品是否有库存变化同时你还有一个全量历史数据爬虫要爬10万页商品历史价格。如果这两种任务放在同一个无优先级队列里高优的库存检查请求可能要排队几十分钟才能被执行这在业务上是不可接受的。我的做法是使用多级队列配合加权轮询出队策略高优队列q:high库存检查、价格变更监听等时效性强的请求中优队列q:medium正常业务请求低优队列q:low历史数据补全、大量翻页请求。Worker在出队时按照高:中:低 5:3:2的比例轮询三个队列在保证高优请求及时处理的同时低优任务也不至于被饿死。如果你用的是纯Redisblpop要注意一点blpop只支持单个key的阻塞弹出所以你需要循环三次分别尝试非阻塞的lpop。为了减少Redis往返次数可以试试用Lua脚本一次性完成加权出队的原子逻辑。4. 调度器核心实现Redis队列 指纹去重 分布式锁4.1 自研还是用scrapy-redis做分布式调度器首先要决定是直接用scrapy-redis还是自研。scrapy-redis是社区最常用的方案它把Scrapy的队列、去重、调度器全部用Redis重写了一遍。使用门槛低配置改几行就能跑。但scrapy-redis有几个明显的局限维护频率较低很多年没有大更新不支持延迟队列重试策略比较简单队列默认就是list优先级支持不够灵活没有成熟的租约/心跳机制Worker崩溃容易丢任务准确说是无法自动恢复。我的建议是中小项目、快速验证阶段直接用scrapy-redis稳定可靠、文档也多核心业务、大规模生产环境在此基础上做二次开发或者干脆基于Redis自研一个轻量调度器组件。自研的代码量并不大核心逻辑300行以内就能实现但可控性会大幅提升。4.2 核心数据结构的落地设计我自研的调度器在Redis里的数据结构是这样的键名按项目名加前缀避免和其他应用冲突数据结构键名示例用途Listspider:req:high/req:mid/req:low待爬请求队列按优先级拆分ZSetspider:req:delay延迟队列member为序列化请求score为可执行时间戳Setspider:seen:queue已入队但未抓取的请求指纹Setspider:seen:done已抓取成功的请求指纹Hashspider:inflight正在处理的请求field为指纹value为Worker ID和取走时间Stringspider:stats:qps/spider:stats:total调度器运行状态指标下面给出一个精简版调度器的核心代码重点是入队和出队两个函数import json import time import hashlib import redis class RedisScheduler: def __init__(self, redis_client): self.r redis_client self.prefix spider self.lock_key f{self.prefix}:lock def _req_key(self, priority): level high if priority 5 else (mid if priority 1 else low) return f{self.prefix}:req:{level} def _fingerprint(self, request_dict): raw f{request_dict[method]}|{request_dict[url]}|{hashlib.md5(request_dict.get(callback,).encode()).hexdigest()} return hashlib.sha1(raw.encode()).hexdigest() def push(self, request_dict): 请求入队先判断去重再压入对应优先级队列 fp self._fingerprint(request_dict) # 只在两个集合都不存在时才入队避免重复 if not self.r.sismember(f{self.prefix}:seen:queue, fp) and \ not self.r.sismember(f{self.prefix}:seen:done, fp): pipe self.r.pipeline() pipe.sadd(f{self.prefix}:seen:queue, fp) pipe.rpush(self._req_key(request_dict.get(priority, 1)), json.dumps(request_dict, ensure_asciiFalse)) pipe.execute() return True return False def pop(self, worker_id, timeout0): Worker取请求带租约机制防崩溃丢任务 # 先处理到期的延迟队列 self._promote_delay_queue() # 加权轮询三档队列 for q_key in [f{self.prefix}:req:high, f{self.prefix}:req:mid, f{self.prefix}:req:low]: item self.r.lpop(q_key) if item: request_dict json.loads(item) fp self._fingerprint(request_dict) # 记录租约谁、什么时候取走的 self.r.hset(f{self.prefix}:inflight, fp, json.dumps({worker: worker_id, ts: time.time()})) return request_dict if timeout 0: time.sleep(timeout) return None def ack(self, request_dict): 请求处理完成从inflight移除加入已抓取集合 fp self._fingerprint(request_dict) pipe self.r.pipeline() pipe.hdel(f{self.prefix}:inflight, fp) pipe.srem(f{self.prefix}:seen:queue, fp) pipe.sadd(f{self.prefix}:seen:done, fp) pipe.execute()4.3 分布式锁在这里的两种用法分布式锁在调度架构里非常关键它的使用场景有两个场景一防止多个Worker同时取同一个请求。lpop是原子操作本身就能保证请求不会被两个Worker同时取走所以这一层的锁其实已经由Redis的list原子性解决了。场景二全局任务管理动作的互斥。比如你有一个定时任务调度器每10分钟要扫描一遍延迟队列把到期的任务回收到待爬队列。如果有两个Worker同时执行这个扫描动作就会产生重复回收把同一个请求入队两遍。这时候就需要一把分布式锁保证同一时间只有一个Worker在执行延迟队列回收这个管理动作。分布式锁最简单的实现是Redis的SET key value NX EX加上一个随机值防止误删别人的锁def acquire_lock(self, lock_key, token, expire10): return self.r.set(lock_key, token, nxTrue, exexpire) def release_lock(self, lock_key, token): # Lua脚本保证原子性只有在value匹配时才删除 lua if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end return self.r.eval(lua, 1, lock_key, token)这里有几个经验教训锁的过期时间不能太短否则任务没执行完锁就释放了别的Worker进来抢锁造成并发问题也不能太长否则持有锁的Worker崩溃了锁要等很久才能释放。我一般根据管理操作的平均耗时设为3-5倍时长。4.4 Worker侧的集成方式调度器本身实现好了怎么跟Scrapy的Worker集成最优雅的方式是自定义Scrapy的Scheduler和DupeFilter。# settings.py SCHEDULER myproject.schedulers.RedisScheduler DUPEFILTER_CLASS myproject.dupefilters.RedisDupeFilter REDIS_URL redis://redis-host:6379/0Scrapy的Scheduler接口要求实现enqueue_request入队、next_request出队、has_pending_requests是否还有待处理请求等方法。你在自定义Scheduler类里把这三个方法映射到上面RedisScheduler的push、pop、has_pending即可。注意has_pending_requests这个方法的实现。Scrapy引擎在判断爬虫是否结束时会调用它。分布式场景下一个Worker判断自己没活儿干了不一定代表整个任务已经结束——可能别的Worker还在跑。所以这里不能简单地判断我的本地队列是否为空而是要判断Redis里所有待爬队列、延迟队列、正在处理队列是否都为空。否则爬虫会提前退出导致任务还没跑完就停了。5. 生产环境踩过的坑调度器假死、重复消费、Redis内存爆炸5.1 爬虫任务跑完但还有一堆URL没抓这是我踩过最深的一个坑。单机Scrapy爬虫跑完后会正常退出但分布式环境下某个Worker负责的某个分类页面解析出来的新URL全部入队到了Redis而它自己本地已经没有任务了。结果这个Worker判断没任务了直接退出剩下Redis里的任务没人处理。排查思路是这样的先看Redis里各队列的llen如果大于0说明还有待爬请求再看Worker数量如果所有Worker都因为无任务退出了Redis队列就陷入了僵死检查has_pending_requests的实现发现它只判断了本地内存队列没有统计Redis。解决办法在has_pending_requests里统一判断Redis各队列的长度之和并配备一个看门狗进程定期检查有队列长度0但活跃Worker数0的情况触发新的Worker拉起。5.2 多个Worker同时执行延迟队列回收请求翻倍延迟队列回收用的是多Worker竞争执行如果没有分布式锁保护两个Worker同时扫描到同一个到期的请求同时把它zrem并rpush到待爬队列这个请求就会被抓两遍。排查时我发现日志里同一个URL在几秒内被两个不同IP的Worker各请求了一次目标网站的详情页出现了两条重复数据。解决办法就是我上面讲的用Redis分布式锁把延迟队列回收操作变成全局互斥。另外在入队时再检查一次去重指纹多一层保险。5.3 Redis内存被去重set撑爆有一个项目跑了大半年Redis内存从2GB涨到40GB最后直接OOM所有Worker全部瘫痪。排查出来的元凶是seen:done集合无限增长——这台爬虫抓了将近2亿个URL每个指纹40个字符光这个集合就占了接近10GB内存而且完全不清理。解决办法分三步给seen:done设置定期清理策略超过30天的指纹用SSCAN分批删除对全量历史数据的抓取去重集合改用Bloom Filter内存降为原来的十分之一将Redis启动参数maxmemory设置为物理内存的70%避免Redis OOM影响整台服务器。5.4 Worker假死导致任务无限期卡在inflight某个Worker在取走请求之后因为网络抖动、内存溢出等原因假死请求一直停留在inflight里租约过期了也没人管后续请求全部卡住。我当时的排查过程是队列里明明有几千个请求但抓取速率降为0。查看inflight发现里面躺着十来个请求全是同一个Worker ID时间戳停留在2小时之前。解决办法是增加一个租约回收定时任务每隔1分钟扫描inflight找到ts超过5分钟的记录把对应请求重新放回待爬队列并删除inflight记录。这个回收任务也要用分布式锁保护防止多Worker同时回收同一个过期请求。5.5 动态页面和iframe导致的假请求无效最后说一个业务层面的坑。很多网站在详情页里嵌了iframe数据是通过异步接口加载的。如果你只是把https://example.com/item/123这个URL入队Worker抓下来发现页面里全是iframe框架解析不出任何数据。这种情况不能靠调度器硬扛需要配合渲染层。我现在的做法是在Downloader Middleware里集成scrapy-playwright对指定域名开启浏览器渲染。但注意带渲染的请求耗时会比普通请求高出5到10倍所以这类请求要单独设置一个低的优先级队列用独立的Worker去消费避免拖慢普通请求的处理速度。另外对于动态URL指纹去重不能只看URL要结合meta中的页面版本号或数据时间戳来生成指纹否则网站更新了数据你的爬虫因为指纹相同直接跳过抓回来的全是旧数据。6. 从伪分布式到真分布式部署演进与监控保障6.1 伪分布式一台机器上先跑通全部逻辑很多人在本地开发的时候习惯用PyCharm直接点运行调Scrapy爬虫。但在分布式调试场景下我建议先在单台机器上伪装出多个Worker来验证调度逻辑。具体操作是在命令行里同时启动多个爬虫进程每个进程绑定不同的SPIDER_NAME和日志文件# 终端1启动Worker A nohup scrapy crawl product_spider --logfileworker_a.log /dev/null 21 # 终端2启动Worker B nohup scrapy crawl product_spider --logfileworker_b.log /dev/null 21 # 终端3启动Worker C nohup scrapy crawl product_spider --logfileworker_c.log /dev/null 21 然后观察Redis里的队列变化情况、三个进程的日志是否出现重复抓取、inflight里是否有超过预期时间的残留。这一步能把调度器80%的逻辑bug暴露出来而不用急着上多台服务器。之所以建议先跑伪分布式是因为分布式系统的bug排查成本非常高。两台机器之间出问题你要同时看两边的日志还要排查网络、Redis连接等外部因素。而在同一台机器上模拟分布式起码可以先把调度器自身的逻辑问题排查干净。6.2 真分布式部署容器化 统一代理出口伪分布式验证通过后就可以上容器了。我现在的标准部署方式是每个Worker是一个Docker容器通过Docker Compose或Kubernetes编排统一从镜像仓库拉取爬虫镜像。这里有个部署细节容易被忽略代理IP的出口问题。多台Worker如果各自使用本机公网IP去请求目标网站很容易形成多个不同IP访问同一个站点的情况反爬系统非常容易识别出这是分布式采集。我的做法是所有Worker的出口流量统一走代理网关代理池管理工具负责维护IP池、检测IP失效、按权重分配IP。Worker只需要从代理网关获取一个可用代理不需要关心IP来自哪里。这样能极大降低被封风险。6.3 监控指标体系怎么判断调度器健不健康调度器跑得好不好不能靠感觉。我整理的监控指标分三层队列层各优先级队列的llen反映任务积压情况持续上涨说明Worker消费能力不足zset延迟队列的长度值过大说明很多请求在重试可能代理IP质量差或网站反爬严格inflight数量长时间大于0且不下降说明有Worker假死。Worker层每个Worker最近一次心跳时间超过1分钟没有心跳就告警每个Worker的抓取QPS和成功率某个Worker成功率骤降大概率是代理IP失效。结果层数据入库延迟从请求入队到结果入库的平均时间重复数据率目标库里同一业务主键出现多次的比例。这些指标我会全部推到Prometheus再通过Grafana做可视化。告警规则就三条最关键的队列积压超过阈值、Worker心跳丢失、重复数据率超过5%每一条都对应一个明确的运维响应动作而不是盲目拉群。最后说点实在的调度器做得好不好直接决定了整个分布式爬虫系统的稳定上限。我见过太多项目一上来就堆机器、堆协程、堆代理结果Redis先崩了也见过不少团队花大量时间优化解析逻辑却忽视了调度器里的一个低级bug导致数据重复率居高不下。我个人在实际操作中的体会是调度器设计的核心矛盾不是多快而是多稳。你宁可每秒只抓50个请求也不要让整个任务在跑了3个小时之后因为一个租约问题前功尽弃。设计的时候多花点时间想清楚队列、去重、锁、租约这四件事上线之后会省下无数个熬夜排查的夜晚。还有一个经验送给正在起步的人先用scrapy-redis跑通最小模型再按本文的思路逐步把延迟队列、租约机制、分布式锁加进去不要一上来就闭门造车写一套类似Kafka的重型调度系统——那不是架构设计那是性格测试。

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

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

免费获取报价