1. 项目概述
在工业供应链领域,数据采集一直是个既关键又头疼的问题。去年我们团队接手了一个工业供应商数据采集项目,需要从200多个行业门户网站抓取设备参数、企业资质、产品规格等结构化数据,最终目标是在3个月内完成千万级数据的采集清洗。传统单机爬虫在面对这种量级时,不仅效率低下,而且一旦中断就得从头再来。经过多轮技术选型,我们最终采用Scrapy-Redis构建分布式爬虫系统,实现了日均50万条稳定采集,完整代码已封装成可复用的组件库。
这个方案最让我满意的不是性能提升(虽然确实很显著),而是系统展现出的工程化特性:任何节点崩溃都不会影响整体任务,新增机器能线性提升吞吐量,去重精度达到99.99%,这些特性让后期运维成本降低了70%。下面我就拆解这套架构的核心设计,包含你一定能用上的实战技巧。
2. 核心架构设计
2.1 为什么选择Scrapy-Redis
在分布式爬虫领域,技术选型直接决定后期开发维护成本。我们对比了三种主流方案:
| 方案 | 开发成本 | 扩展性 | 断点续传 | 去重机制 |
|---|---|---|---|---|
| 纯Scrapy集群 | 高 | 差 | 需自定义 | 内存去重 |
| Scrapy+RabbitMQ | 中 | 优 | 支持 | 需外接 |
| 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 # 使用集合存储初始URL
2.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', f'node_{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_df=2, max_df=0.95) tfidf = vectorizer.fit_transform([text]) return hashlib.md5(tfidf.data).hexdigest()第三层:布隆过滤器
- 误判率设置为0.01%
- 动态扩容机制
class ScalableBloomFilter: def __init__(self, initial_size=1000000): 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 监控看板实现
使用Grafana+Prometheus构建的监控体系包含:
- 实时请求成功率
- 各网站反爬触发频率
- 数据入库速率
- 节点资源占用热力图
核心指标采集代码:
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/industry+json', '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, signal=signals.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主从切换、网站改版等各种意外情况,但核心采集任务从未中断。最让我自豪的是,所有设计决策都经受住了真实生产环境的考验,这也是我分享这些经验的底气所在。