Ep.07 数据归宿:MongoDB 高效存储与 ETL 数据清洗流水线
这是《Python爬虫进阶》专栏的第七篇。
摘要:
分布式集群带来了海量数据,但脏乱差的数据毫无价值。本篇将探讨爬虫数据的终极归宿——为何放弃 MySQL 转投 MongoDB 的怀抱。我们将深入解析 Schema-free 架构在应对频繁变动的网页结构时的绝对优势。实战环节,不仅手把手带你编写 Scrapy 的多级 Item Pipeline,实现数据清洗、去重与格式化(ETL);还会揭秘工业级高并发写入技巧,利用
pymongo的bulk_write(批量写入) 和 Upsert (更新插入),榨干数据库性能,为你的单兵自动化采集系统打下坚实基石。
Ep.07 数据归宿:MongoDB 高效存储与 ETL 数据清洗流水线
在上一篇中,我们构建了基于 Scrapy-Redis 的分布式爬虫集群。成百上千个 Worker 像不知疲倦的工蜂,疯狂地把目标网站(如 Libvio 的视频元数据)搬运回来。
但随之而来的是幸福的烦恼:数据洪流涌入,存哪里?怎么存?
如果你还在用 txt、csv,或者每次抓取前都要痛苦地去 MySQL 里 CREATE TABLE 然后对齐字段,那么你需要升级你的武器库了。今天,我们来聊聊爬虫工程师的最爱——MongoDB,以及如何构建一条稳健的 ETL 数据清洗流水线。对于想要以极低运维成本跑通全链路的独立开发者来说,这是必修课。
一、 为什么爬虫的首选是 MongoDB?
很多后端开发出身的朋友,第一反应是把数据塞进 MySQL 或 PostgreSQL。但在爬虫领域,关系型数据库往往会成为效率的绊脚石。
1. 痛点:网页结构是动态的
假设你今天抓取了视频的 标题、导演、评分,并在 MySQL 里建了这三个字段。
明天,网站改版了,突然增加了一个 演员列表 和 播放量 字段。
用 MySQL,你就得去 ALTER TABLE,修改爬虫代码,处理由于字段缺失导致的 SQL 报错。这就陷入了无穷无尽的运维噩梦。
2. 破局:Schema-free (无模式) 架构
MongoDB 是一种 NoSQL 文档型数据库,数据以 BSON(类似 JSON)的格式存储。
它不需要预先定义表结构!
- 灵活性极高: 爬到什么字段,就存什么字段。今天少个字段,明天多个嵌套列表(比如包含了多条评论的字典),MongoDB 都能照单全收。
- 数据结构同构: Python 爬虫解析出的
Item(字典),可以直接无缝无损地插入 MongoDB,中间不需要任何 ORM 转换。
二、 拒绝脏数据:构建 ETL 流水线
数据存进库之前,必须经过清洗。
ETL 代表 Extract(提取)、Transform(转换)、Load(加载)。
在 Scrapy 中,这个过程完美映射到了 Spider 和 Item Pipeline 的生命周期中。
- Extract (提取):
Spider中用 XPath/CSS 选择器提取原始字符串。 - Transform (转换): 在
Pipeline中剥离 HTML 标签、转换时间格式、计算清洗。 - Load (加载): 在另一个
Pipeline中执行高并发数据库写入。
实战:多级 Pipeline 协同
好的工程实践是职责分离。我们不在一个 Pipeline 里把清洗和存储全做了,而是分成两个串联的卡点。
首先,在 settings.py 里配置优先级(数字越小越先执行):
ITEM_PIPELINES = {
'myproject.pipelines.CleanDataPipeline': 300, # 先清洗
'myproject.pipelines.MongoStoragePipeline': 400, # 后存储
}
阶段 1:Transform - 数据清洗管道
在这个阶段,我们将杂乱的原始字符串变成标准数据。
# pipelines.py
from itemadapter import ItemAdapter
from datetime import datetime
import re
class CleanDataPipeline:
def process_item(self, item, spider):
adapter = ItemAdapter(item)
# 1. 清理标题中的多余空白符和换行
if adapter.get('title'):
adapter['title'] = adapter['title'].strip().replace('\n', '')
# 2. 转换评分:从字符串 "8.5分" 提取浮点数 8.5
if adapter.get('rating'):
match = re.search(r'(\d+\.\d+)', adapter['rating'])
adapter['rating'] = float(match.group(1)) if match else 0.0
# 3. 添加爬取时间戳 (非常重要,用于后期数据追踪)
adapter['crawled_at'] = datetime.utcnow()
return item # 将清洗后的数据传递给下一个 Pipeline
三、 工业级加载:MongoDB 的高并发与 Upsert
经过清洗的数据终于来到了最后一关:入库(Load)。
新手写 MongoDB 存储,通常是一条条 insert_one()。但在分布式集群下,每秒成百上千的数据量会让数据库的网络 I/O 成为瓶颈。
此外,爬虫会重复抓取相同的页面,如何保证数据不重复?
核心技巧:Bulk Write (批量写入) + Upsert (更新插入)
- Upsert (Update or Insert): 如果数据库里已有这条视频(通过唯一 ID 判断),就更新它(比如更新播放量);如果没有,就插入新数据。这就从根本上解决了数据重复问题。
- Bulk Write: 把几十上百个 Upsert 操作打包成一个请求发给数据库,极大降低网络延迟。
阶段 2:Load - 工业级 MongoDB 管道实战
import pymongo
from pymongo import UpdateOne
class MongoStoragePipeline:
def __init__(self, mongo_uri, mongo_db):
self.mongo_uri = mongo_uri
self.mongo_db = mongo_db
self.batch_size = 100 # 累积 100 条再统一写入
self.items_buffer = []
@classmethod
def from_crawler(cls, crawler):
return cls(
mongo_uri=crawler.settings.get('MONGO_URI'),
mongo_db=crawler.settings.get('MONGO_DATABASE', 'crawler_db')
)
def open_spider(self, spider):
self.client = pymongo.MongoClient(self.mongo_uri)
self.db = self.client[self.mongo_db]
self.collection = self.db['videos']
# 建立唯一索引,防止脏数据突破防线
self.collection.create_index("video_id", unique=True)
def process_item(self, item, spider):
# 将 item 加入缓冲池
self.items_buffer.append(dict(item))
# 当缓冲池满了,执行批量写入
if len(self.items_buffer) >= self.batch_size:
self.flush_buffer()
return item
def flush_buffer(self):
if not self.items_buffer:
return
requests = []
for item in self.items_buffer:
# 构造 UpdateOne 操作:条件匹配 video_id,操作为 $set 全量更新,开启 upsert
req = UpdateOne(
filter={"video_id": item.get("video_id")},
update={"$set": item},
upsert=True
)
requests.append(req)
# 执行批量操作,ordered=False 可以让单个错误不影响全局写入
try:
self.collection.bulk_write(requests, ordered=False)
except Exception as e:
# 在生产环境中,这里应该将写入失败的数据记录到死信队列或日志中
print(f"批量写入异常: {e}")
# 清空缓冲池
self.items_buffer.append.clear()
def close_spider(self, spider):
# 爬虫关闭时,把缓冲池里剩余的尾部数据写入
self.flush_buffer()
self.client.close()
四、 总结
建立一条稳定、高效的数据流水线,是区分“业余玩家”和“工程高手”的重要标志。
通过 MongoDB 的 Schema-free 特性,我们解放了繁琐的表结构维护工作;通过 Scrapy Pipeline 的多级拆分,我们实现了逻辑清晰的 ETL 过程;最后,通过 Bulk Write 与 Upsert 技术,我们让存储层能够轻松抗住分布式集群的高并发冲击。
到这一步,从前端逆向、代理调度到后端存储的整套“分布式爬虫引擎”已经搭建完毕了。
下期预告:
一切都在完美运行,但如果你去睡觉了,目标网站突然改变了加密逻辑导致爬虫挂掉,或者 Redis 队列卡死了怎么办?一个成熟的系统离不开监控。
下一篇,我们将探讨如何利用 Docker 容器化部署,并搭建一套属于你的 数据可视化监控看板。
更多推荐



所有评论(0)