爬虫数据入仓实战:将Scrapy数据实时写入ClickHouse
在数据驱动的业务场景中,爬虫是获取外部数据的核心手段,而数据入仓则是实现数据价值的关键环节。Scrapy 作为 Python 生态中成熟的爬虫框架,能高效抓取各类网页数据;ClickHouse 作为列式存储数据库,凭借超高的查询性能和写入吞吐量,成为实时数据分析场景的优选。本文将从实战角度出发,详细讲解如何搭建 “Scrapy 爬虫 + ClickHouse 数据仓” 的实时数据链路,解决数据抓取后即时入仓的核心问题。
一、技术选型背景:为什么选择 Scrapy+ClickHouse?
在开始实战前,我们需要先明确技术选型的合理性 —— 不同工具的组合需匹配业务对 “抓取效率” 和 “数据处理效率” 的双重需求。
1. Scrapy 的核心优势
- 高效抓取能力:内置异步下载器,支持多线程并发抓取,可通过调整
CONCURRENT_REQUESTS等参数优化吞吐量,单爬虫实例日均可抓取百万级数据。 - 灵活的中间件机制:提供下载中间件(处理请求头、代理)、爬虫中间件(过滤重复请求),可快速集成反反爬策略(如 UA 池、IP 池)。
- 结构化数据处理:通过
Item和Item Pipeline规范数据格式,支持自动去重、数据清洗,避免脏数据流入下游。
2. ClickHouse 的核心优势
- 超高写入吞吐量:采用列式存储 + 批量写入机制,单节点写入速度可达每秒数十万行,完全适配爬虫的高并发数据输出场景。
- 实时查询性能:针对聚合查询(如
COUNT、SUM、GROUP 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_id用String(避免 ID 为字符串格式时丢失信息),price用Float32(平衡精度与存储),crawl_time用DateTime(支持时间范围查询)。 - 引擎选择:
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"更多推荐
所有评论(0)