当“也想在我们的仪表板上查看竞争对手数据”传到数据团队时
内部数据源是稳定的。Schema 由我们控制,变更会通过发布说明提前通知,即使失败也有重处理流程。然而,当管理层和业务团队提出“也想在同一个仪表板上查看竞争对手价格、评论和市场数据”时,就会出现一个问题:必须将控制权之外的数据源接入管道。
外部网页的 Schema 并不归我们所有。目标网站会在没有预告的情况下更改页面结构,拦截策略也会随时加强。昨天采集为空,今天仪表板就会出现缺口;通过 Excel 收到的数据也无法接入 DAG。本文讨论的是如何将外部网页采集整合到内部 DW·数据湖管道中——即其整合方式。前提只有一个:采集本身委托出去,数据团队专注于加载·建模·使用。因为要像对待“管理良好的源系统”一样处理外部网页,该数据源就必须以数据团队的语言建立契约。
1. 外部网页数据与内部数据源有何不同
在接入之前,必须先了解差异,才能进行整合设计。内部表的列变更处于我们的发布流程之中,但外部网页可由源网站自行更改。昨天还存在的字段今天可能消失,同样的“价格”可能混杂为 12,900韩元和12900。更新也不是事务,而是按固定周期运行的定期批处理,因此需要让 DAG 调度和仪表板新鲜度 SLA 与该周期相匹配。
最危险的区别在于失败模式。内部任务失败时会抛出异常并停止,但外部爬虫最糟糕的情况是已损坏却假装成功并持续积累空值。当网站结构变化导致选择器失效时,任务会正常结束,但记录可能为 0 条或充满缺失值。仅靠查看数量的验证无法捕获这种问题,而过去的网页页面无法再次采集,因此该缺口将永久存在。因此,外部数据整合并非“仅加载”就结束,而是加载后的验证也必须作为一套流程。
2. 按交付形式划分的整合模式——没有唯一正确答案,应匹配规范
将外部采集接入管道的方式取决于交付形式。大致分为三类,与其说哪一种正确,不如说应选择符合数据团队管道规范的方式。
| 模式 | DAG 触发方式 | 适用场景 | 协商要点 |
|---|---|---|---|
| API 轮询 | 调度+传感器 | 由数据团队控制拉取时机 | 增量游标·分页·认证 |
| DB·S3 直接加载 | 到达传感器 | 大容量、避免轮询负载 | 完成标记·表/前缀规范 |
| 文件接收 | 文件传感器 | 低频、Excel→管道迁移 | 文件名·路径·工作表 Schema |
API 轮询型通过 API 获取采集数据,由 DAG 仅查询最后一个水位线之后的增量并加载至落地区。关键在于数据团队持有增量游标,从而避免重复接收已获取的区间。DB·S3 直接加载型由采集方直接加载到 DW·数据湖,DAG 只检测到达状态并执行后续处理。在这种模式下,需要提前约定完成标记(manifest·_SUCCESS)规则,以免消费只到达一半的数据。文件接收型是指将符合规范的文件投放至标准位置后进行解析和加载,适用于将以 Excel 为中心的运营方式迁移到管道的过渡阶段。
三种模式的前提都是由我们(采集方)与数据团队先对齐接口规范。应根据数据团队的管道协商 API 响应格式·DB 表/S3 前缀·文件规则,并从该规范反向决定采用何种交付形式。只有当交付形式不仅限于 Excel,还支持 API·直接加载·S3 投放时,这种协商才能成立。
3. 落地区——即使外部数据源变更,下游也不会中断的缓冲装置
如果将外部数据直接放入建模层,数据源变更当天,下游(dbt·仪表板)就会整体中断。落地区就是缓冲装置,在此处理三件事。
第一,保留 raw 原始数据。将收到的数据原样保留一份(API 响应原文、原始文件),不进行修改。这样,当解析逻辑出错或日后需要其他字段时,可以无需重新采集,直接从原始数据重新解析。过去的网页页面无法恢复,因此保留原始数据实际上就是安全网。
第二,Schema 版本管理。我们与数据团队协作的标准如下:Schema 变更会提前通知(设置 deprecation 期限),并在该期间同时提供旧/新版本,让下游能够按自己的节奏迁移。例如,持续输出 price_daily.v1的同时,一并提供新增currency列的v2;下游 dbt 准备好后迁移引用,待 deprecation 结束后再下线v1。这样,“数据源昨天改版,今天仪表板就宕了”会变为“并行运行至下个季度,并在此期间完成迁移”。由于变更责任和缓冲期都写入契约,数据团队就不会在毫无预警时被叫去救火。
第三,将 backfill 作为标准流程处理。当采集缺失数日,或需要为历史区间补充新字段时,我们指定期间·条件重新采集并重新加载至落地区,DAG 则只对该区间重新建模。前提是通过分区替换·upsert 键进行幂等设计,确保重复写入同一区间时不会产生重复数据。由于重新采集的窗口在我们这一侧,数据团队只需指定“需要重跑哪个期间”。
4. 加载后验证——在 DAG 中捕获静默缺失
第 1 节所述的“静默缺失”只能通过加载后验证来捕获。在publish之前设置 DQ 门禁,未通过时不晋升到下游。损坏的数据不会一路流入仪表板,而是停留在落地区并触发告警。门禁通常检查四项:记录数是否骤降至预期下限以下(静默缺失)、必填字段缺失率是否低于阈值(部分采集)、是否存在解析异常标记(结构改版)、采集时间是否处于 freshness SLA 范围内(新鲜度)。
为了让数据团队无需从零开始编写 DQ 规则,我们会随数据一并交付用于验证的元数据:采集时间(freshness)、采集条件(哪些搜索词·类别·范围)、源 URL(lineage·审计)、记录数(相较预期是否骤降)、缺失/异常标记等。监控不只设置在爬虫一端,而是同时部署在数据源(采集)和管道(加载)两端。一旦检测到采集异常·延迟,我们会通过数据团队频道(Slack·Webhook)通知,并同步共享预计恢复时间,确保能够回答“为什么为空、何时会补齐”。
5. 数据目录·血缘——回答“这个数字从哪里来的”
一旦外部数据进入仪表板,终有一天会在审计或管理会议上被问到:“这个竞争对手价格的来源是什么?”要回答这个问题,必须在数据目录和血缘中登记来源·周期元数据。为便于直接用于登记,我们会一并提供源域名/URL、采集周期、采集条件(目标范围)、当前 Schema 版本,以及该数据源属于委托运营数据源的说明。这些内容与验证元数据部分重叠,但增加了数据目录视角的字段。
在注册数据集时填充这些元数据后,就可以在血缘图中从仪表板指标反向追溯至源 URL。只需点击几次就能回答“这个数字从哪里来的”——这才是将外部数据以与内部数据源同等地位纳入治理的状态。与数据团队对齐 Schema·backfill·元数据的原因也汇集于此。只有以可正式注册到数据目录的形式交付外部数据源,它才能成为“管理良好的源系统”。
最后是新鲜度。外部采集是按固定周期运行的定期批处理,因此仪表板 SLA 也应根据该批处理周期设定。如果批处理每天一次,仪表板也应设计为“按天更新的指标”,文案最好明确写明批量采集时间,例如“截至○时”。将新鲜度沟通为“基于哪个批处理时间点”,而不是“有多新”,对数据团队和业务团队都更安全。
这种整合为数据团队留下什么
将外部网页采集整合进管道的核心,是将不可控的数据源置于可控的接口之后。通过 Schema 版本契约封装 Schema、通过 DQ 门禁封装失败、通过数据目录元数据封装来源、通过 backfill 标准流程封装重处理,外部数据也可以像内部数据源一样被管理。
最后的决定在于:谁来承担该接口背后的持续运营——应对拦截、跟踪网站改版、采集监控、修复。如果由数据团队承担,负责 DAG 运营的人力将被外部采集维护消耗殆尽。我们将开发·维护·拦截应对·采集监控纳入月费套餐,以委托方式承接,并向数据团队交付具备 Schema 契约·验证元数据·数据目录元数据的数据源。我们通过采集→清洗→AI 分析,交付的不是 raw 数据,而是可加载的结构化数据,因此也能减少整理落地区的负担。
结论只有一个。将外部网页采集·维护委托出去,使其像管理良好的源系统一样被对待;数据团队则专注于加载·建模·使用。实际上,我们正在以定期管道的形式受托运营某监管机构的大规模在线监测,大规模·多部门需求也可通过这种方式集中至一个窗口。
推荐阅读
- 从爬取数据到成为决策依据
- 通过 Excel 接收的数据 vs 在仪表板查看的数据——交付形式改变的内容
- 各部门分别进行爬取的公司——终结重复投入的全公司数据采集治理
立即开始
请告诉我们您当前管道采用何种规范——是 Airflow 编排、S3 数据湖,还是 DW 直接加载——我们将从接口开始,共同设计如何让外部网页采集匹配该规范。包括 Schema 契约·验证元数据·数据目录元数据,全部将以数据团队的语言进行对齐。




