首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >从能跑到能扛:给数据管道装上三道防线

从能跑到能扛:给数据管道装上三道防线

原创
作者头像
dsy
发布2026-08-25 22:31:31
发布2026-08-25 22:31:31
520
举报

上篇我把管道跑通了,Gold 用户特征能算出来,当时还挺得意。

今天我故意让校验挂了一次,想看看系统怎么反应。结果吓一跳:程序报错退出了,但 data/gold/customer_features.parquet 那个旧文件还安安静静躺着。下游要是来取,会以为这是本次跑出来的结果。

那一刻我明白:能算出结果,不等于能放心用。从"能跑"到"能扛",中间差着几道我现在才补上的防线。

一、配置集中:把旋钮收拢到一个面板

最早的脚本里,路径是硬编码的。data/rawdata/bronze 散落在各个文件,改个目录得满世界找。

我把它们收进一个 settings.py,用 Pydantic 兜一层类型:

代码语言:python
复制
from pydantic import BaseModel
from pathlib import Path

class Settings(BaseModel):
    raw_dir: Path = Path("data/raw")
    bronze_dir: Path = Path("data/bronze")
    silver_dir: Path = Path("data/silver")
    gold_dir: Path = Path("data/gold")

类比:硬编码路径像把电器的旋钮装在机器背面,每次调都得拆机箱。现在所有旋钮收在一个面板上,改开发、测试、生产环境,动这一处就行。

下游直接引用,不再写死字符串:

代码语言:python
复制
GOLD_FILE = settings.gold_dir / "customer_features.parquet"

二、校验与 fail-fast:门口的质检员

新增 validate.py。今天最关键的认知是:校验不是维修工,是门卫。

它不负责把脏数据修干净,只负责"不合规就不让进下游"。比如:

代码语言:python
复制
if (df["order_count"] <= 0).any():
    raise ValueError("order_count 必须大于 0")

只要整列里有一个不合格,立刻抛异常,后面的导出不会执行。

这背后是一种叫 fail-fast(快速失败)的策略:发现第一个问题就拉闸,不带着问题继续跑。

代码语言:python
复制
customer_features = build_customer_features(event_features)
validate_customer_features(customer_features)   # 挂了就停在这
export_customer_features(customer_features, gold_file)

它的好处很实在:

  • 坏数据不会继续往下游传;
  • 错误位置好定位,因为程序死在第一个出问题的地方。

也有代价:它不会一次性把全部错误都告诉你。想要"收集所有问题再统一报错",得换个写法:

代码语言:python
复制
errors = []
if condition:
    errors.append("错误一")
if another_condition:
    errors.append("错误二")
if errors:
    raise ValueError(";".join(errors))

两样各有所用。数据管道里我更想要 fail-fast——宁可早死,别把半成品喂给下游。

三、日志与原子写入:留字条,别骗人

新增 logging_config.py,统一格式:

代码语言:python
复制
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s | %(levelname)s | %(name)s | %(message)s",
)

现在每次跑 Gold,都会留下这么几行:

代码语言:shell
复制
开始构建 Gold 用户特征
购买事件数量:4
事件特征数量:4
用户特征数量:3
Gold 文件已生成

类比:日志是给几个月后的维护者留的字条。那时候你已经忘了这段代码怎么写的,字条会告诉你任务何时开始、读了多少、卡在哪一步。

但今天还扒出一个更阴险的隐患。

校验失败时,旧 Gold 文件不会自动删除。于是出现一种诡异状态:本次任务实际失败了,旧文件却还在,下游误以为那是本次成果。这比直接报错更危险——它不报错,只是 silently 骗你。

生产环境的正确写法是原子替换:先写临时文件,写成功了再替换正式的,绝不拿半成品去覆盖正牌。

代码语言:python
复制
temp_path = output_path.with_suffix(".tmp.parquet")
features.write_parquet(temp_path)
temp_path.replace(output_path)

顺序很关键:临时文件写成功,才 replace。中间任何一步炸了,正式文件完好无损。

心得

今天折腾完,我给"能扛的管道"画了像。它至少有四样东西:配置、校验、日志、分层数据,加上明确的输入和输出。

但今天最值钱的一句话不是某个技术点,而是一个判断标准:

数据管道不仅要计算出结果,还要知道配置从哪来、数据合不合格、失败发生在哪里,以及失败后会不会留下误导性的旧结果。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 一、配置集中:把旋钮收拢到一个面板
  • 二、校验与 fail-fast:门口的质检员
  • 三、日志与原子写入:留字条,别骗人
  • 心得
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档