前言
当我们从入门爬虫走到进阶之路,单机单线程的爬虫已经远远无法满足生产环境的需求。面对动辄百万级的数据量、频繁的反爬策略变更、以及实时性的要求,我们必须掌握高并发与分布式爬虫架构的设计与实现。这篇文章是我在实际项目中沉淀的笔记与思考,希望能帮到正在进阶的你。
一、高并发爬虫基础:同步 / 异步 / 协程
在开始写代码之前,先理清三个核心概念。
| 模型 | IO 等待处理方式 | 资源开销 | 典型场景 |
|------|----------------|----------|----------|
| 同步阻塞 | 等待 IO 完成,CPU 空转 | 高(每请求一线程) | 简单脚本、快速验证 |
| 多线程/多进程 | 并发执行,操作系统调度 | 中高(线程/进程切换) | IO 密集型任务 |
| 异步 + 协程 | 事件循环驱动,用户态切换 | 极低(单线程可管理万级连接) | 高并发网络请求 |
对于爬虫这种典型的 IO 密集型任务,协程是性价比最高的方案。一个线程内可以跑成千上万个协程,上下文切换成本几乎可以忽略。
# 伪代码对比:同步 vs 异步
# --- 同步爬虫 ---
def fetch_sync(url):
resp = requests.get(url) # 阻塞在此,CPU 空转
return resp.text
# --- 异步协程爬虫 ---
async def fetch_async(url):
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
return await resp.text() # 挂起协程,事件循环处理其他任务
核心区别在于:同步版本在等待网络响应时线程是阻塞的;异步版本在 await 处自动挂起,事件循环去执行其他就绪的协程,真正做到了 "人等事 -> 事等人"。
二、asyncio + aiohttp 异步爬虫实战
下面是一个完整的异步爬虫示例,爬取一个分页 API 并解析 JSON 数据。
import asyncio
import aiohttp
import json
BASE_URL = "https://api.example.com/items"
CONCURRENCY = 20 # 最大并发数
async def fetch_one(session, url, sem):
"""带信号量控制的单次请求"""
async with sem:
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=10)) as resp:
if resp.status != 200:
print(f"请求失败: {url}, status={resp.status}")
return None
return await resp.json()
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
print(f"请求异常: {url}, error={e}")
return None
async def main():
sem = asyncio.Semaphore(CONCURRENCY)
connector = aiohttp.TCPConnector(limit=CONCURRENCY, limit_per_host=10)
async with aiohttp.ClientSession(connector=connector) as session:
# 构造 100 个分页请求
tasks = []
for page in range(1, 101):
url = f"{BASE_URL}?page={page}"
tasks.append(fetch_one(session, url, sem))
results = await asyncio.gather(*tasks)
# 过滤失败的结果
valid = [r for r in results if r is not None]
print(f"成功获取 {len(valid)} 页数据")
if __name__ == "__main__":
asyncio.run(main())
2.1 关键点解析
- **`asyncio.Semaphore`**:限制并发数,防止瞬间打爆目标服务器,也避免被 ban IP。
- **`TCPConnector(limit=...)`**:控制连接池大小,避免文件描述符耗尽。
- **`asyncio.gather`**:并发执行所有协程,等待全部完成。
- **超时处理**:`ClientTimeout` 设置请求超时,防止某个慢请求拖垮整个队列。
三、高并发控制:信号量、连接池、请求限速
高并发爬虫的核心不是 "快",而是 "又快又稳"。下面是三种常用的控制手段。
3.1 令牌桶算法(请求限速)
import asyncio
class TokenBucket:
"""简单的令牌桶限速器"""
def __init__(self, rate: float, capacity: int):
self.rate = rate # 每秒补充的令牌数
self.capacity = capacity # 桶容量
self.tokens = capacity
self.last_refill = asyncio.get_event_loop().time()
self._lock = asyncio.Lock()
async def acquire(self):
async with self._lock:
now = asyncio.get_event_loop().time()
elapsed = now - self.last_refill
self.tokens = min(self.capacity, self.tokens + elapsed * self.rate)
self.last_refill = now
if self.tokens < 1:
wait_time = (1 - self.tokens) / self.rate
await asyncio.sleep(wait_time)
self.tokens = 0
else:
self.tokens -= 1
# 使用示例
bucket = TokenBucket(rate=10, capacity=20) # 每秒 10 个请求,突发可到 20
async def controlled_request(url):
await bucket.acquire()
# 发送请求...
3.2 重试机制与退避策略
async def fetch_with_retry(session, url, retries=3):
for attempt in range(retries):
try:
async with session.get(url) as resp:
if resp.status == 429: # Too Many Requests
wait = 2 ** attempt # 指数退避
await asyncio.sleep(wait)
continue
resp.raise_for_status()
return await resp.text()
except (aiohttp.ClientError, asyncio.TimeoutError):
if attempt == retries - 1:
raise
await asyncio.sleep(2 ** attempt)
return None
四、分布式爬虫架构设计:Master-Worker 模式
单机性能终究有上限——网络带宽、CPU、内存、IP 封禁阈值。要突破这些瓶颈,必须走向分布式。
4.1 整体架构
┌──────────────────────────────────────────────────┐
│ Master │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ URL 调度 │ │ 去重管理 │ │ 结果收集 │ │
│ └─────┬────┘ └────┬─────┘ └─────┬────┘ │
└─────────┼─────────────┼───────────────┼──────────┘
│ │ │
┌─────▼─────────────▼───────────────▼─────┐
│ 消息队列 (Redis/Kafka) │
└─────▲─────────────▲───────────────▲─────┘
│ │ │
┌─────────┼─────────────┼───────────────┼──────────┐
│ ┌──────┴──────┐ ┌────┴──────┐ ┌────┴──────┐ │
│ │ Worker 1 │ │ Worker 2 │ │ Worker N │ │
│ │ (爬虫节点) │ │ (爬虫节点) │ │ (爬虫节点) │ │
│ └─────────────┘ └───────────┘ └───────────┘ │
│ Worker 集群 │
└─────────────────────────────────────────────────┘
4.2 各组件职责
| 组件 | 职责 | 技术选型建议 |
|------|------|-------------|
| Master | URL 调度、任务分配、去重、结果收集 | Python + Redis |
| 消息队列 | 解耦 Master 和 Worker,削峰填谷 | Redis / RabbitMQ / Kafka |
| Worker | 实际执行爬取与解析,无状态,可水平扩展 | Scrapy / aiohttp |
| 结果存储 | 持久化爬取结果 | MongoDB / Elasticsearch / HDFS |
4.3 Redis 数据结构在爬虫中的应用
# URL 待爬队列 (List)
LPUSH spider:start_urls https://example.com/page1
RPOP spider:start_urls
# 已去重集合 (Set)
SADD spider:visited https://example.com/page1
# 布隆过滤器去重 (Redis 模块)
BF.ADD spider:bloom https://example.com/page1
BF.EXISTS spider:bloom https://example.com/page1 # 返回 0/1
五、消息队列在爬虫中的应用
不同的消息队列适应不同的场景,选型时需权衡吞吐量、持久化需求和运维成本。
5.1 Redis List — 轻量级任务队列
适合中小规模爬虫(日百万级),天然支持 LPUSH/RPUSH + BRPOP 阻塞消费。
import redis
import json
r = redis.Redis(host='redis-master', decode_responses=True)
# Master:推送任务
def push_task(url, meta=None):
task = {"url": url, "meta": meta or {}}
r.lpush("spider:task_queue", json.dumps(task))
# Worker:阻塞消费
def get_task(timeout=30):
_, data = r.brpop("spider:task_queue", timeout=timeout)
return json.loads(data)
5.2 RabbitMQ — 可靠的消息投递
适合需要消息确认、死信队列、路由策略的复杂场景。
import pika
# Worker 端
connection = pika.BlockingConnection(pika.URLParameters("amqp://guest:guest@rabbitmq:5672"))
channel = connection.channel()
channel.queue_declare(queue='crawl_tasks', durable=True)
def callback(ch, method, properties, body):
task = json.loads(body)
print(f"爬取: {task['url']}")
ch.basic_ack(delivery_tag=method.delivery_tag) # 确认完成
channel.basic_qos(prefetch_count=1)
channel.basic_consume(queue='crawl_tasks', on_message_callback=callback)
channel.start_consuming()
5.3 Kafka — 海量日志与高吞吐
适合每天数亿 URL 的极端场景,分区机制天然支持并行消费。
from kafka import KafkaConsumer, KafkaProducer
import json
# Producer (Master)
producer = KafkaProducer(
bootstrap_servers=['kafka1:9092'],
value_serializer=lambda v: json.dumps(v).encode()
)
producer.send('crawl-topic', {"url": "https://example.com"})
# Consumer (Worker)
consumer = KafkaConsumer(
'crawl-topic',
bootstrap_servers=['kafka1:9092'],
group_id='crawl-workers',
enable_auto_commit=True,
auto_offset_reset='earliest'
)
for msg in consumer:
task = msg.value
# 爬取逻辑...
六、Scrapy 框架深入:中间件、管道与分布式部署
6.1 自定义下载中间件——代理轮换 + 随机 UA
# middlewares.py
import random
from scrapy import signals
class RotateProxyMiddleware:
"""动态代理中间件"""
def __init__(self, proxy_list):
self.proxy_list = proxy_list
@classmethod
def from_crawler(cls, crawler):
proxy_list = crawler.settings.get('PROXY_LIST', [])
return cls(proxy_list)
def process_request(self, request, spider):
proxy = random.choice(self.proxy_list)
request.meta['proxy'] = proxy
spider.logger.debug(f"使用代理: {proxy}")
def process_response(self, request, response, spider):
if response.status in (403, 429):
# 代理被封,重试
request.meta['proxy'] = random.choice(self.proxy_list)
return request
return response
class RandomUserAgentMiddleware:
"""随机 User-Agent 中间件"""
def __init__(self, user_agents):
self.user_agents = user_agents
@classmethod
def from_crawler(cls, crawler):
uas = crawler.settings.get('USER_AGENT_LIST', [])
return cls(uas)
def process_request(self, request, spider):
request.headers['User-Agent'] = random.choice(self.user_agents)
6.2 自定义 Pipeline——数据去重与存储
# pipelines.py
import hashlib
import redis
import pymongo
class DuplicateFilterPipeline:
"""基于 Redis 的布隆过滤器去重"""
def open_spider(self, spider):
self.redis_client = redis.Redis(
host=spider.settings.get('REDIS_HOST', 'localhost'),
port=spider.settings.get('REDIS_PORT', 6379),
decode_responses=True
)
def process_item(self, item, spider):
key = hashlib.md5(str(item).encode()).hexdigest()
if self.redis_client.sadd("spider:duplicate_keys", key):
return item
else:
raise DropItem(f"重复数据已丢弃: {item}")
class MongoPipeline:
"""MongoDB 持久化管道"""
def open_spider(self, spider):
self.client = pymongo.MongoClient(
host=spider.settings.get('MONGO_URI', 'mongodb://localhost:27017')
)
self.db = self.client[spider.settings.get('MONGO_DB', 'crawler')]
self.collection = self.db[spider.name]
def process_item(self, item, spider):
self.collection.insert_one(dict(item))
return item
def close_spider(self, spider):
self.client.close()
6.3 settings.py 中的并发配置
# settings.py
CONCURRENT_REQUESTS = 32 # 全局并发
CONCURRENT_REQUESTS_PER_DOMAIN = 8 # 每域名并发
DOWNLOAD_DELAY = 0.5 # 请求间隔(秒)
RANDOMIZE_DOWNLOAD_DELAY = True # 随机化间隔
# 中间件优先级
DOWNLOADER_MIDDLEWARES = {
'mycrawler.middlewares.RandomUserAgentMiddleware': 400,
'mycrawler.middlewares.RotateProxyMiddleware': 500,
}
# Pipeline 优先级
ITEM_PIPELINES = {
'mycrawler.pipelines.DuplicateFilterPipeline': 100,
'mycrawler.pipelines.MongoPipeline': 300,
}
七、分布式爬虫实战:基于 Scrapy-Redis
[Scrapy-Redis](https://github.com/rmax/scrapy-redis) 是 Scrapy 最成熟的分布式扩展,它用 Redis 替换了 Scrapy 默认的内存队列,让多个爬虫节点共享任务队列和去重集合。
7.1 安装与配置
pip install scrapy-redis
# settings.py — 启用 Scrapy-Redis
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
REDIS_HOST = 'redis-master' # Redis 服务地址
REDIS_PORT = 6379
REDIS_PARAMS = {'password': 'your_password'}
SCHEDULER_PERSIST = True # 爬取结束后保持 Redis 中的任务队列
7.2 Spider 实现
# spiders/news_spider.py
from scrapy_redis.spiders import RedisSpider
from mycrawler.items import NewsItem
class NewsDistributedSpider(RedisSpider):
"""继承 RedisSpider,从 Redis 队列读取 URL"""
name = 'news_distributed'
redis_key = 'news:start_urls' # Master 向此 key 推送起始 URL
def parse(self, response):
# 解析列表页
for article_url in response.css('a.article-link::attr(href)').getall():
yield scrapy.Request(
url=response.urljoin(article_url),
callback=self.parse_article
)
# 处理分页
next_page = response.css('a.next::attr(href)').get()
if next_page:
yield scrapy.Request(
url=response.urljoin(next_page),
callback=self.parse
)
def parse_article(self, response):
item = NewsItem()
item['title'] = response.css('h1.article-title::text').get()
item['content'] = response.css('div.article-content::text').getall()
item['url'] = response.url
item['crawled_at'] = datetime.now().isoformat()
yield item
7.3 部署与运行
# Master 节点:推送起始 URL
redis-cli LPUSH news:start_urls "https://example.com/news"
# 任意 Worker 节点启动爬虫(可启动 N 个)
scrapy crawl news_distributed
# 停止后恢复(SCHEUDLER_PERSIST=True 时)
scrapy crawl news_distributed # 自动从 Redis 恢复队列
这种架构下,Worker 节点是无状态的,你可以随意增加或减少节点数量,所有节点通过 Redis 共享状态,实现了真正的水平扩展。
八、高并发场景下的数据一致性保障
分布式爬虫中,数据一致性往往被忽视,但却是最致命的坑。
8.1 幂等性设计
每条数据应具备天然的唯一标识,即使重复爬取也不会产生脏数据。
# 使用 URL + 时间戳 hash 作为 _id
item['_id'] = hashlib.md5(f"{item['url']}_{item['title']}".encode()).hexdigest()
# MongoDB upsert 保证幂等
self.collection.update_one(
{"_id": item['_id']},
{"$set": dict(item)},
upsert=True
)
8.2 任务确认与防丢失
| 场景 | 解决方案 |
|------|---------|
| Worker 崩溃导致任务丢失 | 使用 RabbitMQ ack 机制 / Redis BRPOPLPUSH |
| 重复消费 | 使用 Redis Set / 布隆过滤器去重 |
| 部分爬取成功 | 维护 URL 状态(pending/done/failed)|
# 使用 BRPOPLPUSH 实现安全消费
def safe_get_task(redis_client, queue_name, backup_queue, timeout=30):
"""从主队列取任务,同时备份到 backup 队列"""
task_data = redis_client.brpoplpush(queue_name, backup_queue, timeout=timeout)
if task_data:
task = json.loads(task_data)
try:
yield crawl(task['url'])
# 成功:从 backup 队列移除
redis_client.lrem(backup_queue, 0, task_data)
except Exception:
# 失败:重新放回主队列,等待重试
redis_client.lpush(queue_name, task_data)
redis_client.lrem(backup_queue, 0, task_data)
8.3 分布式去重——布隆过滤器
当 URL 规模达到亿级时,Redis Set 的内存开销会变得不可接受。此时布隆过滤器是更好的选择。
import redis
from pybloom_live import ScalableBloomFilter
class BloomFilterDupeFilter:
"""基于布隆过滤器的分布式去重"""
def __init__(self, redis_client, key, capacity=100_000_000, error_rate=0.001):
self.redis = redis_client
self.key = key
self.capacity = capacity
self.error_rate = error_rate
self.local = ScalableBloomFilter(
initial_capacity=capacity,
error_rate=error_rate
)
def request_seen(self, request):
fp = request.url # 实际建议对 URL 做 hash
if fp in self.local:
return True
# 同步到 Redis(定期 persist 减少网络开销)
self.local.add(fp)
return False
九、架构选型建议与总结
9.1 选型矩阵
| 规模 | 并发量 | 建议方案 |
|------|--------|---------|
| 小(千级/天) | 不需要 | 单机 requests + BeautifulSoup |
| 中(万~百万/天) | 50~500 | asyncio + aiohttp + Redis 队列 |
| 大(千万/天) | 500~5000 | Scrapy + Scrapy-Redis + 多 Worker |
| 超大(亿级/天) | 5000+ | Kafka + 自研异步引擎 + Kubernetes |
9.2 写在最后
高并发分布式爬虫远不止是 "写得快",它涉及以下几个维度的平衡:
1. 效率与优雅:不要为了并发而并发,评估目标网站的承载能力
2. 稳定性:重试、退避、降级、熔断,一个都不能少
3. 可观测性:日志、指标(QPS、成功率、延迟)、告警是生产环境的基础设施
4. 合规性:尊重 robots.txt、控制频率、注意数据使用的法律边界
爬虫的本质是 "程序化的网络请求 + 结构化解析",掌握分布式爬虫架构的核心在于理解 「任务如何分发、状态如何共享、数据如何汇聚」 这三个问题。
希望这篇文章能帮你搭建起自己的高并发分布式爬虫系统。如果有什么地方写得不够清楚,欢迎在评论区交流讨论。
*如果你对本文有任何疑问或想了解更多实战细节,欢迎留言,我会持续更新补充。*
评论