在数据驱动的业务场景中,爬虫是获取外部数据的核心手段,而数据入仓则是实现数据价值的关键环节。Scrapy 作为 Python 生态中成熟的爬虫框架,能高效抓取各类网页数据;ClickHouse 作为列式存储数据库,凭借超高的查询性能和写入吞吐量,成为实时数据分析场景的优选。本文将从实战角度出发,详细讲解如何搭建 “Scrapy 爬虫 + ClickHouse 数据仓” 的实时数据链路,解决数据抓取后即时入仓的核心问题。

一、技术选型背景:为什么选择 Scrapy+ClickHouse?

在开始实战前,我们需要先明确技术选型的合理性 —— 不同工具的组合需匹配业务对 “抓取效率” 和 “数据处理效率” 的双重需求。

1. Scrapy 的核心优势

  • 高效抓取能力:内置异步下载器,支持多线程并发抓取,可通过调整CONCURRENT_REQUESTS等参数优化吞吐量,单爬虫实例日均可抓取百万级数据。
  • 灵活的中间件机制:提供下载中间件(处理请求头、代理)、爬虫中间件(过滤重复请求),可快速集成反反爬策略(如 UA 池、IP 池)。
  • 结构化数据处理:通过ItemItem Pipeline规范数据格式,支持自动去重、数据清洗,避免脏数据流入下游。

2. ClickHouse 的核心优势

  • 超高写入吞吐量:采用列式存储 + 批量写入机制,单节点写入速度可达每秒数十万行,完全适配爬虫的高并发数据输出场景。
  • 实时查询性能:针对聚合查询(如COUNTSUMGROUP BY)优化,毫秒级响应分析需求,无需等待数据 “离线批处理”。
  • 轻量易部署:无复杂依赖,支持单机部署和集群扩展,中小规模数据场景下无需搭建复杂的分布式架构。

3. 链路匹配性

Scrapy 的Item Pipeline天然支持 “数据输出扩展”,可在 Pipeline 中直接集成 ClickHouse 写入逻辑;而 ClickHouse 的HTTP接口Python客户端(如clickhouse-driver),能无缝对接 Scrapy 的 Python 运行环境,无需额外中间件(如 Kafka)即可实现 “抓取 - 入仓” 实时联动,降低技术栈复杂度。

二、实战准备:环境搭建与依赖安装

在编写代码前,需完成本地环境配置,确保 Scrapy 能正常运行、ClickHouse 能正常连接。

1. 环境要求

  • 操作系统:Windows 10/11(需安装 WSL2)、Linux(Ubuntu 20.04+)、macOS(12+)
  • Python 版本:3.8-3.11(Scrapy 对 Python 3.12 + 兼容性待优化)
  • ClickHouse 版本:22.3+(推荐 LTS 版本,稳定性更高)

2. 依赖安装

(1)Scrapy 及相关依赖

通过 pip 安装 Scrapy 框架,以及数据处理所需的库:

bash

pip install scrapy  # 核心爬虫框架
pip install pandas  # 可选,用于数据清洗
pip install requests # 可选,用于调用ClickHouse HTTP接口
(2)ClickHouse 客户端

Python 操作 ClickHouse 主要依赖clickhouse-driver(原生 TCP 连接,性能更优):

bash

pip install clickhouse-driver
(3)ClickHouse 服务部署(以 Linux 为例)

若本地无 ClickHouse 服务,可通过 Docker 快速部署(推荐新手使用):

bash

# 拉取ClickHouse镜像
docker pull yandex/clickhouse-server:22.3

# 启动容器(映射端口8123:HTTP接口,9000:TCP接口)
docker run -d --name clickhouse-server -p 8123:8123 -p 9000:9000 yandex/clickhouse-server:22.3

# 进入容器,测试服务是否正常
docker exec -it clickhouse-server clickhouse-client
# 若进入命令行(显示clickhouse-client>),则服务启动成功

三、核心实战步骤:从爬虫搭建到数据入仓

本节将分 3 个关键步骤实现 “数据抓取 - 实时入仓”:①创建 Scrapy 爬虫项目;②设计 ClickHouse 数据表;③编写实时写入 Pipeline。

步骤 1:创建 Scrapy 爬虫项目,定义数据结构

以 “抓取某电商平台商品列表” 为例(仅作演示,实际需遵守网站robots.txt协议),搭建基础爬虫。

(1)初始化 Scrapy 项目

bash

# 创建项目(项目名:ecommerce_spider)
scrapy startproject ecommerce_spider

# 进入项目目录
cd ecommerce_spider

# 创建爬虫(爬虫名:product_spider,目标域名:example.com)
scrapy genspider product_spider example.com
(2)定义 Item 数据结构

修改ecommerce_spider/items.py,定义需抓取的商品字段(如商品 ID、名称、价格、销量):

python

import scrapy

class EcommerceProductItem(scrapy.Item):
    # 商品唯一ID(主键)
    product_id = scrapy.Field()
    # 商品名称
    product_name = scrapy.Field()
    # 商品价格(单位:元)
    price = scrapy.Field()
    # 商品销量(格式:100+)
    sales = scrapy.Field()
    # 抓取时间(自动填充,无需爬虫解析)
    crawl_time = scrapy.Field()
(3)编写爬虫解析逻辑

修改ecommerce_spider/spiders/product_spider.py,模拟商品列表页面解析(实际需替换为目标网站的真实 XPath/CSS 选择器):

python

import scrapy
from datetime import datetime
from ecommerce_spider.items import EcommerceProductItem

class ProductSpider(scrapy.Spider):
    name = 'product_spider'
    # 实际抓取时替换为目标网站的URL(如https://www.example.com/products?page=1)
    start_urls = ['https://example.com/product-list?page=1']

    def parse(self, response):
        # 1. 解析商品列表(假设每页10条商品)
        product_list = response.xpath('//div[@class="product-item"]')
        for product in product_list:
            item = EcommerceProductItem()
            # 提取商品字段(XPath需根据目标网站调整)
            item['product_id'] = product.xpath('./@data-id').get()
            item['product_name'] = product.xpath('./div[@class="name"]/text()').get().strip()
            item['price'] = product.xpath('./div[@class="price"]/span/text()').get().replace('¥', '')
            item['sales'] = product.xpath('./div[@class="sales"]/text()').get()
            # 自动填充抓取时间(格式:YYYY-MM-DD HH:MM:SS)
            item['crawl_time'] = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
            # 将Item传递给Pipeline处理(入仓逻辑在Pipeline中实现)
            yield item

        # 2. 分页抓取(模拟翻页,实际需解析下一页URL)
        next_page = response.xpath('//a[@class="next-page"]/@href').get()
        if next_page:
            yield response.follow(next_page, callback=self.parse)

步骤 2:设计 ClickHouse 数据表,适配爬虫数据

ClickHouse 的表结构需与 Scrapy Item 的字段一一对应,同时结合数据特性选择合适的字段类型和引擎。

(1)连接 ClickHouse 客户端

通过 Docker 或本地客户端连接 ClickHouse:

bash

# Docker容器内连接
docker exec -it clickhouse-server clickhouse-client

# 本地客户端连接(若已安装)
clickhouse-client
(2)创建数据库与数据表

sql

-- 1. 创建数据库(用于存储爬虫数据)
CREATE DATABASE IF NOT EXISTS ecommerce_spider_db;

-- 2. 切换数据库
USE ecommerce_spider_db;

-- 3. 创建商品数据表(核心参数说明)
CREATE TABLE IF NOT EXISTS product_data (
    product_id String COMMENT '商品唯一ID',
    product_name String COMMENT '商品名称',
    price Float32 COMMENT '商品价格(元)',
    sales String COMMENT '商品销量',
    crawl_time DateTime COMMENT '抓取时间'
) 
ENGINE = MergeTree()  -- ClickHouse默认引擎,支持高效查询
ORDER BY (product_id, crawl_time)  -- 排序键:按商品ID+抓取时间排序,优化查询
PARTITION BY toDate(crawl_time)  -- 分区键:按抓取日期分区,方便数据清理(如删除旧数据)
TTL toDateTime(crawl_time) + INTERVAL 30 DAY  -- 数据过期时间:保留30天数据,自动清理
COMMENT '电商商品爬虫数据表';
  • 字段类型选择product_idString(避免 ID 为字符串格式时丢失信息),priceFloat32(平衡精度与存储),crawl_timeDateTime(支持时间范围查询)。
  • 引擎选择MergeTree适合批量写入 + 实时查询场景,若需更高写入性能,可选择StripeLog(但查询性能略低)。

步骤 3:编写 ClickHouse Pipeline,实现实时写入

Scrapy 的Pipeline是数据处理的核心环节,我们将在 Pipeline 中集成 ClickHouse 写入逻辑,实现 “抓取一条、写入一条” 的实时效果(若需更高吞吐量,可改为批量写入)。

(1)编写 ClickHouse Pipeline 类

ecommerce_spider/pipelines.py中添加ClickHousePipeline

python

from clickhouse_driver import Client
from scrapy.exceptions import DropItem

class ClickHousePipeline:
    # 1. 初始化:连接ClickHouse(爬虫启动时执行一次)
    def open_spider(self, spider):
        # 配置ClickHouse连接参数(根据实际部署地址调整)
        self.client = Client(
            host='localhost',  # ClickHouse服务地址(Docker部署时填宿主机IP)
            port=9000,         # TCP端口(默认9000)
            database='ecommerce_spider_db',  # 目标数据库
            user='default',    # 默认用户名(无密码)
            password=''        # 默认密码为空
        )
        # 验证连接(可选)
        try:
            self.client.execute('SELECT 1')
            spider.logger.info('ClickHouse连接成功!')
        except Exception as e:
            spider.logger.error(f'ClickHouse连接失败:{str(e)}')
            raise e

    # 2. 数据处理:将Item写入ClickHouse(每条Item执行一次)
    def process_item(self, item, spider):
        try:
            # 1. 数据校验:过滤必填字段为空的Item(避免脏数据)
            if not item.get('product_id') or not item.get('product_name'):
                raise DropItem(f'缺失必填字段:{dict(item)}')
            
            # 2. 构造SQL插入语句(字段与ClickHouse表一致)
            insert_sql = """
                INSERT INTO product_data (product_id, product_name, price, sales, crawl_time)
                VALUES (%(product_id)s, %(product_name)s, %(price)s, %(sales)s, %(crawl_time)s)
            """
            
            # 3. 执行插入(将Item转换为字典,适配参数化查询)
            self.client.execute(insert_sql, dict(item))
            
            # 4. 日志记录(可选,便于调试)
            spider.logger.debug(f'数据写入成功:product_id={item["product_id"]}')
            
            return item  # 若后续有其他Pipeline,可继续传递Item
        
        except Exception as e:
            # 捕获异常,记录错误日志,丢弃异常Item
            spider.logger.error(f'数据写入失败:{str(e)},Item:{dict(item)}')
            raise DropItem(f'数据写入失败:{str(e)}')

    # 3. 关闭连接:爬虫停止时执行一次
    def close_spider(self, spider):
        self.client.disconnect()
        spider.logger.info('ClickHouse连接已关闭!')
(2)启用 Pipeline

修改项目配置文件ecommerce_spider/settings.py,启用ClickHousePipeline(注:Pipeline 优先级数字越小,执行越靠前):

python

# 启用Pipeline(替换默认的空列表)
ITEM_PIPELINES = {
    'ecommerce_spider.pipelines.ClickHousePipeline': 300,
}

# 可选:增加爬虫并发数(根据目标网站抗压能力调整)
CONCURRENT_REQUESTS = 16  # 默认16,最高可设为64(需避免触发反爬)
DOWNLOAD_DELAY = 0.5      # 每个请求延迟0.5秒,降低抓取压力

# 可选:设置User-Agent(模拟浏览器请求,避免被识别为爬虫)
USER_AGENT = 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/114.0.0.0 Safari/537.36'

四、测试与验证:确保数据实时入仓

完成代码编写后,需通过实际运行爬虫,验证数据是否成功写入 ClickHouse。

1. 启动 Scrapy 爬虫

在项目根目录执行以下命令,启动爬虫并查看日志:

bash

# 启动爬虫,输出日志到控制台(-L INFO:只显示INFO及以上级别日志)
scrapy crawl product_spider -L INFO

若日志中出现以下信息,说明爬虫正常运行且数据写入成功:

plaintext

[product_spider] INFO: ClickHouse连接成功!
[product_spider] DEBUG: 数据写入成功:product_id=123456
[product_spider] DEBUG: 数据写入成功:product_id=123457

2. 验证 ClickHouse 数据

连接 ClickHouse 客户端,执行查询语句,确认数据是否存在:

sql

-- 切换数据库
USE ecommerce_spider_db;

-- 查询最新写入的数据(前10条)
SELECT * FROM product_data ORDER BY crawl_time DESC LIMIT 10;

-- 统计总写入条数
SELECT COUNT(*) AS total_count FROM product_data;

若查询结果能显示爬虫抓取的商品数据,且crawl_time为当前时间附近,则说明 “实时入仓” 链路完全打通。

五、性能优化:应对高并发抓取场景

当爬虫抓取量达到 “每秒数十条” 以上时,单条写入 ClickHouse 可能出现性能瓶颈。可通过以下 3 个方向优化:

1. 批量写入 ClickHouse

修改ClickHousePipeline,将单条写入改为 “批量缓存 + 定时写入”,减少 SQL 执行次数:

python

class ClickHouseBatchPipeline:
    def open_spider(self, spider):
        self.client = Client(...)  # 同前
        self.batch_size = 50  # 每积累50条数据批量写入
        self.item_buffer = []  # 数据缓存列表

    def process_item(self, item, spider):
        self.item_buffer.append(dict(item))
        # 当缓存达到批量大小时,执行写入
        if len(self.item_buffer) >= self.batch_size:
            self._batch_insert()
            self.item_buffer = []  # 清空缓存
        return item

    # 批量插入逻辑
    def _batch_insert(self):
        insert_sql = """
            INSERT INTO product_data (product_id, product_name, price, sales, crawl_time)
            VALUES
        """
        # 构造批量参数(列表推导式,每个元素为一条数据的字典)
        self.client.execute(insert_sql, self.item_buffer)

    # 爬虫停止时,写入剩余缓存数据
    def close_spider(self, spider):
        if self.item_buffer:
            self._batch_insert()
        self.client.disconnect()

2. 调整 ClickHouse 写入参数

在 ClickHouse 客户端执行以下语句,优化写入性能(适用于集群环境,单机可省略):

sql

-- 临时关闭分区合并(写入期间减少IO消耗,写入后再开启)
SET optimize_on_insert = 0;

-- 调整批量写入大小(默认65536,可根据内存调整)
SET max_insert_block_size = 131072;

3. 启用 Scrapy 分布式爬虫

若单爬虫实例性能不足,可通过scrapy-redis实现分布式抓取,多实例同时向 ClickHouse 写入数据(需确保product_id唯一,避免重复数据):

bash

# 安装scrapy-redis
pip install scrapy-redis

修改settings.py,配置分布式参数:

python

# 启用Redis调度器(分布式核心)
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# 启用Redis去重(避免多实例抓取重复URL)
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# Redis连接配置(需部署Redis服务)
REDIS_URL = "redis://localhost:6379/0"

更多推荐