这是《Python爬虫进阶》专栏的第七篇。

摘要:

分布式集群带来了海量数据,但脏乱差的数据毫无价值。本篇将探讨爬虫数据的终极归宿——为何放弃 MySQL 转投 MongoDB 的怀抱。我们将深入解析 Schema-free 架构在应对频繁变动的网页结构时的绝对优势。实战环节,不仅手把手带你编写 Scrapy 的多级 Item Pipeline,实现数据清洗、去重与格式化(ETL);还会揭秘工业级高并发写入技巧,利用 pymongobulk_write (批量写入)Upsert (更新插入),榨干数据库性能,为你的单兵自动化采集系统打下坚实基石。


Ep.07 数据归宿:MongoDB 高效存储与 ETL 数据清洗流水线

在上一篇中,我们构建了基于 Scrapy-Redis 的分布式爬虫集群。成百上千个 Worker 像不知疲倦的工蜂,疯狂地把目标网站(如 Libvio 的视频元数据)搬运回来。

但随之而来的是幸福的烦恼:数据洪流涌入,存哪里?怎么存?

如果你还在用 txtcsv,或者每次抓取前都要痛苦地去 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 中,这个过程完美映射到了 SpiderItem 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 容器化部署,并搭建一套属于你的 数据可视化监控看板

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐