ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Scrapy-Redis构建工业级分布式爬虫实战

2026/8/3 16:10:47 拓冰建站 浏览量
Scrapy-Redis构建工业级分布式爬虫实战 1. 项目概述在工业供应链领域数据采集一直是个既关键又头疼的问题。去年我们团队接手了一个工业供应商数据采集项目需要从200多个行业门户网站抓取设备参数、企业资质、产品规格等结构化数据最终目标是在3个月内完成千万级数据的采集清洗。传统单机爬虫在面对这种量级时不仅效率低下而且一旦中断就得从头再来。经过多轮技术选型我们最终采用Scrapy-Redis构建分布式爬虫系统实现了日均50万条稳定采集完整代码已封装成可复用的组件库。这个方案最让我满意的不是性能提升虽然确实很显著而是系统展现出的工程化特性任何节点崩溃都不会影响整体任务新增机器能线性提升吞吐量去重精度达到99.99%这些特性让后期运维成本降低了70%。下面我就拆解这套架构的核心设计包含你一定能用上的实战技巧。2. 核心架构设计2.1 为什么选择Scrapy-Redis在分布式爬虫领域技术选型直接决定后期开发维护成本。我们对比了三种主流方案方案开发成本扩展性断点续传去重机制纯Scrapy集群高差需自定义内存去重ScrapyRabbitMQ中优支持需外接Scrapy-Redis低优原生支持布隆过滤器选择Scrapy-Redis的核心原因是其零改造特性——在保留Scrapy所有优点的前提下仅通过替换调度器组件就实现了分布式特性。具体来说调度器扩展使用RedisSpider替代原生调度器所有爬虫共享同一个Redis队列去重优化原生使用Python的set结构我们升级为Redis的Bloom Filter状态同步通过Redis的pub/sub机制实现节点间心跳检测关键技巧在settings.py中启用以下配置这是大多数教程不会告诉你的优化项SCHEDULER_PERSIST True # 保持任务队列不丢失 SCHEDULER_IDLE_BEFORE_CLOSE 10 # 空队列等待时长(秒) REDIS_START_URLS_AS_SET True # 使用集合存储初始URL2.2 断点续传实现细节工业数据采集最怕的就是中途崩溃。我们的方案在以下三个层面确保任务可恢复1. 请求指纹持久化# 自定义去重过滤器 class BloomDupeFilter(RFPDupeFilter): def __init__(self, server, key): self.bf BloomFilter(server, key, 10000000, 0.01) def request_seen(self, request): fp request_fingerprint(request) if self.bf.exists(fp): return True self.bf.insert(fp) return False2. 爬虫状态快照每小时将以下数据持久化到Redis已爬取URL计数当前深度优先搜索路径异常重试次数统计3. 分级恢复策略根据中断原因自动选择恢复点网络中断从最后成功请求继续解析失败回退到上一个URL层级反爬触发切换备用User-Agent池3. 千万级数据处理实战3.1 负载均衡方案当20个爬虫节点同时工作时传统的轮询调度会导致某些节点饿死。我们的解决方案是动态权重分配节点性能画像# 在爬虫启动时注册节点信息 redis_client.hset(node_status, fnode_{uuid}, json.dumps({ cpu: psutil.cpu_percent(), memory: psutil.virtual_memory().percent, bandwidth: speedtest().download }))任务分配算法def get_task_weight(): node_info get_current_node_status() # 计算综合负载系数 load_factor (node_info[cpu]*0.4 node_info[memory]*0.3 (100-node_info[bandwidth])*0.3) return max(1, int(100 - load_factor))动态调整机制每5分钟更新一次节点权重高负载节点自动减少任务获取频率新增节点自动加入调度池3.2 数据去重优化面对千万级数据传统MD5去重会消耗15GB以上内存。我们采用分层去重策略第一层URL去重使用CRC32算法生成64位指纹Redis Set存储命中率约70%第二层内容特征去重对正文提取TF-IDF特征向量相似度90%视为重复def content_fingerprint(text): vectorizer TfidfVectorizer(min_df2, max_df0.95) tfidf vectorizer.fit_transform([text]) return hashlib.md5(tfidf.data).hexdigest()第三层布隆过滤器误判率设置为0.01%动态扩容机制class ScalableBloomFilter: def __init__(self, initial_size1000000): self.filters [BloomFilter(initial_size)] self.current 0 def add(self, item): if self.filters[self.current].count self.filters[self.current].capacity*0.8: new_filter BloomFilter(self.filters[self.current].capacity*2) self.filters.append(new_filter) self.current 1 self.filters[self.current].add(item)4. 工业数据采集专项处理4.1 反反爬策略组合工业网站的反爬往往比电商更复杂我们总结出这类网站的三个特点验证码触发频率高特别是图片验证码基于IP的行为分析严格动态参数加密普遍我们的应对方案# 在middlewares.py中实现智能切换 class AntiAntiSpiderMiddleware: def process_request(self, request, spider): if request.meta.get(retry_times, 0) 3: # 1. 自动切换代理 request.meta[proxy] self.proxy_pool.get() # 2. 降低请求频率 time.sleep(random.uniform(1, 3)) # 3. 更换浏览器指纹 request.headers.update(self.gen_new_headers()) # 4. 触发验证码识别 if captcha in response.text: return self.handle_captcha(request)4.2 数据清洗管道工业数据的脏数据率通常高达30%我们设计了三级清洗管道结构化处理# 处理各种日期格式 def normalize_date(date_str): for fmt in (%Y-%m-%d, %m/%d/%Y, %d.%m.%Y): try: return datetime.strptime(date_str, fmt).date() except ValueError: continue return None单位统一化# 将各种功率单位转为千瓦 def convert_power(value): if kW in value: return float(value.replace(kW,)) elif MW in value: return float(value.replace(MW,))*1000 elif HP in value: return float(value.replace(HP,))*0.7457异常值过滤# 基于行业标准范围校验 def validate_machine_param(param, value): ranges { voltage: (220, 10000), rotation: (500, 3000), weight: (50, 100000) } return ranges[param][0] value ranges[param][1]5. 性能优化实录5.1 Redis调优参数这些配置让我们的Redis吞吐量提升了3倍# redis.conf关键修改 maxmemory 16gb maxmemory-policy allkeys-lru hash-max-ziplist-entries 512 hash-max-ziplist-value 64 set-max-intset-entries 512 activerehashing yes5.2 爬虫节点参数在scrapy.cfg中增加机器特定配置[settings:node1] CONCURRENT_REQUESTS 32 DOWNLOAD_DELAY 0.25 AJAXCRAWL_ENABLED True AUTOTHROTTLE_TARGET_CONCURRENCY 85.3 监控看板实现使用GrafanaPrometheus构建的监控体系包含实时请求成功率各网站反爬触发频率数据入库速率节点资源占用热力图核心指标采集代码from prometheus_client import Counter, Gauge # 自定义指标 requests_total Counter(spider_requests_total, Total requests) items_scraped Counter(spider_items_scraped, Items scraped) request_latency Gauge(spider_request_latency, Request latency) # 在回调函数中埋点 def parse(self, response): start_time response.meta.get(start_time) request_latency.set(time.time() - start_time) items_scraped.inc()6. 完整代码结构项目采用模块化设计关键目录结构如下scrapy_industrial/ ├── spiders/ │ ├── __init__.py │ ├── base.py # 所有爬虫的基类 │ └── machinery/ # 按设备类型分类 ├── middlewares/ │ ├── proxies.py # 代理中间件 │ └── useragents.py # UA轮换 ├── pipelines/ │ ├── validation.py # 数据验证 │ └── dedupe.py # 去重管道 ├── utils/ │ ├── bloomfilter.py # 布隆过滤器 │ └── throttling.py # 动态限速 └── config/ ├── redis.conf # Redis优化配置 └── scrapy.cfg # 环境配置核心基类代码片段class IndustrialSpider(RedisSpider): custom_settings { ITEM_PIPELINES: { scrapy_industrial.pipelines.ValidationPipeline: 300, scrapy_industrial.pipelines.DedupePipeline: 800, }, DOWNLOADER_MIDDLEWARES: { scrapy_industrial.middlewares.ProxiesMiddleware: 543, } } def make_requests_from_url(self, url): # 统一添加工业网站必要请求头 request super().make_requests_from_url(url) request.headers.update({ Accept: application/industryjson, X-Requested-With: IndustrialDataCollector }) return request7. 踩坑经验总结Redis连接池泄露 初期没有正确关闭Redis连接导致系统出现大量TIME_WAIT状态连接。正确做法是在spider_closed信号中显式释放classmethod def from_crawler(cls, crawler): spider super().from_crawler(crawler) crawler.signals.connect(spider.spider_closed, signalsignals.spider_closed) return spider def spider_closed(self): self.server.connection_pool.disconnect()布隆过滤器误判 当数据量超过初始容量时误判率会急剧上升。我们最终实现了自动扩容方案当插入失败率达到5%时自动创建新的更大的过滤器并将旧数据批量迁移。代理IP失效 工业网站对代理IP的检测非常严格我们开发了智能检测机制每15分钟测试一次代理可用性自动屏蔽连续失败3次的IP段不同网站使用不同的代理池数据分片技巧 千万级数据如果直接写入单个文件会导致后续处理困难。我们的方案按时间分片每小时生成一个新文件按数据特征分片不同设备类型写入不同目录使用Parquet格式存储比CSV节省40%空间这套系统已经稳定运行11个月累计采集工业设备数据2300万条期间经历过服务器迁移、Redis主从切换、网站改版等各种意外情况但核心采集任务从未中断。最让我自豪的是所有设计决策都经受住了真实生产环境的考验这也是我分享这些经验的底气所在。