前言

当我们从入门爬虫走到进阶之路,单机单线程的爬虫已经远远无法满足生产环境的需求。面对动辄百万级的数据量、频繁的反爬策略变更、以及实时性的要求,我们必须掌握高并发与分布式爬虫架构的设计与实现。这篇文章是我在实际项目中沉淀的笔记与思考,希望能帮到正在进阶的你。


一、高并发爬虫基础:同步 / 异步 / 协程

在开始写代码之前,先理清三个核心概念。

| 模型 | 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、控制频率、注意数据使用的法律边界

爬虫的本质是 "程序化的网络请求 + 结构化解析",掌握分布式爬虫架构的核心在于理解 「任务如何分发、状态如何共享、数据如何汇聚」 这三个问题。

希望这篇文章能帮你搭建起自己的高并发分布式爬虫系统。如果有什么地方写得不够清楚,欢迎在评论区交流讨论。


*如果你对本文有任何疑问或想了解更多实战细节,欢迎留言,我会持续更新补充。*