我在天津滨海新区做工业自动化供应链分析时,遇到过一个典型的大数据采集痛点:
- 需要抓取全国1000+工业自动化供应商的产品信息、价格、库存、联系方式,总数据量超过1500万条;
- 之前用单台Scrapy爬虫,一天最多抓20万条,全部抓完要75天,效率极低;
- 单台机器IP被封后整个爬虫停摆,没有断点续传,重新抓又要从头开始;
- 数据去重全靠内存,抓了100万条后内存溢出,程序直接崩溃。
后来我们用Scrapy-Redis分布式架构,搭了5台云服务器节点,彻底解决了这个问题:
- 5台节点并发抓取,一天抓120万条,1500万条12天就搞定了,效率提升了6倍;
- Redis做分布式队列和去重,单台节点被封不影响其他节点,断点续传随时重启;
- 千万级去重放在Redis集合里,内存占用只有单台的1/5,稳定运行无溢出;
- 代理IP池+随机User-Agent+随机延迟,反爬能力拉满,封IP率从80%降到5%。
本文就把这套经过工业级验证的分布式爬虫架构完整分享出来,从架构设计、环境搭建、核心开发、反爬应对、数据持久化到避坑指南,全流程带可复用代码,附千万级数据采集的性能优化方案。
一、为什么选「Scrapy-Redis」做千万级分布式爬虫?
千万级数据采集,单台机器、多线程、Selenium都不是最优解,Scrapy-Redis的组合是工业级分布式爬虫的事实标准:
二、整体架构设计(工业供应商数据真实架构)
我们针对千万级工业供应商数据采集,设计了4层核心架构,兼顾稳定性、扩展性、反爬能力和性能:
| 代理层 | 付费代理IP池 | 每个节点用不同代理,降低封IP率 | 快代理/阿布云/自建代理池 |
| 调度层 | Redis 7.x | 分布式请求队列、URL去重集合、爬虫状态存储 | Redis 7.x + RDB/AOF持久化 |
| 爬虫层 | 5台Scrapy节点(云服务器) | 并发抓取供应商数据、解析HTML、提取Item | Scrapy 2.11 + Scrapy-Redis 0.7 |
| 存储层 | MongoDB 6.x + MySQL 8.0 | MongoDB存非结构化产品描述,MySQL存结构化供应商信息 | PyMongo + PyMySQL + 批量写入 |
核心设计原则
三、前置环境与核心选型
1. 开发与运行环境
- 开发工具:PyCharm Professional / VS Code
- Python版本:Python 3.10+(推荐,Scrapy 2.11+要求)
- 爬虫节点:5台2核4G云服务器(Ubuntu 22.04)
- Redis服务器:1台4核8G云服务器(Ubuntu 22.04,Redis 7.x)
- 数据库服务器:1台4核8G云服务器(Ubuntu 22.04,MongoDB 6.x + MySQL 8.0)
2. 核心库安装
(1)Redis服务器安装(Ubuntu 22.04)
# 更新源
sudo apt update
# 安装Redis
sudo apt install redis-server -y
# 配置Redis:允许远程访问,设置密码
sudo nano /etc/redis/redis.conf
# 修改以下配置:
# bind 0.0.0.0 # 允许远程访问
# requirepass your_redis_password # 设置Redis密码
# 重启Redis
sudo systemctl restart redis-server
# 测试连接
redis-cli -a your_redis_password
(2)Scrapy节点环境安装(所有节点都要装)
# 安装Python 3.10+
sudo apt install python3 python3-pip -y
# 安装Scrapy、Scrapy-Redis、相关库
pip install scrapy scrapy-redis fake-useragent pymongo pymysql redis -i https://pypi.tuna.tsinghua.edu.cn/simple
四、核心开发实战:从单台到分布式
第一步:创建Scrapy项目并配置
在开发机上创建Scrapy项目,配置好后同步到所有爬虫节点:
# 创建Scrapy项目
scrapy startproject industrial_supplier
cd industrial_supplier
# 创建Spider
scrapy genspider supplier example.com
核心配置:settings.py(分布式关键)
# Scrapy settings for industrial_supplier project
import os
from fake_useragent import UserAgent
ua = UserAgent()
BOT_NAME = "industrial_supplier"
SPIDER_MODULES = ["industrial_supplier.spiders"]
NEWSPIDER_MODULE = "industrial_supplier.spiders"
# ————————– Scrapy-Redis核心配置 ————————–
# 调度器:换成Scrapy-Redis的调度器
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# 去重器:换成Scrapy-Redis的Redis集合去重器
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# Redis连接配置:替换为你的Redis服务器IP和密码
REDIS_HOST = "192.168.1.200" # Redis服务器IP
REDIS_PORT = 6379
REDIS_PARAMS = {
"password": "your_redis_password",
"db": 0
}
# 队列持久化:爬虫关闭后保留Redis队列,实现断点续传
SCHEDULER_PERSIST = True
# 队列类型:使用Redis列表(FIFO),也可以用SortedSet做优先级队列
SCHEDULER_QUEUE_CLASS = "scrapy_redis.queue.SpiderQueue"
# ————————– Scrapy基础配置 ————————–
# 随机User-Agent
USER_AGENT = ua.random
# 请求头
DEFAULT_REQUEST_HEADERS = {
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*;q=0.8",
"Accept-Language": "zh-CN,zh;q=0.9,en;q=0.8",
"Referer": "https://www.baidu.com/",
}
# 并发数:每个节点16个并发(根据服务器配置调整)
CONCURRENT_REQUESTS = 16
# 下载延迟:随机0.5-2秒,避免被封
DOWNLOAD_DELAY = 0.5
RANDOMIZE_DOWNLOAD_DELAY = True
# 超时时间:30秒
DOWNLOAD_TIMEOUT = 30
# 重试次数:3次
RETRY_ENABLED = True
RETRY_TIMES = 3
# 中间件:添加代理中间件、随机User-Agent中间件
DOWNLOADER_MIDDLEWARES = {
"industrial_supplier.middlewares.RandomUserAgentMiddleware": 543,
"industrial_supplier.middlewares.ProxyMiddleware": 544,
}
# Item Pipeline:数据清洗和存储
ITEM_PIPELINES = {
"industrial_supplier.pipelines.DataCleanPipeline": 300,
"industrial_supplier.pipelines.MongoDBPipeline": 400,
"industrial_supplier.pipelines.MySQLPipeline": 500,
}
第二步:编写中间件(反爬核心)
(1)随机User-Agent中间件
from fake_useragent import UserAgent
class RandomUserAgentMiddleware:
def __init__(self):
self.ua = UserAgent()
def process_request(self, request, spider):
# 每次请求换一个随机User-Agent
request.headers["User-Agent"] = self.ua.random
return None
(2)代理IP中间件
import random
# 代理IP池:替换为你的代理IP列表(格式:http://username:password@ip:port)
PROXY_POOL = [
"http://user:pass@123.123.123.123:8888",
"http://user:pass@456.456.456.456:8888",
]
class ProxyMiddleware:
def process_request(self, request, spider):
# 随机选一个代理IP
if PROXY_POOL:
proxy = random.choice(PROXY_POOL)
request.meta["proxy"] = proxy
return None
第三步:编写Spider(继承RedisSpider)
关键:继承RedisSpider,而不是scrapy.Spider,从Redis队列取请求:
import scrapy
from scrapy_redis.spiders import RedisSpider
from industrial_supplier.items import SupplierItem, ProductItem
class SupplierSpider(RedisSpider):
name = "supplier"
# Redis队列的key:爬虫从这个key取请求
redis_key = "supplier:start_urls"
# 可选:Redis队列的默认起始URL(如果Redis队列空的话)
# redis_batch_size = 16 # 每次从Redis取16个请求
def parse(self, response):
"""
解析供应商列表页,提取供应商详情页URL和下一页URL
"""
# 提取供应商详情页URL
supplier_urls = response.css("a.supplier-link::attr(href)").getall()
for url in supplier_urls:
# 构造详情页请求,交给parse_supplier解析
yield scrapy.Request(url, callback=self.parse_supplier)
# 提取下一页URL,继续爬列表页
next_page = response.css("a.next-page::attr(href)").get()
if next_page:
yield scrapy.Request(next_page, callback=self.parse)
def parse_supplier(self, response):
"""
解析供应商详情页,提取供应商信息和产品信息
"""
# 提取供应商信息
supplier_item = SupplierItem()
supplier_item["supplier_id"] = response.css("div.supplier-id::text").get().strip()
supplier_item["name"] = response.css("h1.supplier-name::text").get().strip()
supplier_item["contact"] = response.css("div.contact::text").get().strip()
supplier_item["phone"] = response.css("div.phone::text").get().strip()
supplier_item["address"] = response.css("div.address::text").get().strip()
supplier_item["url"] = response.url
yield supplier_item
# 提取产品信息
product_items = response.css("div.product-item")
for item in product_items:
product = ProductItem()
product["product_id"] = item.css("div.product-id::text").get().strip()
product["supplier_id"] = supplier_item["supplier_id"]
product["name"] = item.css("div.product-name::text").get().strip()
product["price"] = item.css("div.product-price::text").get().strip()
product["stock"] = item.css("div.product-stock::text").get().strip()
product["description"] = item.css("div.product-desc::text").get().strip()
yield product
第四步:编写Item和Pipeline(数据清洗和存储)
(1)items.py
import scrapy
class SupplierItem(scrapy.Item):
supplier_id = scrapy.Field()
name = scrapy.Field()
contact = scrapy.Field()
phone = scrapy.Field()
address = scrapy.Field()
url = scrapy.Field()
class ProductItem(scrapy.Item):
product_id = scrapy.Field()
supplier_id = scrapy.Field()
name = scrapy.Field()
price = scrapy.Field()
stock = scrapy.Field()
description = scrapy.Field()
(2)pipelines.py(数据清洗+MongoDB+MySQL批量写入)
import pymongo
import pymysql
from pymysql.cursors import DictCursor
class DataCleanPipeline:
"""数据清洗:去除空格、处理空值"""
def process_item(self, item, spider):
for key, value in item.items():
if isinstance(value, str):
item[key] = value.strip()
if not value:
item[key] = None
return item
class MongoDBPipeline:
"""存储非结构化产品描述到MongoDB"""
def __init__(self, mongo_uri, mongo_db):
self.mongo_uri = mongo_uri
self.mongo_db = mongo_db
self.client = None
self.db = None
self.product_buffer = [] # 批量写入缓冲区
self.buffer_size = 1000 # 每1000条批量写入一次
@classmethod
def from_crawler(cls, crawler):
return cls(
mongo_uri=crawler.settings.get("MONGO_URI", "mongodb://localhost:27017/"),
mongo_db=crawler.settings.get("MONGO_DB", "industrial_supplier")
)
def open_spider(self, spider):
self.client = pymongo.MongoClient(self.mongo_uri)
self.db = self.client[self.mongo_db]
# 创建唯一索引,避免重复存储
self.db["products"].create_index("product_id", unique=True)
def close_spider(self, spider):
# 爬虫关闭时写入剩余的缓冲区数据
if self.product_buffer:
self.db["products"].insert_many(self.product_buffer)
self.client.close()
def process_item(self, item, spider):
if isinstance(item, ProductItem):
# 产品描述存MongoDB
self.product_buffer.append(dict(item))
if len(self.product_buffer) >= self.buffer_size:
try:
self.db["products"].insert_many(self.product_buffer)
except pymongo.errors.BulkWriteError:
# 忽略重复键错误
pass
self.product_buffer = []
return item
class MySQLPipeline:
"""存储结构化供应商信息到MySQL,批量写入"""
def __init__(self, mysql_host, mysql_port, mysql_user, mysql_password, mysql_db):
self.mysql_host = mysql_host
self.mysql_port = mysql_port
self.mysql_user = mysql_user
self.mysql_password = mysql_password
self.mysql_db = mysql_db
self.conn = None
self.cursor = None
self.supplier_buffer = []
self.buffer_size = 500
@classmethod
def from_crawler(cls, crawler):
return cls(
mysql_host=crawler.settings.get("MYSQL_HOST", "localhost"),
mysql_port=crawler.settings.get("MYSQL_PORT", 3306),
mysql_user=crawler.settings.get("MYSQL_USER", "root"),
mysql_password=crawler.settings.get("MYSQL_PASSWORD", ""),
mysql_db=crawler.settings.get("MYSQL_DB", "industrial_supplier")
)
def open_spider(self, spider):
self.conn = pymysql.connect(
host=self.mysql_host,
port=self.mysql_port,
user=self.mysql_user,
password=self.mysql_password,
database=self.mysql_db,
charset="utf8mb4",
cursorclass=DictCursor
)
self.cursor = self.conn.cursor()
# 创建供应商表
self.cursor.execute("""
CREATE TABLE IF NOT EXISTS suppliers (
id INT AUTO_INCREMENT PRIMARY KEY,
supplier_id VARCHAR(50) UNIQUE NOT NULL,
name VARCHAR(200) NOT NULL,
contact VARCHAR(100),
phone VARCHAR(20),
address TEXT,
url TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
""")
self.conn.commit()
def close_spider(self, spider):
if self.supplier_buffer:
self._batch_insert()
self.cursor.close()
self.conn.close()
def _batch_insert(self):
sql = """
INSERT INTO suppliers (supplier_id, name, contact, phone, address, url)
VALUES (%s, %s, %s, %s, %s, %s)
ON DUPLICATE KEY UPDATE
name=VALUES(name), contact=VALUES(contact), phone=VALUES(phone),
address=VALUES(address), url=VALUES(url)
"""
values = [
(item["supplier_id"], item["name"], item["contact"], item["phone"], item["address"], item["url"])
for item in self.supplier_buffer
]
self.cursor.executemany(sql, values)
self.conn.commit()
self.supplier_buffer = []
def process_item(self, item, spider):
if isinstance(item, SupplierItem):
self.supplier_buffer.append(item)
if len(self.supplier_buffer) >= self.buffer_size:
self._batch_insert()
return item
第五步:启动分布式爬虫
(1)向Redis队列推送起始URL
在Redis服务器上执行,或者用Python脚本推送:
import redis
# 连接Redis
r = redis.Redis(host="192.168.1.200", port=6379, password="your_redis_password", db=0)
# 推送起始URL到Redis队列
start_urls = [
"https://www.example.com/suppliers/page1.html",
"https://www.example.com/suppliers/page2.html",
]
for url in start_urls:
r.lpush("supplier:start_urls", url)
print(f"已推送 {len(start_urls)} 个起始URL到Redis队列")
(2)启动所有爬虫节点
在每台爬虫节点上执行:
cd industrial_supplier
scrapy crawl supplier
五、千万级数据采集性能优化方案
- 设置Redis的maxmemory和maxmemory-policy,避免内存溢出;
- 去重集合的URL可以压缩存储,减少内存占用;
- 定期清理已完成的请求队列,释放内存。
- MongoDB和MySQL都用批量写入,每500-1000条写一次,减少IO;
- 数据库连接池复用,避免频繁创建和关闭连接。
- 根据服务器配置和目标网站反爬强度,调整CONCURRENT_REQUESTS,建议16-32;
- 用Redis的redis_batch_size调整每次取的请求数,建议16-32。
- 用付费代理IP池,质量高,封IP率低;
- 定期检测代理IP的可用性,剔除失效IP。
六、避坑指南(千万级数据踩过的坑)
- 原因:千万级URL去重集合占用内存太大;
- 解决:设置Redis的maxmemory为服务器内存的70%,maxmemory-policy为allkeys-lru,定期清理已完成的队列。
- 原因:URL有参数顺序不同但内容相同的情况;
- 解决:在Spider里对URL参数排序,生成唯一指纹,再交给Scrapy-Redis去重。
- 原因:部分节点快,部分节点慢;
- 解决:用Redis的列表做队列,所有节点从同一个列表取,自动负载均衡。
- 原因:爬虫关闭时没有持久化Redis队列;
- 解决:在settings.py里设置SCHEDULER_PERSIST = True,Redis开启RDB/AOF持久化。
- 原因:网络波动导致同一个请求被处理多次;
- 解决:在数据库里创建唯一索引,Pipeline里用INSERT ON DUPLICATE KEY UPDATE。
七、真实项目效果对比
| 抓取1500万条时间 | 75天 | 12天 |
| 日均抓取量 | 20万条 | 125万条 |
| 封IP率 | 80% | 5% |
| 内存占用(去重) | 10GB+ | 2GB |
| 断点续传 | 不支持 | 支持 |
| 扩展性 | 差 | 好(随时加节点) |
八、总结
这套「Scrapy-Redis+MongoDB+MySQL」的分布式爬虫架构,已经帮我们抓取了超过3000万条工业供应商数据,稳定运行无故障,非常适合千万级以上的大数据采集场景。
核心要点总结:




