做爬虫的人,只要采集量过了千万条这个坎儿,基本都会遇到一个尴尬:数据越存越多,但越来越不敢用。原始 JSON 堆了一堆,字段说改就改,重复采集的数据散落各处,领导问一个数据要查半天,临时写脚本捞数据的时候永远在猜“这个目录里到底是哪天的数据”。我自己的经验是,爬虫本身不是难点,难点在于爬下来的数据怎么组织,才能让你在三个月后还能用。今天这篇就聊聊我落地的一套方案,核心思路是构建一个基于 source / date / schema_version 三层分区的轻量数据湖,配合一套元数据表做治理,让爬虫数据从“能存”变成“好用”。这套方案不需要上 Hadoop 那套重资产,也不会让代码复杂度失控,适合爬虫数据量在 GB 到 TB 级、有一定多源采集任务、又不想被存储结构拖垮的团队或个人。
整个方案的落点就一句话:用目录结构做索引,用元数据表做治理,让爬虫数据在采集端就天然有序,查询端做分区裁剪,治理端做到可追溯。
1. 为什么爬虫数据需要数据湖,而不是继续堆 JSON 文件
1.1 原始文件堆叠的三个致命问题
很多爬虫项目跑着跑着就变成了一堆 JSON 文件的坟场。按日期建文件夹、按站点建文件夹、甚至直接按采集批次命名的情况我都见过。表面上看数据都在,实际上已经接近不可用了。
第一个问题是重复采集无法去重。爬虫失败重试、定时增量、补数据,同一批内容被拉了三四遍,每遍都生成一份新文件,谁也不知道哪份是最新的,哪份是不完整的。问就是“后面处理的时候会去重”,但真正处理的时候你会发现,去重逻辑比爬虫本身还难写,尤其当数据没有天然主键的时候。
第二个问题是 schema 变化失控。目标网站改版、接口加字段、字段类型变化,这些都是爬虫的日常。如果你只是把响应体原封不动存下来,那还好;但凡你做了字段抽取,老的抽取逻辑产出的字段和新逻辑对不上,下游想用数据就得写一堆 if else 来兼容各种历史版本。
第三个问题是查询效率极低。没有索引意味着每次用数据都要全量扫描。数据量小还能忍,到了几十 GB 甚至上 TB,光扫描一遍就要等半天,更别提清洗加工了。
1.2 数据湖不是大厂专属,轻量方案也能做
一说数据湖,很多人第一反应就是 Iceberg、Hudi、Delta Lake 这一套,觉得要上 Spark、要搭集群、要考虑存储计算分离,对一个小爬虫项目来说简直杀鸡用牛刀。这个认知其实是个误区。
数据湖的本质是"存储 + 元数据"的分层架构:存储层放原始文件,元数据层记录这些文件的位置、结构、产生时间、数据量等信息。Iceberg 们做的事情本质上也是维护这么一套元数据,只是它们把元数据管理的内核做成了表格式规范,自动帮你在写入时维护事务、快照和时间旅行。
对于爬虫场景来说,完全可以用更轻的方案达到九成效果。核心就是自己控制文件目录的分层规则,再用元数据表记录分区信息。这样既不需要 Spark 集群,也不需要理解快照隔离、MVCC 这些东西,却能拿到数据湖最核心的价值——把无序的文件组织成可跟踪、可裁剪、可演进的数据集。
提示:理解这一点很重要。你引入数据湖不是跟风,而是为了解决多源数据的组织问题和 schema 演进的混乱问题。如果你的数据量只有一百万条,接口也不经常变,那直接存 CSV 也够用,不需要折腾这套。
1.3 这套方案的架构思考
我这套数据湖方案在存储层上用的是普通的文件系统或对象存储(本地磁盘、NFS、MinIO、S3 都行),管理层是一张或多张元数据表(存在 MySQL 或 SQLite 里),查询层则通过 Python 脚本或 Trino/DuckDB 做分区裁剪和数据分析。
架构上非常朴素,没有引入任何重量级组件。文件解析、字段校验、分区路径拼装、元数据登记,这些全部在 Python 代码里完成,写起来也不复杂。
但正是因为这层“朴素”,它的适应面反而极广。小到一台 Linux 服务器跑采集,大到多台机器分布式采集写入 MinIO,方案都不需要变。唯一要变的是元数据表的存储位置,从 SQLite 换成 MySQL 或者 PostgreSQL 而已。
2. source/date/schema_version 三层索引的设计逻辑
2.1 为什么偏偏选这三个字段做索引
分区字段的选择是整个方案的灵魂。不是随便拍脑袋选三个字段,而是针对爬虫数据的特点,逐一分析出来的。
source 是业务维度的第一隔离。爬虫通常同时采集多个站点或接口,每个源的数据格式、采集频率、字段语义都不同。如果所有数据混在一起,后面做分析时必须强依赖“这个字段到底是哪个源产出的”,但数据本身又没有天然标记,就会出现严重的口径冲突。把 source 作为第一层分区,可以让不同源的数据从物理上隔离,互不干扰。
date 是时间维度的锚点。爬虫数据天然带时间属性,你需要回答“某月某日的数据是否完整”“某天某个指标为什么变高”。如果数据没有按时间做分区,这类问题每次都要全量算一遍。按日分区的另一个好处是,采集任务天然按日/按小时运行,写入路径可以直接和调度时间对齐,增量处理逻辑会非常自然。
schema_version 是结构演进的版本控制。这个字段是我踩过坑之后才加上的。早期做过一个爬虫,目标网站改版了好几次,我每次改完字段就直接覆盖写,结果三个月后发现自己根本说不清某些字段是什么时候加的,哪些旧数据没有这个字段。后来我规定,凡是字段有增删改,schema_version 必须递增,数据写入到对应版本的目录下,这样下游消费时既能按版本过滤,又能追溯某个字段从哪个版本开始存在。
2.2 三层索引的层级关系和查询模式
这三个字段之间存在明确的依赖关系,目录层级的设计也是按照这个依赖来的:
{base_path}/{source}/{date}/{schema_version}/为什么要这样排而不是把 date 放最前或者把 schema_version 插中间?核心原因是分区裁剪的效率。
查询的时候,最常见的过滤条件组合是“给我某个源某段时间的数据”,其次是“某个版本之后的数据”。把 source 放最前面,可以保证在扫描文件时把不相关的源直接跳过;date 放第二层,可以把时间范围压缩到几天乃至一天;schema_version 放最后一层,用于精确定位结构匹配的数据集。
反过来,如果你把 date 放在最前面,那跨源统计某天数据的时候会方便一点,但同时,任何查询都必须先限定时间范围,这对做跨时间窗口的分析就不太友好了。爬虫场景里,source 的隔离优先级高于 date 的聚合优先级,所以 source 必须在最外层。
注意:Partition 字段的选择没有绝对的对错,只有适不适合你的查询特征。你这个项目的核心查询是"按源管理数据",就以 source 为根;如果核心查询是"按天看全量趋势",就考虑往 date 放在第一层。别照搬,要思考。
2.3 分区的粒度:为什么按天而不是按小时
三层索引中,date 这一层我选择按天分区,而不是按小时。虽然有些爬虫是小时级甚至分钟级调度,但我在实际使用中发现,小时分区会产生大量碎片文件(每小时的目录里可能只有几 MB 甚至几百 KB),夸大了元数据管理的规模,也让后续数据处理任务的数量暴增。
按天分区的另一个好处是“天然的对账单位”。我每天睡觉前检查一下当天分区下有没有 _SUCCESS 标记文件(后面会讲写入流程),就能确认当天数据是否完整。如果按小时分区,就得检查 24 个目录,对账成本直线上涨。
当然,如果你的数据量已经大到单日分区都有几十 GB 甚至更大,那就需要考虑按小时分区了。判断标准很简单:单日数据量超过 5 GB,或者单日文件数超过 500 个,就可以考虑下钻到小时级。
3. 分层目录的工程落地:从采集到入湖的完整流转
3.1 目录结构定义和统一命名规范
先把最终的目录结构摆出来,后面所有代码和流程都会围绕这个结构展开:
/data/lake/ └── wechat/ # source: 微信公众平台 ├── dt=2024-05-20/ │ ├── schema_version=1/ │ │ ├── part-00001.parquet │ │ ├── part-00002.parquet │ │ └── _SUCCESS # 写入完成标记 │ └── schema_version=2/ │ ├── part-00003.parquet │ └── _SUCCESS └── dt=2024-05-21/ ├── schema_version=2/ │ ├── part-00004.parquet │ └── _SUCCESS └── schema_version=3/ ├── part-00005.parquet └── _SUCCESS源目录我用的是业务代号而不是中文名或带版本号的站点名,比如wechat而不是wechat_2023_old。原因很简单,代码里写路径时越短越省事,而且业务代号只在元数据表里做一次映射,源改名了只改元数据,不用改存储路径。
日期分区我保留了dt=前缀。这个前缀是 Hive 风格的分区目录命名法,好处是查询引擎(Trino、Spark、DuckDB)可以直接识别为分区字段,不需要额外配置。如果你后面的查询直接用 Python 扫目录,前缀也没坏处,用glob匹配dt=*/schema_version=*还能少写一层路径拆解逻辑。
文件我用 Parquet 格式存。为什么不用 JSON?JSON 的好处是保留原始字段、不做类型推断、方便人类阅读,但坏处是查询性能差、文件体积大、没有内嵌 schema。Parquet 在这三个维度上全面优于 JSON,而且 Python 生态里pandas、pyarrow、duckdb都是原生支持,写起来并不比 JSON 复杂。
提示:如果是无法转成结构化表格的数据(比如嵌套极深的用户评论树,或者包含大量自由文本的响应体),Parquet 就不好使了,这种情况建议按原始 JSON 存放,但同样要套用三层分区结构,元数据表里记录 val 文件的格式标记即可。
3.2 Python 写入客户端的核心实现
爬虫采集到数据之后,经过清洗和字段规范化,由写入客户端完成"路径拼装 → 写入 Parquet → 登记元数据"三步。核心代码如下:
import datetime as dt import hashlib import json import os from pathlib import Path import pandas as pd from pyarrow import parquet as pq class LakeWriter: def __init__(self, base_path: str, source: str, schema_version: int): self.base_path = Path(base_path) self.source = source self.schema_version = schema_version def write_partition(self, df: pd.DataFrame, biz_date: dt.date) -> Path: # 1. 拼装分区路径: base/source/dt=xxx/schema_version=n/ partition_dir = ( self.base_path / self.source / f"dt={biz_date.isoformat()}" / f"schema_version={self.schema_version}" ) partition_dir.mkdir(parents=True, exist_ok=True) # 2. 生成文件名,包含批次号与数据指纹,保证幂等 batch_id = dt.datetime.now().strftime("%Y%m%d%H%M%S") content_md5 = hashlib.md5( pd.util.hash_pandas_object(df).values.tobytes() ).hexdigest()[:8] file_path = partition_dir / f"part-{batch_id}-{content_md5}.parquet" # 3. 写 Parquet,开启压缩 pq.write_table( pa.Table.from_pandas(df, preserve_index=False), file_path, compression="snappy", ) # 4. 返回文件路径,后续用于元数据登记 return file_path这段代码看着简单,但有几个细节值得展开讲。
幂等设计。文件名里带了content_md5,如果同一批次的数据被重放写入,文件内容一致、文件名一致,自然会被覆盖而不是产生重复文件。配合_SUCCESS标记文件的原子创建,可以实现“重复采集但不重复污染”。
类型规范化。写入前,所有字段需要做一次统一的类型规范化,比如时间字段统一转 ISO 格式、金额字段统一转 Decimal、空值统一为 None 而不是 NaN。这一步在爬虫端做掉,比在查询端再做要好得多,因为清洗逻辑归一到了源头。
3.3 _SUCCESS 标记文件与对账机制
目录下的_SUCCESS文件是整个写入流程的“提交信号”。原理借鉴了 Hive 和 Spark 的约定——只有_SUCCESS文件存在,才表示这个分区下的数据是完整、可用的。
实现上,我通常在写完所有数据文件之后,最后一步创建一个名为_SUCCESS的空文件。因为空文件创建在大多数文件系统上是原子操作,不会出现“文件写了一半”的情况。
对账任务每天跑一遍,重点检查两类问题:
- 缺失分区检查:今天应该有数据的源,
dt=今天的目录下_SUCCESS是否存在。如果不存在,说明采集任务失败了或者是数据的写入逻辑挂掉了。 - 残留分区检查:元数据表里登记过
is_active=1的分区,对应的目录是否都还在。如果目录被误删,元数据标记要同步更新。
这套对账机制帮我避免了多次“数据少了一半但没人发现”的尴尬局面。爬虫挂了不可怕,可怕的是挂了三天没人知道,后续每一天的增量都是建立在脏数据之上的。
3.4 增量与全量采集的策略统一
爬虫采集一般分两种模式:增量采集和全量重拉。在分层索引体系下,两种模式可以统一处理:
- 全量重拉:直接写入当天分区,schema_version 用当前结构版本。历史分区不动,保证时间线不可变。
- 增量采集:只抓取新增或变更的条目,写入当天分区,与已有文件合并为同一分区的多个文件。
实际上我在设计时遵循了一个更重要的原则:写入只有增量,没有更新。任何一次采集任务,永远只向“当前时间对应的分区”写入新文件,不修改历史分区中的数据文件。比如需要修复历史数据的采集错误,不是去改写旧文件,而是把修正后的数据写入一个新的分区(通常加一个修正标记或使用更高版本的 schema_version),在元数据表里通过映射关系指向正确的数据版本。这种方式在数据湖理念里叫作“不可变数据,可变元数据”,代价是存储会多一点冗余,好处是每次写入都像发布一个不可变快照,事务处理和审计追溯都变得极其简单。
4. 元数据表的实现与版本演化管理
4.1 为什么要元数据表,文件系统本身不够吗
有人会问:目录结构就是索引,直接在文件系统上 glob 扫描不就行了,为什么还要维护一张元数据表?我自己最开始也觉得多此一举,直到踩了几个坑:
- 文件扫描慢。目录数量到了几十万个(多源、多年、多版本),递归遍历一次要几分钟。
- 信息不够。目录名只告诉你“有什么”,不告诉你“有多少条、多大、校验值是多少、有没有被处理过”。
- 无法做状态追踪。某个分区是否已经完成 ETL、是否已经推送到下游,文件系统上无从体现。
元数据表的本质是“从文件系统的树形结构里提取出结构化信息,让你可以用 SQL 来回答问题”。比如我想知道“微信源的数据里,每个版本的 schema 各覆盖了哪些日期范围”,这个查询在文件系统上几乎无法高效实现,在元数据表里一条GROUP BY就出来了。
4.2 核心元数据表设计
以下是我实现的核心表结构,也是这套系统最关键的设计文档:
CREATE TABLE dataset_partition_info ( id BIGINT AUTO_INCREMENT PRIMARY KEY, source VARCHAR(64) NOT NULL, -- 数据来源,例如 wechat / douyin dt DATE NOT NULL, -- 数据日期 schema_version INT NOT NULL, -- 结构版本号 base_path VARCHAR(512) NOT NULL, -- 相对根路径的目录路径 file_count INT NOT NULL DEFAULT 0, -- 文件数 row_count BIGINT NOT NULL DEFAULT 0, -- 行数 byte_size BIGINT NOT NULL DEFAULT 0, -- 总字节数 file_format VARCHAR(16) NOT NULL DEFAULT 'parquet', created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, is_active BOOLEAN NOT NULL DEFAULT TRUE, -- 逻辑删除标记 checksum VARCHAR(64) DEFAULT NULL, -- 分区内容指纹,用于校验完整性 UNIQUE KEY uk_source_dt_version (source, dt, schema_version) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;每一条记录都代表一个独一无二的“分区”,也就是source/dt/schema_version三个字段的唯一组合。
字段的选择讲究如下:
- base_path不是主键,因为同一个目录可能因为修正数据被重新登记(逻辑删除+新记录),但物理路径不变。加上唯一索引
(source, dt, schema_version)保证同一分区不会重复登记。 - row_count 和 byte_size用来判断数据量突变。某天行数突然从十万掉到一千,大概率是采集逻辑漏了数据,光看目录结构发现不了这种问题。
- file_count是大小文件问题的观测指标。如果单分区文件数持续增加,就需要做文件合并操作。
- checksum登记的是分区内容的总和校验,用于跨环境复制数据时验证完整性。
4.3 schema_version 的生成与登记流程
schema_version 的生成比较讲究,不能随手写,必须有一个明确的规则。我用的是“变更才递增”的原则:只有当采集字段列表、字段类型、嵌套结构或枚举取值发生变化时,schema_version 才加一。
实现上,我在代码里维护了一个SCHEMA_REGISTRY字典:
SCHEMA_REGISTRY = { "wechat": { 1: {"title": "string", "author": "string", "publish_time": "timestamp", "content": "string"}, 2: {"title": "string", "author": "string", "publish_time": "timestamp", "content": "string", "read_count": "int"}, 3: {"title": "string", "author": "string", "publish_time": "timestamp", "content": "string", "read_count": "int", "like_count": "int"}, } }每次写数据前,拉取最新注册的 version 对应的 schema,对采集到的 DataFrame 做字段校验。如果字段和注册的 schema 对不上,直接报错而不是自作主张写入,这样倒逼开发者在改采集字段时同步更新SCHEMA_REGISTRY。
4.4 三种常见的 schema 演化场景及处理策略
梳理我见过的 schema 演化场景,核心其实是三种:
添加字段。比如公众号文章新增了“评论数”。处理策略是递增一个版本,新版本数据使用新字段结构写入新分区;历史分区不动,历史数据没有的字段在查询时做一个COALESCE(field, NULL)处理即可。不建议为了统一把历史数据全量刷一遍,既耗时又没必要。
删除字段。比如接口下线了一个字段。策略同样是递增版本,新版本不再采集该字段。这里要特别说明,历史数据的这个字段并不会消失,只是查询时需要注意版本差异——schema_version=1有某个字段,>=2的版本没有,查询 SQL 需要兼容。我在元数据表里额外记录了每个版本的字段列表 JSON,就是为了让这类兼容查询可以自动生成。
类型变更。这个最麻烦。比如publish_time从字符串变更为时间戳。处理策略有三种:一是解析时统一转成老类型(牺牲精度);二是写入新版本字段,物理类型不同但业务含义相同;三是分成两个字段,老字段保留原值、新字段存新解析结果。我自己的默认选择是第三种,因为查询时最直观,也不会因为自动转换而丢失原生值的信息。
5. 查询优化与生命周期管理:索引如何反哺数据治理
5.1 分区裁剪的原理和实际效果
三层索引体系的最大受益者其实是查询侧。当你用 Trino、DuckDB 或 Spark 查询数据时,只需要在 SQL 里指定source、dt、schema_version三个过滤条件,查询引擎就能通过元数据直接跳过大部分数据文件,只读取匹配的分区内容。这个机制叫分区裁剪(Partition Pruning)。
实际项目中,我的一个典型查询长这样:
SELECT dt, COUNT(*) AS article_count, AVG(read_count) AS avg_read FROM lake.wechat_articles WHERE source = 'wechat' AND dt BETWEEN DATE '2024-05-01' AND DATE '2024-05-31' AND schema_version >= 2 GROUP BY dt ORDER BY dt在没有索引时,这个查询需要扫描该源全网文件;加了分区裁剪后,它只读取五月份且 version>=2 的目录下的 Parquet 文件。实测数据量 500 GB 规模的查询,可以做到从十几分钟降到几十秒。
5.2 元数据表反哺治理的三个实际场景
除了加速查询,元数据表还能用来做更细粒度的数据治理。这里说三个我用得最多的场景。
数据完整性对账。我之前提到每天跑对账任务,其实核心就是一条 SQL:
SELECT source, COUNT(DISTINCT dt) AS missing_days FROM dataset_partition_info WHERE dt BETWEEN DATE_SUB(CURDATE(), INTERVAL 7 DAY) AND CURDATE() AND is_active = TRUE GROUP BY source HAVING COUNT(DISTINCT dt) < 7这条语句能找出最近 7 天中数据不完整的源。审核团队和数据需求方都很吃这一套,因为数据质量变得可量化了。
冷热数据识别与归档。元数据表里的updated_at和row_count字段能帮我判断哪些分区很久没被访问了。结合文件系统的访问时间(atime),我可以识别出“近三个月没被查询的源”和“某天之后再也没人看的 schema 版本”,把这些冷数据迁移到低频存储或压缩归档,节省不少存储成本。
数据血缘追踪。我在元数据表里增加了一个parent_partition_id字段,表示当前分区是由哪个上游分区加工而来。爬虫采集的原始数据落在raw层,经过清洗后的数据落在clean层,加工后的特征数据落在feature层,每层之间都通过这个字段建立血缘关系。排查数据问题时,顺着血缘链就能找到是哪个环节出了问题。
5.3 与查询引擎的对接
Python 生态下最省事的查询引擎我推荐 DuckDB。它可以直接查询 Parquet 文件,不需要像 Trino 那样起一个服务,嵌入式使用的方式让它尤其适合单机/小集群场景。下面是一个典型的查询片段:
import duckdb conn = duckdb.connect() conn.execute(f""" SELECT source, schema_version, COUNT(*) AS cnt FROM read_parquet('{base_path}/wechat/dt=2024-05-*/schema_version=*/*.parquet') GROUP BY source, schema_version """).fetchdf()如果数据量已经到 TB 级或者需要多引擎共享,推荐直接用 Trino,它原生支持 Hive 风格的分区目录,把/data/lake挂载成外表后,所有分区字段自动识别,无需额外声明。
5.4 生命周期管理:数据过了多久该清理
数据湖最容易被忽视的问题就是数据膨胀。我给每个源都定义了不同的保留周期,这个规则写在元数据配置表里:
| 源 | 保留周期 | 清理策略 |
|---|---|---|
| 日志类 | 90 天 | 直接删除过期分区,更新元数据 is_active |
| 业务快照类 | 1 年 | 过期后压缩为汇总表,再删原始 Parquet |
| 全量资料类 | 永久 | 只保留最新 schema_version,历史版本按季度归档 |
清理任务的核心思路是:先更新元数据表,再删文件。千万不要反着来,否则文件删了但元数据还在,下游会去读一个不存在的路径,白白浪费排查时间。
6. 数据湖治理的六个常见坑与解决路径
6.1 小文件泛滥:存量分析尽力,增量合并必做
Parquet 文件如果几十 MB 一个,查询性能尚可。但爬虫常常是持续写入,几分钟一个小文件,一个月下来产生了上万个文件,查询引擎光打开文件就要花费大量时间。
根治思路其实就一句话:存量尽力而为,增量合并必做。增量合并的关键是写路径上的 buffer 逻辑,不建议直接在上面的write_partition方法里做合并,而是单独跑合并任务,把当日的多个小文件读进来重新写入一个大文件,然后更新元数据表的file_count和byte_size。合并之后,旧文件做逻辑删除(用is_active标记),不要立刻物理删除,留足观察期以防合并逻辑有问题。
6.2 时区错位:分区日期必须统一基准时区
这个坑踩得最深。早期采集任务部署在多台机器上,有的机器用的 UTC,有的用 CST,结果同一个 "2024-05-20" 分区的数据在不同机器上实际对应的时间差了一个时区,日期对不上,对账天天报错。
统一方案并不复杂:在代码里凡是格式化dt字段的地方,都强制使用一个统一的时区常量。我用的pytz.timezone('Asia/Shanghai'),业务侧需求都是北京时间,就以它为基准。任务部署的机器必须明确设置系统时区,并在采集日志里打印当前时间带时区信息,排查问题的时候一眼就能看出来是否错位。
6.3 schema_version 管理失控
理论上 schema_version 只有变更时递增,但实践中多人协作很容易出现冲突:A 在wechat源加了read_count字段,B 同时加了like_count字段,两个人都用的 version 2,最后写入的分区数据结构和字段全部乱掉。
我的解决方法是把SCHEMA_REGISTRY集中管理,每个人改 schema 前必须先在版本目录里提交变更请求,评审通过后统一发新版本号。虽然增加了流程开销,但换来的是 schema 的血缘可追溯。个人项目或者小团队如果不想引入太重流程,至少要做到在代码仓库里注明这个版本的字段列表是哪个 commit 更新上去的,确保变更可回溯。
6.4 目录重命名导致指针失效
文件系统里可以通过 MVCC 或类似机制实现“目录硬链接”,只要有一个原子操作把 commit 指针从一个目录切换到另一个目录,就能保证目录切换对读取方是瞬时的。但如果跨文件系统或者对象存储不支持硬链接,就只能在元数据表里做逻辑切换:新写入的数据放到全新的分区目录,确认完成后把元数据表里的is_active标记切换指向。这个过程在元数据表里是通过两次UPDATE完成的(先置新分区为is_active=1,再置旧分区为is_active=0),不是严格意义上的原子切换,所以我在代码里加了一个状态字段switching,处于切换过程中的分区会被打上这个标记,生产环境里禁止这个分区的写入操作,避免不一致。
6.5 元数据与文件系统不一致:恢复策略与机制
做这套系统的人都会经历一个噩梦:元数据表有记录,但文件被误删了;或者文件还在,但元数据表的is_active被错误置为 0。这类不一致问题一旦出现,整个数据体系的信任度就会崩塌。
我的经验是定期做一次“文件系统扫描”与元数据表比对,把两边不一致的地方汇总出来:
for source in sources: for partition_dir in Path(base_path, source).glob("dt=*/schema_version=*/"): dt_str = partition_dir.parent.name.replace("dt=", "") version = partition_dir.name.replace("schema_version=", "") ... # 比对: 目录是否在元数据表中 / 元数据记录是否对应真实目录解决不一致的策略是“以文件系统为准重建元数据”。元数据表本质上只是缓存和索引,文件系统中的目录结构才是真相的唯一来源,这个原则不能颠倒。一旦发现不一致,先做全量盘点,再重建对应记录。
6.6 补数据机制:旧分区不能动,新分区来承接
补历史数据这个场景很多人处理得特别粗暴:直接把修正文件写到原来的日期目录下,覆盖老文件。这个过程一旦失败,轻则丢数据,重则把原来可用的数据都搞坏。
我建议的补数策略是:修正数据永远写入一个新目录,比如dt=2024-05-20/schema_version=3/(时间还是原来的业务日期,但 schema_version 升一档),然后通过元数据表将“业务日期 2024-05-20 的共识版本”重定向到版本 3。这样物理上原本的 version=2 目录显然还在,万一版本 3 有问题还可以切回版本 2,无损回滚。
如果修正涉及的内容不是结构变化,而只是数据本身错了,可以考虑用dt=2024-05-20-fix1/这类目录来命名修正版本。但注意,用额外标记目录的话一定要在元数据表里说明这个分区的业务日期是哪个,不然查询引擎无法自动识别它是哪一天的数据。
回到开头那个问题,“爬虫数据越存越不敢用”,其实本质是缺少一套让数据“自然有序”的机制。我落地这套方案后最大的体感是:每天对数据对账从人工变自动,上游字段变更对下游的影响从全量排查变成精准定位,新增数据源从改代码改成加配置。这套数据湖架构虽然不是最重的,但它给爬虫数据提供了一个足够结实的骨架,让后续的 ETL、分析、甚至机器学习特征工程,都站在了有序数据的肩膀上。对我来说,做数据采集的最高境界,不是“把数据拿回来”,而是“让数据在拿回来的那一刻就准备好了被使用”。才算真正把爬虫做成了数据工程。