进阶 dagster.io 2026-10-10 01:48:56 · 6 阅读

第23章 独自用 Dagster 在 9 天内对百万级 Snowflake 列进行分类

第23章:独自用 Dagster 在9天内完成百万级 Snowflake 列分类

我们正在把十年积累的遗留系统整合进 Snowflake,结果留下了超过一百万个需要分类打标签、用于数据治理的列。我搭建了完成这项工作的系统:一个分级分类引擎、一个人工修正流程,以及把它们串联起来的编排层。这是一篇工程实战记录——讲述让它变得有意思的那些挑战,以及 Dagster 如何充当编排层。一个人独立开发,9天上线。

我在 Group 1001 工作,这是一家科技驱动的综合性金融服务与保险集团。如果你看过 Indy 500® 赛车赛,就会在赛车上看到我们的品牌 Gainbridge℠。集团旗下还有 Delaware Life Insurance Company、Gainbridge Life Insurance Company、Clear Spring Life and Annuity Company、Clear Spring Property and Casualty Group、RVI Group 等公司,管理资产规模超过 800 亿美元。

我们正在进行一场大型迁移:把原本各自独立跑在遗留 SQL Server 上的多个业务域(销售与市场、精算、投资、运营、人力资源)整合进同一个 Snowflake 账户。每个域都带着自己的历史、schema、命名规范,以及几十年前对客户记录、精算表、HR 字段如何建模的种种决定。合并之后,就是横跨 62 个数据库、超过一百万列的规模,而且没有任何两个源系统对同一种敏感字段的命名是相同的。

作为一家旗下有受监管保险子公司的集团,分类不是可选项。每一列都需要一个标签(SPI 表示敏感个人信息、PII、或 NON_SENSITIVE),以及 Snowflake 中对应的治理标签,因为脱敏和访问策略都建立在这些标签之上。到了百万列的规模,唯一可行的方案就是全自动化、覆盖全账户,并能在新数据落地时持续跟进。

而我是独自一人、在兼顾其他工作的情况下构建它的。

技术挑战

在选工具之前,先看看系统真正要解决的问题。

大规模、低成本地准确分类。 列名承载了大部分信号,但不是全部,而那些模糊的案例恰恰是不能出错的。语言模型可以裁决它们,但对百万列逐一调用既慢又贵。引擎需要分级:简单的大多数交给廉价的确定性规则,昂贵的判断只留给模糊的少数。

多阶段流水线,但一个人要看得懂。 发现步骤喂给规则,规则扇出到两个独立的分类引擎,结果再合并回来。这个图由一个人维护,而这个人可能一离开就是好几周,所以它必须一眼可读,而不是每次都要靠回忆重建。

两个独立触发器,一个打标签动作。 标签的应用有两种方式:定时任务,以及人工修正引擎。这两条路径毫无共享(不同触发器、不同节奏、不同数据路径),却必须收敛到完全相同的打标签逻辑。一旦拆成两份代码,就会漂移,你就是在调试两个写标签的程序而不是一个。

轮询循环下的安全修正。 业务专家需要在不能接触 SQL 的情况下修正引擎。如果系统每分钟轮询一次待处理修正队列,而一次执行超过一分钟,轮询就会重叠。处理不当,你要么把同一个修正应用两次,要么悄悄丢掉别人正在等待的那一条。

可观测性,但不自建可观测性平台。 一个人干活,没有预算专门搭日志和仪表盘来判断某次执行有没有做对事。

解决方案

系统分两部分:分类引擎——主要是 Snowflake 层面的事,以及驱动它的编排层——这就是 Dagster 的用武之地。

分类引擎

引擎分三级。

Tier 1 是一组按优先级排序、对列名执行的正则规则,更具体的匹配优先于泛化匹配。每条规则把名称模式映射到一个标签,并且刻意防范一个天真实现容易踩的坑:某个列名看起来敏感,但该列的数据类型根本存不下这类值。数据类型过滤器解决了这个问题——一个名为 LAST_CONTACT_DT 的 TIMESTAMP_NTZ 列会被识别为元数据时间戳而非联系人字段,而同样的名字出现在文本列上才算真正命中。规则列表还包含 EXCLUDE 规则,用于那些看起来敏感实际不然的名称,直接短路跳过检查。EMAIL、PHONE、NAME 这类模式在这里就解决了。凭借优先级排序、类型过滤和排除规则,这一层以零成本、确定性地覆盖了大多数列。

少数规则模糊到仅凭名称匹配不可信。命中这些规则的列不会被打标签,而是被留下来交给后续层级做二次判断。

Tier 2 调用 Snowflake 原生的 SYSTEM$CLASSIFY 处理抽样数据库,并将其隐私类别映射到我们的标签。

Tier 3 只把模糊案例发给 Snowflake Cortex。Prompt 是刻意写成对抗式的,因为语言模型在这里的失败模式是过度分类:它看到 customer_email_verified 这样的名字就想标记,哪怕这列存的是布尔值而不是邮箱。Prompt 反向施压,指示模型拒绝任何只是描述或追踪敏感数据而非存储它的列:"REJECT:该列是标志位、日期、计数、ID、描述、审计字段或权限字段;它追踪敏感数据但并不存储它。"把任务框架设定为"拒绝"而非"检测",正是压低误报率的关键。由于规则层已经处理了 95% 以上的匹配,这一层规模小、成本低——这正是整个设计在百万列规模下依然可负担的原因。

这就解决了第一个挑战。下面的一切都是编排,也正是我用 Dagster 的地方。

用 Dagster 做编排

流水线形态。 Dagster 的 software-defined assets 让依赖图直接体现在函数签名里:

``` @dg.asset def rule_based_classification(discover_columns): ...

@dg.asset def snowflake_classification(rule_based_classification): ... # branch A

@dg.asset def cortex_classification(rule_based_classification): ... # branch B

@dg.asset def store_classification_results( snowflake_classification, cortex_classification # fan-in ): ... ```

两个 asset 把 rule_based_classification 声明为输入,于是形成扇出;一个 asset 同时引用两个分支,于是形成扇入并等待两者完成。Dagster 从签名中读出这些关系,并发执行各分支,并把每个 asset 的输出路由给它声明的消费者。不需要维护任何边,也没有胶水代码。下面是这条分类主干在 Dagster 中的呈现:

classification_job 的 asset 图:discover_columns 喂给 rule_based_classification,扇出到 cortex_classification 和 snowflake_classification,两者再喂给 store_classification_results,最后是 update_classification_run。

对可读性挑战来说,回报就在这里:依赖图本身就是文档。离开几周后回来,读一遍签名就知道数据流向,不用去追边找数据从哪来。(对比一下,Airflow 的 DAG 描述的是任务间的执行顺序;任务之间的数据传递是另一件事,得自己接线。)

两个触发器,一个打标签动作。 收敛问题是整个方案里最有意思的。两个生产者——定时的 classification_job 和人工驱动的 override_job——都要汇入同一个 tagging_job,且任一生产者都不需要知道消费者的存在。Dagster 的 run_status_sensor 做的正是这件事。下游订阅上游的成功事件;上游对此一无所知:

``` @dg.run_status_sensor( monitored_jobs=[classification_job], run_status=dg.DagsterRunStatus.SUCCESS, request_job=tagging_job, ) def tagging_sensor(context): return dg.RunRequest()

@dg.run_status_sensor( monitored_jobs=[override_job], run_status=dg.DagsterRunStatus.SUCCESS, request_job=tagging_job, ) def tagging_after_override_sensor(context): return dg.RunRequest() ```

tagging_job 有两个上游触发器,但不知道是哪一个触发的。生产者也不知道有任何东西在监视它们。一个写标签的程序,两条可达路径,组件之间零耦合。

值得看看没有这个原语要付出什么代价。如果 classification_job 直接调用打标签代码,分类就得知道打标签的存在,override job 也一样,等于把每个生产者硬编码到消费者上。用定时器在其他任务之后调度打标签则很脆弱,一旦某次执行拖长就断了。队列或事件总线——为一个扇入引入一整个子系统。Airflow 只能走到一半:它的 TriggerDagRunOperator 耦合方向是反的,因为上游 DAG 必须知道并触发下游;它的 Datasets 特性(下游可订阅数据更新)是最接近的类比,但响应的是数据变化而非运行结果。run_status_sensor 把触发完全放在消费者一侧,几行代码,专门响应成功事件。

Dagster 的 jobs 列表,展示 classification_job、health_check_job、override_job 和 tagging_job,各自带有最近的成功运行。

轮询循环下的安全修正。 业务专家通过一个小型 React UI 提交修正,它的唯一职责是把意图写进一张 COLUMN_OVERRIDES 表。他们从不直接碰 Snowflake 标签。一个 sensor 把这个队列转成运行,幂等性难题靠一个 run key 的选择解决:

``` @dg.sensor(job=override_job, minimum_interval_seconds=60) def override_sensor(context: dg.SensorEvaluationContext): rows = conn.cursor().execute(""" SELECT OVERRIDE_ID FROM COLUMN_OVERRIDES WHERE STATUS = 'pending' ORDER BY CREATED_AT ASC """).fetchall()

if not rows: yield dg.SkipReason("No pending overrides") return

# 排序后的待处理 ID 作为 run_key:同一批永远不会触发两次 run_key = "_".join(sorted(str(r[0]) for r in rows)) yield dg.RunRequest(run_key=run_key) ```

run_key 是排序后的待处理 ID 集合,Dagster 不会启动两个相同 key 的运行。如果 sensor 在一批还在处理中时又触发,它算出的是同一个 key,Dagster 会丢弃重复。如果有新修正到来,待处理集合变化,key 随之变化,就成了一次等待下一个 tick 的新运行。批处理、去重和重叠处理全都源于这一个决定。如果用 cron 循环手写,同样的保证意味着要建一张锁表、显式定义两个 tick 重叠时的行为、记录哪些修正已应用,还要一套覆盖所有情况的测试。在这里,这些是框架契约加一个 sorted()。

连接与可观测性。 每个 asset 都通过 Dagster resource 访问 Snowflake,而不是各自构建连接——环境感知的客户端只定义一次,参数注入,测试时可换成 fake。每个 asset 在每次运行时都会上报自己的元数据:

``` context.add_output_metadata({ "match_count": matched, "label_distribution": dg.MetadataValue.json(dist), }) ```

结果

现在这个系统按固定节奏运行,覆盖 62 个数据库、超过一百万列,并附带每日健康检查。每次运行处理全账户清单。业务专家的修正提交后约一分钟内就落地到 Snowflake,走的是与定时任务相同的打标签路径。每次运行都把计数作为元数据上报,持续审查只需一瞥而非一次排查,asset 图则成为整个系统如何拼合的活文档。

九天换来了什么

脚手架早已就位:EKS、Snowflake 账户和角色、基础规则库结构。这段时间里我构建的是分类引擎、打标签流水线、人工修正流程和 React UI。

时间的分配才是重点。分类引擎——发现步骤、规则分级、两个引擎分支——大约占了头两天,难点在规则本身,而不在围绕它们的编排。这一周剩下的时间大部分花在人工修正流程和 React UI 上,主要是应用开发和 SQL 工作。而本文讲的编排部分——sensor、run-status 链接、run_key 幂等性——只占很小的时间比例,因为每块就是几十行代码。

这就是一个人能在九天内上线的原因。Dagster 没有让我写 UI 更快、规则写得更好,它把编排、耦合、幂等性和可观测性这些本会围绕真正问题堆砌起来的工作压缩掉了,让时间花在问题上,而不是管道上。

本文观点仅代表个人,不一定代表雇主立场。本内容仅供信息参考,不构成且不应被解读为个性化投资、法律或税务建议。

©2026 Group 1001 IP Properties, LLC | Group 1001. 保留所有权利。有反馈或问题?在 Slack 或 Github 上发起讨论。有意与我们合作?查看我们的在招职位。想看更多类似内容?在 LinkedIn 上关注我们。

评论 (0)