定时爬虫与数据采集系统: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可以组成高效的分布式爬虫集群。

THE END