定时爬虫与数据采集系统:Scrapy + Redis去重 + 香港VPS分布式部署实战
香港服务器访问境外网站无需特殊配置,网络延迟低,是部署数据采集系统的优质平台。本文基于Scrapy框架,给出从单机爬虫到分布式爬虫的完整工程化方案。
一、Scrapy项目结构
<code">myspider/ ├── myspider/ │ ├── spiders/ │ │ ├── product_spider.py # 商品数据爬虫 │ │ ├── news_spider.py # 新闻爬虫 │ │ └── price_monitor.py # 价格监控爬虫 │ ├── middlewares.py # 中间件(代理/UA轮换) │ ├── pipelines.py # 数据处理管道 │ ├── settings.py # 全局配置 │ └── items.py # 数据结构定义 ├── requirements.txt └── docker-compose.yml
二、安装与项目初始化
<code"># 安装依赖
pip install scrapy scrapy-redis redis pymysql \
fake-useragent scrapy-user-agents \
itemadapter --break-system-packages
# 创建Scrapy项目
scrapy startproject myspider
cd myspider
# 创建第一个爬虫
scrapy genspider product_spider target-site.com三、爬虫示例:商品价格监控
<code"># myspider/spiders/price_monitor.py
import scrapy
from scrapy_redis.spiders import RedisSpider
from myspider.items import ProductItem
class PriceMonitorSpider(RedisSpider):
"""
分布式价格监控爬虫
从Redis队列读取待爬URL,支持多节点并发
"""
name = 'price_monitor'
redis_key = 'price_monitor:start_urls' # Redis队列键名
custom_settings = {
'DOWNLOAD_DELAY': 2, # 请求间隔2秒(礼貌爬取)
'RANDOMIZE_DOWNLOAD_DELAY': True,
'CONCURRENT_REQUESTS_PER_DOMAIN': 2,
}
def parse(self, response):
item = ProductItem()
item['url'] = response.url
item['name'] = response.css('h1.product-title::text').get('').strip()
item['price'] = self.extract_price(response)
item['stock'] = response.css('.stock-status::text').get('')
item['crawled_at'] = response.headers.get('Date', b'').decode()
yield item
# 跟进分页
next_page = response.css('a.next-page::attr(href)').get()
if next_page:
yield response.follow(next_page, self.parse)
def extract_price(self, response) -> float:
price_str = response.css('span.price::text').get('0')
# 清理价格字符串:$1,299.00 → 1299.00
import re
price_clean = re.sub(r'[^\d.]', '', price_str)
return float(price_clean) if price_clean else 0.0四、核心中间件:代理与UA轮换
<code"># myspider/middlewares.py
import random
import time
from fake_useragent import UserAgent
ua = UserAgent()
class RotateUserAgentMiddleware:
"""随机轮换User-Agent,降低被识别风险"""
def process_request(self, request, spider):
request.headers['User-Agent'] = ua.random
return None
class ProxyMiddleware:
"""代理IP池中间件"""
def __init__(self, proxy_list: list):
self.proxies = proxy_list
self.failed_proxies = set()
@classmethod
def from_crawler(cls, crawler):
# 从设置中读取代理列表
proxies = crawler.settings.getlist('PROXY_LIST', [])
return cls(proxies)
def process_request(self, request, spider):
if not self.proxies:
return None
# 过滤掉失败的代理
available = [p for p in self.proxies if p not in self.failed_proxies]
if not available:
self.failed_proxies.clear() # 全部失败时重置
available = self.proxies
proxy = random.choice(available)
request.meta['proxy'] = proxy
request.meta['proxy_used'] = proxy
return None
def process_exception(self, request, exception, spider):
proxy = request.meta.get('proxy_used')
if proxy:
self.failed_proxies.add(proxy)
spider.logger.warning(f'代理失败,已标记: {proxy}')
return None
class RetryWithDelayMiddleware:
"""智能重试:指数退避,避免频繁触发反爬"""
MAX_RETRY = 3
def process_response(self, request, response, spider):
if response.status in [429, 503, 403]:
retry_count = request.meta.get('retry_count', 0)
if retry_count < self.MAX_RETRY:
wait_time = 2 ** retry_count * 5 # 指数退避:5s, 10s, 20s
spider.logger.info(f'触发限流,等待{wait_time}秒后重试')
time.sleep(wait_time)
new_request = request.copy()
new_request.meta['retry_count'] = retry_count + 1
new_request.dont_filter = True
return new_request
return response五、数据管道:存储到MySQL + 价格告警
<code"># myspider/pipelines.py
import pymysql
import redis
import json
from itemadapter import ItemAdapter
from datetime import datetime
class MySQLPipeline:
"""将爬取数据存入MySQL,并检测价格变化"""
def __init__(self, mysql_config, redis_config):
self.mysql_config = mysql_config
self.redis_client = redis.Redis(**redis_config)
@classmethod
def from_crawler(cls, crawler):
return cls(
mysql_config=crawler.settings.get('MYSQL_CONFIG'),
redis_config=crawler.settings.get('REDIS_CONFIG'),
)
def open_spider(self, spider):
self.conn = pymysql.connect(**self.mysql_config)
self.cursor = self.conn.cursor()
def close_spider(self, spider):
self.conn.commit()
self.conn.close()
def process_item(self, item, spider):
adapter = ItemAdapter(item)
# 查询上次价格
url = adapter['url']
prev_key = f"prev_price:{url}"
prev_price = self.redis_client.get(prev_key)
current_price = float(adapter.get('price', 0))
# 价格变化检测
if prev_price is not None:
prev_price_float = float(prev_price)
change_pct = (current_price - prev_price_float) / prev_price_float * 100
if abs(change_pct) >= 5: # 价格变动超过5%
self.send_price_alert(adapter, prev_price_float, current_price, change_pct)
# 更新Redis中的价格缓存
self.redis_client.set(prev_key, str(current_price), ex=86400 * 7)
# 写入数据库
sql = """
INSERT INTO product_prices (url, name, price, stock, crawled_at)
VALUES (%s, %s, %s, %s, %s)
ON DUPLICATE KEY UPDATE
name=VALUES(name), price=VALUES(price),
stock=VALUES(stock), crawled_at=VALUES(crawled_at)
"""
self.cursor.execute(sql, (
url, adapter.get('name'), current_price,
adapter.get('stock'), adapter.get('crawled_at'),
))
self.conn.commit()
return item
def send_price_alert(self, item, prev, current, pct):
direction = "📉 降价" if pct < 0 else "📈 涨价"
msg = f"{direction} {abs(pct):.1f}%\n商品:{item['name']}\n旧价:{prev}\n新价:{current}"
# 发送Telegram告警(参见第79篇)
self.redis_client.lpush('alerts:queue', json.dumps({'msg': msg}))六、settings.py核心配置
<code"># myspider/settings.py
BOT_NAME = 'myspider'
SPIDER_MODULES = ['myspider.spiders']
# 遵守robots.txt
ROBOTSTXT_OBEY = True
# 并发控制
CONCURRENT_REQUESTS = 16
DOWNLOAD_DELAY = 1
AUTOTHROTTLE_ENABLED = True # 自动限速(根据服务器响应调整)
AUTOTHROTTLE_START_DELAY = 1
AUTOTHROTTLE_MAX_DELAY = 10
AUTOTHROTTLE_TARGET_CONCURRENCY = 2.0
# 中间件
DOWNLOADER_MIDDLEWARES = {
'myspider.middlewares.RotateUserAgentMiddleware': 400,
'myspider.middlewares.ProxyMiddleware': 410,
'myspider.middlewares.RetryWithDelayMiddleware': 550,
'scrapy.downloadermiddlewares.retry.RetryMiddleware': None, # 禁用默认重试
}
# 数据管道
ITEM_PIPELINES = {
'myspider.pipelines.MySQLPipeline': 300,
}
# Scrapy-Redis分布式配置
SCHEDULER = 'scrapy_redis.scheduler.Scheduler'
DUPEFILTER_CLASS = 'scrapy_redis.dupefilter.RFPDupeFilter' # Redis去重
SCHEDULER_PERSIST = True # 爬虫关闭后不清空队列(断点续爬)
REDIS_HOST = '127.0.0.1'
REDIS_PORT = 6379
REDIS_PARAMS = {'password': '你的Redis密码', 'db': 2}
# MySQL配置
MYSQL_CONFIG = {
'host': 'localhost', 'port': 3306,
'user': 'spider', 'password': '数据库密码',
'db': 'spider_data', 'charset': 'utf8mb4',
}
# 请求缓存(开发调试用,防止重复请求)
# HTTPCACHE_ENABLED = True
# HTTPCACHE_EXPIRATION_SECS = 3600七、定时调度:Crontab + Scrapyd
<code"># 方式1:简单Crontab定时运行
# 每天凌晨3点爬取商品价格
0 3 * * * cd /srv/myspider && /srv/venv/bin/scrapy crawl price_monitor >> /var/log/scrapy/price.log 2>&1
# 方式2:Scrapyd(推荐,支持远程部署和状态查询)
pip install scrapyd scrapyd-client --break-system-packages
# 启动scrapyd服务
scrapyd &
# 部署爬虫到Scrapyd
scrapyd-deploy default
# 通过API调度爬虫
curl http://localhost:6800/schedule.json \
-d project=myspider \
-d spider=price_monitor
# 查看爬虫运行状态
curl http://localhost:6800/listjobs.json?project=myspider八、向Redis队列添加初始URL
<code"># 通过Python脚本批量添加待爬URL
import redis
r = redis.Redis(host='127.0.0.1', port=6379, password='你的密码', db=2)
urls = [
'https://target-site.com/products/page/1',
'https://target-site.com/products/page/2',
# ... 更多URL
]
for url in urls:
r.lpush('price_monitor:start_urls', url)
print(f"已添加 {len(urls)} 个URL到爬虫队列")九、总结
Scrapy + Scrapy-Redis的分布式方案将URL去重和队列管理统一到Redis,多台爬虫节点共享同一个队列,横向扩展非常方便。香港服务器访问境外目标站点无需代理,大幅降低爬虫的维护成本。IDC.Net的香港VPS网络稳定、IP质量好,是运行数据采集系统的优质平台,多台VPS可以组成高效的分布式爬虫集群。
版权声明:
作者:后浪云
链接:https://idc.net/help/442858/
文章版权归作者所有,未经允许请勿转载。
THE END
