Scrapy框架实战教程(下):高级功能与实战优化,打造百万级爬虫系统
这是《Scrapy框架实战教程》的下篇。如果你还没读过上篇,建议先看上篇,那里有Scrapy的基础概念和第一个实战案例。
上篇我们完成了从0到1的入门,搭建了一个可以跑通的基本爬虫。但这只是开始。在实际项目中,你会遇到各种挑战:
• 反爬机制越来越强,需要随机UA、代理IP、动态渲染
• 单机性能有限,需要分布式架构
• 数据规模扩大,需要高效的存储和监控
• 代码需要维护,需要模块化和扩展性
这篇下篇,我会带你深入Scrapy的高级功能,解决这些真实场景的问题。
高级功能一:中间件
中间件是Scrapy最强大的扩展机制。简单说,它可以在请求发送前、响应返回后插入自定义逻辑,像一个个"过滤器"处理所有进出爬虫的请求和响应。
中间件的工作原理
中间件分为两种:
• 下载中间件(Downloader Middleware): 处理请求和响应
• 爬虫中间件(Spider Middleware): 处理爬虫输出的Item和Request
我们主要讲下载中间件,因为它更常用。
请求流程:
引擎 → 中间件(入站) → 下载器 → 中间件(出站) → 爬虫
响应流程:
爬虫 → 中间件(入站) → 引擎 → 中间件(出站) → 下载器
每个中间件必须实现三个方法:
• process_request(request, spider): 请求发送前调用
• process_response(request, response, spider): 响应返回后调用
• process_exception(request, exception, spider): 请求异常时调用
实战1:随机UA中间件
最基础的反爬是User-Agent检测。我们来写一个随机UA中间件。
# middlewares.py
import random
from scrapy import signals
class RandomUserAgentMiddleware:
"""随机User-Agent中间件"""
def __init__(self, user_agent_list):
self.user_agent_list = user_agent_list
@classmethod
def from_crawler(cls, crawler):
"""从settings读取配置"""
return cls(
user_agent_list=crawler.settings.get('USER_AGENT_LIST')
)
def process_request(self, request, spider):
"""为每个请求随机分配UA"""
request.headers['User-Agent'] = random.choice(self.user_agent_list)
spider.logger.debug(f"User-Agent: {request.headers['User-Agent']}")
在settings.py中配置:
# 启用中间件
DOWNLOADER_MIDDLEWARES = {
'book_scraper.middlewares.RandomUserAgentMiddleware': 400,
}
# UA列表
USER_AGENT_LIST = [
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
'Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:121.0) Gecko/20100101 Firefox/121.0',
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.1 Safari/605.1.15',
'Mozilla/5.0 (iPhone; CPU iPhone OS 17_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Mobile/15E148 Safari/604.1',
]
实战2:代理IP池中间件
当IP被封时,需要自动切换代理。我们来实现一个带重试的代理池中间件。
# middlewares.py
import random
from scrapy import signals
from scrapy.exceptions import NotConfigured
class ProxyPoolMiddleware:
"""代理IP池中间件"""
def __init__(self, proxy_list):
self.proxy_list = proxy_list
if not self.proxy_list:
raise NotConfigured("Proxy list is empty")
@classmethod
def from_crawler(cls, crawler):
return cls(
proxy_list=crawler.settings.get('PROXY_LIST', [])
)
def process_request(self, request, spider):
"""为每个请求随机分配代理"""
proxy = random.choice(self.proxy_list)
request.meta['proxy'] = proxy
spider.logger.debug(f"Using proxy: {proxy}")
def process_exception(self, request, exception, spider):
"""请求异常时记录"""
spider.logger.warning(f"Proxy failed: {request.meta.get('proxy')}, Exception: {exception}")
配置settings.py:
DOWNLOADER_MIDDLEWARES = {
'book_scraper.middlewares.ProxyPoolMiddleware': 410,
}
# 代理列表(实际项目中建议从API动态获取)
PROXY_LIST = [
'http://proxy1.example.com:8080',
'http://proxy2.example.com:8080',
'http://proxy3.example.com:8080',
]
实战3:重试中间件
当请求失败时自动重试,提高成功率。
# middlewares.py
from scrapy.downloadermiddlewares.retry import RetryMiddleware
from scrapy.utils.response import response_status_message
class CustomRetryMiddleware(RetryMiddleware):
"""自定义重试中间件"""
def __init__(self, settings):
super().__init__(settings)
self.max_retry_times = settings.getint('RETRY_TIMES', 3)
self.retry_http_codes = set(int(x) for x in settings.getlist('RETRY_HTTP_CODES'))
def process_response(self, request, response, spider):
"""检查响应状态码,需要重试时重新请求"""
if response.status in self.retry_http_codes:
reason = response_status_message(response.status)
return self._retry(request, reason, spider) or response
return response
配置:
# 重试设置
RETRY_ENABLED = True
RETRY_TIMES = 3
RETRY_HTTP_CODES = [500, 502, 503, 504, 408, 429]
# 启用重试中间件(默认已启用,优先级越高越先执行)
DOWNLOADER_MIDDLEWARES = {
'scrapy.downloadermiddlewares.retry.RetryMiddleware': 90,
'book_scraper.middlewares.ProxyPoolMiddleware': 410,
}
实战4:动态渲染中间件
很多网站使用JavaScript动态加载内容,普通的HTTP请求无法获取完整数据。这时需要用Selenium或Playwright渲染页面。
我们用Playwright(比Selenium更快、更稳定):
# middlewares.py
from scrapy.http import HtmlResponse
from playwright.sync_api import sync_playwright
class PlaywrightMiddleware:
"""Playwright动态渲染中间件"""
def __init__(self):
self.playwright = None
self.browser = None
@classmethod
def from_crawler(cls, crawler):
middleware = cls()
crawler.signals.connect(middleware.spider_opened, signals.spider_opened)
crawler.signals.connect(middleware.spider_closed, signals.spider_closed)
return middleware
def spider_opened(self, spider):
"""爬虫启动时初始化Playwright"""
self.playwright = sync_playwright().start()
self.browser = self.playwright.chromium.launch(headless=True)
def spider_closed(self, spider):
"""爬虫关闭时清理资源"""
if self.browser:
self.browser.close()
if self.playwright:
self.playwright.stop()
def process_request(self, request, spider):
"""使用Playwright渲染页面"""
# 只对需要渲染的请求启用
if not request.meta.get('render'):
return None
page = self.browser.new_page()
try:
page.goto(request.url, wait_until='networkidle')
# 等待特定元素出现
if request.meta.get('wait_for'):
page.wait_for_selector(request.meta['wait_for'])
# 获取渲染后的HTML
html = page.content()
# 返回HtmlResponse,后续解析流程不变
return HtmlResponse(
url=request.url,
body=html.encode('utf-8'),
encoding='utf-8',
request=request
)
finally:
page.close()
在爬虫中使用:
def parse(self, response):
# 渲染后的解析逻辑
pass
def start_requests(self):
for url in self.start_urls:
yield scrapy.Request(
url,
callback=self.parse,
meta={
'render': True, # 启用渲染
'wait_for': '.product' # 等待元素出现
}
)
安装Playwright:
pip install playwright
playwright install chromium
中间件执行顺序
中间件的优先级(数字越小越先执行):
下载中间件执行顺序:
0-400: Scrapy内置中间件
400-700: 自定义中间件
700+: Scrapy内置中间件
请求流程(从大到小):
引擎 → 中间件(400) → 中间件(410) → 下载器
响应流程(从小到大):
下载器 → 中间件(410) → 中间件(400) → 引擎
高级功能二:管道高级用法
上篇我们用了单个Pipeline做数据清洗和保存。实际项目中,你可能需要:
• 同时保存到多种存储
• 批量处理提高性能
• 数据验证和补全
• 实时去重和报警
实战1:多存储Pipeline
同时保存到JSON、CSV、MySQL三种存储。
# pipelines.py
import json
import csv
import pymysql
from scrapy.exceptions import DropItem
class MultipleStoragePipeline:
"""多存储管道"""
def __init__(self):
self.json_file = None
self.csv_writer = None
self.csv_file = None
self.conn = None
self.cursor = None
self.items = []
@classmethod
def from_crawler(cls, crawler):
pipeline = cls()
crawler.signals.connect(pipeline.spider_opened, signals.spider_opened)
crawler.signals.connect(pipeline.spider_closed, signals.spider_closed)
return pipeline
def spider_opened(self, spider):
"""爬虫启动时初始化连接"""
# JSON
self.json_file = open('books_data.json', 'w', encoding='utf-8')
self.json_file.write('[')
# CSV
self.csv_file = open('books_data.csv', 'w', encoding='utf-8', newline='')
self.csv_writer = csv.writer(self.csv_file)
self.csv_writer.writerow(['title', 'price', 'rating', 'stock', 'url', 'crawl_time'])
# MySQL
self.conn = pymysql.connect(
host='localhost',
user='root',
password='password',
database='scrapy_db'
)
self.cursor = self.conn.cursor()
# 创建表
self.cursor.execute('''
CREATE TABLE IF NOT EXISTS books (
id INT AUTO_INCREMENT PRIMARY KEY,
title VARCHAR(255),
price DECIMAL(10,2),
rating VARCHAR(50),
stock VARCHAR(100),
url TEXT,
crawl_time DATETIME
)
''')
def process_item(self, item, spider):
"""处理每条数据"""
# 数据验证
if not item.get('title'):
raise DropItem("Missing title")
# 添加到列表
self.items.append(item)
# JSON: 逐条写入(逗号分隔)
json_str = json.dumps(dict(item), ensure_ascii=False)
self.json_file.write('\n ' + json_str + ',')
# CSV: 逐行写入
self.csv_writer.writerow([
item['title'],
item['price'],
item['rating'],
item['stock'],
item['url'],
item['crawl_time']
])
# MySQL: 插入数据
self.cursor.execute('''
INSERT INTO books (title, price, rating, stock, url, crawl_time)
VALUES (%s, %s, %s, %s, %s, %s)
''', (
item['title'],
item['price'],
item['rating'],
item['stock'],
item['url'],
item['crawl_time']
))
self.conn.commit()
spider.logger.info(f"Saved: {item['title']}")
return item
def spider_closed(self, spider):
"""爬虫关闭时清理"""
# JSON: 修正最后的逗号并闭合数组
self.json_file.write('\n]')
self.json_file.close()
# CSV: 关闭文件
self.csv_file.close()
# MySQL: 关闭连接
self.cursor.close()
self.conn.close()
spider.logger.info(f"Total saved: {len(self.items)} items")
配置:
ITEM_PIPELINES = {
'book_scraper.pipelines.MultipleStoragePipeline': 300,
}
实战2:批量处理Pipeline
逐条保存到数据库效率太低,改为批量插入。
# pipelines.py
import pymysql
from scrapy.exceptions import DropItem
class BatchMySQLPipeline:
"""批量插入MySQL管道"""
def __init__(self, batch_size=100):
self.batch_size = batch_size
self.conn = None
self.cursor = None
self.batch = []
@classmethod
def from_crawler(cls, crawler):
return cls(
batch_size=crawler.settings.getint('BATCH_SIZE', 100)
)
def open_spider(self, spider):
"""爬虫启动时初始化连接"""
self.conn = pymysql.connect(
host='localhost',
user='root',
password='password',
database='scrapy_db'
)
self.cursor = self.conn.cursor()
def process_item(self, item, spider):
"""暂存数据,达到批量大小后插入"""
if not item.get('title'):
raise DropItem("Missing title")
self.batch.append(item)
# 达到批量大小时插入
if len(self.batch) >= self.batch_size:
self._insert_batch()
return item
def close_spider(self, spider):
"""爬虫关闭时插入剩余数据"""
if self.batch:
self._insert_batch()
self.cursor.close()
self.conn.close()
def _insert_batch(self):
"""批量插入数据"""
if not self.batch:
return
sql = '''
INSERT INTO books (title, price, rating, stock, url, crawl_time)
VALUES (%s, %s, %s, %s, %s, %s)
'''
values = [
(item['title'], item['price'], item['rating'],
item['stock'], item['url'], item['crawl_time'])
for item in self.batch
]
self.cursor.executemany(sql, values)
self.conn.commit()
self.batch = [] # 清空批次
配置:
ITEM_PIPELINES = {
'book_scraper.pipelines.BatchMySQLPipeline': 300,
}
BATCH_SIZE = 100 # 每批100条
实战3:数据验证Pipeline
确保数据完整性和质量。
# pipelines.py
from scrapy.exceptions import DropItem
class DataValidationPipeline:
"""数据验证管道"""
def process_item(self, item, spider):
"""验证数据质量"""
errors = []
# 必填字段检查
required_fields = ['title', 'price', 'rating']
for field in required_fields:
if not item.get(field):
errors.append(f"Missing field: {field}")
# 价格范围检查
try:
price = float(item.get('price', 0))
if price < 0 or price > 10000:
errors.append(f"Invalid price: {price}")
except ValueError:
errors.append(f"Invalid price format: {item.get('price')}")
# URL格式检查
url = item.get('url', '')
if not url.startswith('http'):
errors.append(f"Invalid URL: {url}")
# 有错误则丢弃数据
if errors:
spider.logger.warning(f"Validation failed for {item.get('title', 'unknown')}: {errors}")
raise DropItem(f"Validation failed: {errors}")
return item
Pipeline执行顺序
Pipeline按优先级(数字越小越先执行):
ITEM_PIPELINES = {
'book_scraper.pipelines.DataValidationPipeline': 100, # 先验证
'book_scraper.pipelines.BookPipeline': 200, # 再清洗
'book_scraper.pipelines.BatchMySQLPipeline': 300, # 最后存储
}
高级功能三:分布式爬虫(Scrapy-Redis)
单机爬虫受限于硬件性能,当数据量达到百万级时,需要分布式架构。Scrapy-Redis是基于Redis的分布式爬虫框架,可以轻松扩展到多台机器。
Scrapy-Redis的工作原理
传统Scrapy架构:
调度器 → 爬虫1, 爬虫2, 爬虫3 (每个爬虫独立调度)
Scrapy-Redis架构:
Redis调度器 → 爬虫1, 爬虫2, 爬虫3 (共享调度器)
核心改进:
• 共享调度器:所有爬虫实例共享一个Redis请求队列
• 共享去重:使用Redis集合去重,避免重复请求
• 增量爬取:支持断点续传,爬取过程中可随时停止恢复
安装Scrapy-Redis
pip install scrapy-redis
# 启动Redis
redis-server
创建分布式爬虫
修改原有的爬虫,继承RedisSpider:
# spiders/books_distributed.py
import scrapy
from scrapy_redis.spiders import RedisSpider
from book_scraper.items import BookItem
class BooksDistributedSpider(RedisSpider):
"""分布式商品爬虫"""
name = 'books_distributed'
redis_key = 'books:start_urls' # Redis中存储起始URL的key
# 原有的allowed_domains在分布式爬虫中不需要指定
# allowed_domains = ['books.toscrape.com']
def parse(self, response):
"""解析列表页"""
books = response.css('article.product_pod')
for book in books:
item = BookItem()
item['title'] = book.css('h3 a::attr(title)').get()
item['price'] = book.css('p.price_color::text').get().replace('£', '')
item['rating'] = book.css('p.star-rating::attr(class)').get().split()[-1]
item['stock'] = book.css('p.instock availability::text').re_first(r'\w+.*').strip()
item['url'] = response.urljoin(book.css('h3 a::attr(href)').get())
import datetime
item['crawl_time'] = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')
yield item
# 处理下一页
next_page = response.css('li.next a::attr(href)').get()
if next_page:
yield scrapy.Request(response.urljoin(next_page), callback=self.parse)
修改settings.py
# 启用Scrapy-Redis调度器
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# 启用Redis去重
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# 不清除Redis队列,支持断点续传
SCHEDULER_PERSIST = True
# Redis配置
REDIS_HOST = 'localhost'
REDIS_PORT = 6379
# 其他配置保持不变
ITEM_PIPELINES = {
'book_scraper.pipelines.BookPipeline': 300,
}
启动分布式爬虫
方式1:手动推送起始URL
# 在Redis中添加起始URL
redis-cli
> LPUSH books:start_urls "http://books.toscrape.com/"
# 启动爬虫(可以启动多个实例)
scrapy crawl books_distributed
方式2:使用RedisSpider的自动模式
爬虫启动后会自动从books:start_urls列表中获取起始URL。
启动多个实例
在多台机器上同时运行:
# 机器1
scrapy crawl books_distributed
# 机器2
scrapy crawl books_distributed
# 机器3
scrapy crawl books_distributed
所有实例会共享Redis中的请求队列,协同工作。
监控Redis状态
# 查看请求队列长度
redis-cli
> LLEN books:requests
# 查看已处理的URL数量
> SCARD books:dupefilter
# 查看当前正在爬取的URL
> LINDEX books:requests 0
Scrapy-Redis vs 传统Scrapy
|
特性 |
传统Scrapy |
Scrapy-Redis |
|
调度器 |
内存队列 |
Redis队列 |
|
去重 |
内存集合 |
Redis集合 |
|
扩展性 |
单机 |
多机分布式 |
|
断点续传 |
不支持 |
支持 |
|
性能 |
受限于单机 |
可横向扩展 |
|
部署复杂度 |
简单 |
需要Redis |
性能优化
当爬取规模达到十万、百万级时,性能优化变得至关重要。我总结了几个最有效的优化策略。
1. 并发调优
Scrapy默认的并发配置比较保守,需要根据实际情况调整。
# settings.py
# 全局并发数(所有域名)
CONCURRENT_REQUESTS = 32 # 默认16,可提高到64-128
# 单域名并发数
CONCURRENT_REQUESTS_PER_DOMAIN = 16 # 默认8,可提高到32-64
# 下载延迟(秒)
DOWNLOAD_DELAY = 0.5 # 根据目标网站承受能力调整,建议0.1-1.0
# 随机化延迟(避免规律性)
RANDOMIZE_DOWNLOAD_DELAY = True
# 超时设置
DOWNLOAD_TIMEOUT = 15 # 默认180秒,太快可能导致超时
并发数设置原则:
• CPU密集型爬虫:并发数 = CPU核心数 × 2
• IO密集型爬虫:并发数可以更高,64-128都没问题
• 具体数值需要通过压力测试确定
2. 内存优化
Scrapy默认所有请求都在内存中,爬取大量URL时容易内存溢出。
# settings.py
# 限制内存中请求的数量
CONCURRENT_REQUESTS = 32
# 禁用调试模式(减少内存占用)
DEBUG = False
# 使用更高效的数据结构(需要自定义中间件)
3. 请求队列优化
Scrapy-Redis的请求队列可以优化:
# 使用更高效的队列算法
SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.PriorityQueue' # 优先级队列
# SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.FifoQueue' # FIFO队列
# SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.LifoQueue' # LIFO队列
4. 数据库优化
批量插入比逐条插入快10-20倍:
# pipelines.py
# 详见上面的BatchMySQLPipeline示例
5. 禁用不必要的功能
# settings.py
# 禁用Telnet控制台
TELNETCONSOLE_ENABLED = False
# 禁用统计信息收集(不推荐,但可以提升性能)
STATS_CLASS = None
# 减少日志输出
LOG_LEVEL = 'INFO'
6. 使用连接池
# settings.py
# 下载器连接池大小
REACTOR_THREADPOOL_MAXSIZE = 20
# 数据库连接池
DBPOOL_SIZE = 10
反爬策略应对
反爬是爬虫开发的永恒话题。这里总结几个实用的反爬应对方案。
1. User-Agent轮换
我们已经实现了RandomUserAgentMiddleware,这里补充一个更高级的版本:
# middlewares.py
import random
from fake_useragent import UserAgent
class AdvancedUserAgentMiddleware:
"""高级UA中间件(使用fake_useragent库)"""
def __init__(self):
self.ua = UserAgent()
def process_request(self, request, spider):
request.headers['User-Agent'] = self.ua.random
spider.logger.debug(f"User-Agent: {request.headers['User-Agent']}")
安装fake-useragent:
pip install fake-useragent
2. IP代理池
我们已经实现了基础版本,实际项目中建议使用专业代理服务(如快代理、芝麻代理等):
# middlewares.py
import requests
class DynamicProxyMiddleware:
"""动态代理池中间件(从API获取)"""
def __init__(self, proxy_api):
self.proxy_api = proxy_api
self.proxies = []
self._fetch_proxies()
def _fetch_proxies(self):
"""从API获取代理列表"""
try:
response = requests.get(self.proxy_api)
self.proxies = response.json().get('proxies', [])
except Exception as e:
self.proxies = []
def process_request(self, request, spider):
"""使用代理"""
if self.proxies:
proxy = random.choice(self.proxies)
request.meta['proxy'] = proxy['url']
spider.logger.debug(f"Using proxy: {proxy['url']}")
3. 请求频率控制
通过时间窗口限制请求频率:
# middlewares.py
import time
from collections import defaultdict
class RateLimitMiddleware:
"""频率限制中间件"""
def __init__(self, max_requests=10, window=60):
self.max_requests = max_requests # 时间窗口内最大请求数
self.window = window # 时间窗口(秒)
self.request_times = defaultdict(list)
def process_request(self, request, spider):
"""检查请求频率"""
domain = request.url.split('/')[2]
now = time.time()
# 移除过期的请求记录
self.request_times[domain] = [
t for t in self.request_times[domain] if now - t < self.window
]
# 检查是否超限
if len(self.request_times[domain]) >= self.max_requests:
# 等待直到有名额
wait_time = self.window - (now - self.request_times[domain][0])
if wait_time > 0:
time.sleep(wait_time)
# 记录本次请求
self.request_times[domain].append(now)
4. Cookie池
有些网站需要登录才能访问数据,需要维护Cookie池。
# middlewares.py
import random
import json
class CookiePoolMiddleware:
"""Cookie池中间件"""
def __init__(self, cookie_file):
with open(cookie_file, 'r', encoding='utf-8') as f:
self.cookies = json.load(f)
def process_request(self, request, spider):
"""随机分配Cookie"""
if self.cookies:
cookie = random.choice(self.cookies)
request.cookies = cookie
spider.logger.debug(f"Using cookie")
Cookie文件格式:
[
{
"sessionid": "abc123",
"csrftoken": "xyz456"
},
{
"sessionid": "def789",
"csrftoken": "uvw012"
}
]
5. 行为模拟
模拟真实用户的浏览行为:
# middlewares.py
import random
class HumanBehaviorMiddleware:
"""模拟人类行为中间件"""
def process_request(self, request, spider):
"""随机延迟"""
if request.meta.get('human_behavior'):
delay = random.uniform(1, 3)
time.sleep(delay)
6. 验证码处理
遇到验证码时,可以使用第三方服务(如超级鹰、打码兔):
# middlewares.py
import requests
class CaptchaMiddleware:
"""验证码处理中间件"""
def __init__(self, captcha_api):
self.captcha_api = captcha_api
def process_request(self, request, spider):
"""遇到验证码时自动识别"""
if request.meta.get('captcha'):
captcha_image = self._get_captcha_image(request)
captcha_code = self._recognize_captcha(captcha_image)
request.meta['captcha_code'] = captcha_code
监控与告警
长期运行的爬虫需要完善的监控和告警系统。
1. 内置统计信息
Scrapy自带统计功能:
# 查看统计信息
scrapy crawl books --stats
# 保存统计信息到文件
scrapy crawl books --statsfile=stats.json
2. 自定义监控中间件
# middlewares.py
import time
from scrapy import signals
class MonitorMiddleware:
"""监控中间件"""
def __init__(self):
self.start_time = None
self.request_count = 0
self.error_count = 0
self.item_count = 0
@classmethod
def from_crawler(cls, crawler):
middleware = cls()
crawler.signals.connect(middleware.spider_opened, signals.spider_opened)
crawler.signals.connect(middleware.spider_closed, signals.spider_closed)
crawler.signals.connect(middleware.item_scraped, signals.item_scraped)
crawler.signals.connect(middleware.spider_error, signals.spider_error)
return middleware
def spider_opened(self, spider):
"""爬虫启动时记录时间"""
self.start_time = time.time()
spider.logger.info(f"Spider started: {spider.name}")
def spider_closed(self, spider):
"""爬虫关闭时输出统计"""
duration = time.time() - self.start_time
spider.logger.info(f"""
Spider finished: {spider.name}
Duration: {duration:.2f}s
Requests: {self.request_count}
Items: {self.item_count}
Errors: {self.error_count}
Speed: {self.item_count / duration:.2f} items/s
""")
def item_scraped(self, item, response, spider):
"""每爬取一个item计数"""
self.item_count += 1
def spider_error(self, failure, response, spider):
"""每次错误计数"""
self.error_count += 1
3. 告警通知
# pipelines.py
import requests
class AlertPipeline:
"""告警管道"""
def __init__(self, webhook_url):
self.webhook_url = webhook_url
def process_item(self, item, spider):
"""每爬取1000条发送通知"""
if spider.crawler.stats.get_value('item_scraped_count', 0) % 1000 == 0:
self._send_alert(f"已爬取 {spider.crawler.stats.get_value('item_scraped_count')} 条数据")
return item
def _send_alert(self, message):
"""发送告警(钉钉/企业微信)"""
data = {
"msgtype": "text",
"text": {
"content": message
}
}
requests.post(self.webhook_url, json=data)
4. 日志集中管理
使用ELK(Elasticsearch + Logstash + Kibana)或Sentry进行日志管理。
常见问题解决
问题1:内存溢出
现象:爬取一段时间后程序崩溃,提示"MemoryError"
原因:
• 请求队列太大
• Pipeline缓存了太多数据
• 未正确清理资源
解决方案:
# 1. 限制并发数
CONCURRENT_REQUESTS = 16
# 2. 分批处理数据
# 详见BatchMySQLPipeline
# 3. 及时清理资源
def close_spider(self, spider):
self.items.clear()
self.cursor.close()
问题2:反爬严重
现象:频繁403/404,IP被封
解决方案:
• 使用代理IP池
• 降低请求频率
• 随机化UA和Cookie
• 模拟人类行为
问题3:数据重复
现象:数据库中有大量重复数据
解决方案:
• 使用Pipeline去重
• 在数据库表加唯一索引
• 启用Scrapy的URL去重(Redis)
问题4:速度太慢
现象:爬取速度远低于预期
解决方案:
• 提高并发数
• 降低下载延迟
• 使用批量插入
• 检查网络连接
完整实战案例:分布式电商爬虫
现在我们把前面的知识点整合起来,做一个完整的分布式电商爬虫。
项目结构
ecommerce_crawler/
├── scrapy.cfg
├── requirements.txt
└── ecommerce_crawler/
├── __init__.py
├── items.py
├── middlewares.py
├── pipelines.py
├── settings.py
└── spiders/
├── jd_spider.py
├── taobao_spider.py
└── pinduoduo_spider.py
requirements.txt
scrapy==2.11.0
scrapy-redis==0.7.3
pymysql==1.1.0
fake-useragent==1.4.0
playwright==1.40.0
redis==5.0.1
items.py
import scrapy
class ProductItem(scrapy.Item):
"""商品数据模型"""
title = scrapy.Field()
price = scrapy.Field()
sales = scrapy.Field()
rating = scrapy.Field()
shop = scrapy.Field()
url = scrapy.Field()
image_url = scrapy.Field()
platform = scrapy.Field() # JD/Taobao/Pinduoduo
crawl_time = scrapy.Field()
middlewares.py
整合所有中间件:
from fake_useragent import UserAgent
import random
class ComprehensiveMiddleware:
"""综合中间件:UA+代理+限速+动态渲染"""
def __init__(self, proxy_list):
self.ua = UserAgent()
self.proxy_list = proxy_list
@classmethod
def from_crawler(cls, crawler):
return cls(
proxy_list=crawler.settings.get('PROXY_LIST', [])
)
def process_request(self, request, spider):
# 随机UA
request.headers['User-Agent'] = self.ua.random
# 随机代理
if self.proxy_list:
request.meta['proxy'] = random.choice(self.proxy_list)
# 动态渲染(京东/淘宝需要)
if request.meta.get('render'):
request.meta['download_timeout'] = 30
pipelines.py
多存储+批量插入+验证+告警:
import pymysql
from scrapy.exceptions import DropItem
class ProductPipeline:
"""综合管道:验证+批量插入+告警"""
def __init__(self):
self.batch = []
self.batch_size = 100
self.conn = None
self.cursor = None
def open_spider(self, spider):
self.conn = pymysql.connect(host='localhost', user='root', password='password', database='ecommerce')
self.cursor = self.conn.cursor()
def process_item(self, item, spider):
# 验证
if not item.get('title') or not item.get('price'):
raise DropItem("Invalid item")
# 批量处理
self.batch.append(item)
if len(self.batch) >= self.batch_size:
self._insert_batch()
# 告警(每1000条)
if spider.crawler.stats.get_value('item_scraped_count', 0) % 1000 == 0:
spider.logger.info(f"已爬取 {spider.crawler.stats.get_value('item_scraped_count')} 条数据")
return item
def close_spider(self, spider):
if self.batch:
self._insert_batch()
self.cursor.close()
self.conn.close()
def _insert_batch(self):
sql = 'INSERT INTO products (title, price, sales, rating, shop, url, platform) VALUES (%s, %s, %s, %s, %s, %s, %s)'
values = [(item['title'], item['price'], item.get('sales'), item.get('rating'), item.get('shop'), item['url'], item['platform']) for item in self.batch]
self.cursor.executemany(sql, values)
self.conn.commit()
self.batch.clear()
settings.py
# 基础配置
BOT_NAME = 'ecommerce_crawler'
SPIDER_MODULES = ['ecommerce_crawler.spiders']
NEWSPIDER_MODULE = 'ecommerce_crawler.spiders'
# Scrapy-Redis配置
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
SCHEDULER_PERSIST = True
# 并发配置
CONCURRENT_REQUESTS = 64
CONCURRENT_REQUESTS_PER_DOMAIN = 32
DOWNLOAD_DELAY = 0.5
# 中间件
DOWNLOADER_MIDDLEWARES = {
'ecommerce_crawler.middlewares.ComprehensiveMiddleware': 400,
}
# 管道
ITEM_PIPELINES = {
'ecommerce_crawler.pipelines.ProductPipeline': 300,
}
# Redis
REDIS_HOST = 'localhost'
REDIS_PORT = 6379
# 代理列表(示例)
PROXY_LIST = [
'http://proxy1.example.com:8080',
'http://proxy2.example.com:8080',
]
JD爬虫示例
import scrapy
from scrapy_redis.spiders import RedisSpider
from ecommerce_crawler.items import ProductItem
class JDSpider(RedisSpider):
name = 'jd'
redis_key = 'jd:start_urls'
def parse(self, response):
"""解析京东商品列表"""
products = response.css('.gl-item')
for product in products:
item = ProductItem()
item['title'] = product.css('.p-name em::text').get()
item['price'] = product.css('.p-price i::text').get()
item['sales'] = product.css('.p-commit::text').get()
item['url'] = response.urljoin(product.css('.p-name a::attr(href)').get())
item['platform'] = 'JD'
import datetime
item['crawl_time'] = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')
yield item
# 下一页
next_page = response.css('.pn-next::attr(href)').get()
if next_page:
yield scrapy.Request(response.urljoin(next_page), callback=self.parse)
启动分布式爬虫
# 启动Redis
redis-server
# 推送起始URL
redis-cli LPUSH jd:start_urls "https://search.jd.com/Search?keyword=手机"
# 启动爬虫(多台机器同时运行)
scrapy crawl jd
学习资源与扩展
官方文档
• Scrapy官方文档: https://docs.scrapy.org/
• Scrapy-Redis文档: https://github.com/rmax/scrapy-redis
推荐书籍
• 《Python网络数据采集》Ryan Mitchell
• 《精通Scrapy网络爬虫》刘伟
实战项目
• GitHub上的优秀爬虫项目: https://github.com/topics/scrapy
• 各种电商爬虫开源项目
进阶方向
1. 机器学习:用ML识别验证码、预测反爬策略
2. 大数据:结合Spark处理海量数据
3. 流处理:使用Kafka进行实时数据处理
4. 可视化:构建爬虫监控Dashboard
小结
下篇到这里就结束了。我们深入讲解了:
• 中间件的高级用法(UA、代理、动态渲染)
• 管道的高级用法(多存储、批量处理、验证、告警)
• Scrapy-Redis分布式架构
• 性能优化策略
• 反爬应对方案
• 监控与告警系统
• 完整的分布式电商爬虫实战
总结
上下两篇完整覆盖了Scrapy从入门到精通的全流程:
上篇(入门):
• Scrapy的核心概念和优势
• 环境搭建和项目创建
• 基础爬虫开发流程
• 实战案例:电商商品爬取
• 调试技巧
下篇(进阶):
• 中间件开发(UA、代理、动态渲染)
• 管道高级用法(多存储、批量处理)
• Scrapy-Redis分布式爬虫
• 性能优化策略
• 反爬应对方案
• 监控与告警
• 完整的分布式爬虫实战
掌握了这些内容,你就可以独立开发百万级的爬虫系统了。记住,爬虫开发的核心不是技术,而是理解和解决问题的能力。
祝你在爬虫的道路上越走越远!
更多推荐
所有评论(0)