第26章 Dagster 社区精选:构建垂直数据栈
Dagster 社区里诞生了不少最有意思的项目。这篇文章就来展示一些社区成员的创意作品。我们最喜欢的 Dagster 使用案例,往往是我们从未预料到的。
社区成员正在用 Dagster 探索公开数据集、监控基础设施、自动化科研工作流、构建内部工具,以及尝试全新类型的数据应用。有些项目技术含量很高,有些则非常小众,但都体现了社区用户的创造力。
本文精选了几个引人注目的社区项目,聊聊它们背后的故事、创作者为什么选择 Dagster,以及构建过程中的体验。
如果你也在用 Dagster 做有趣的东西,欢迎告诉我们。
Christian Casazza Linkedin | Github
介绍一下你自己和你在做的事情
我是在工作中自学成才的数据工程师,因为我逐渐意识到数据工程是放大数据价值的利器。在一个系统中打好数据工程的基础,能让下游的一切都变得更容易。我最初和开发者合作,是为电商运营分析搭建 pipeline。我很享受构建 pipeline 逻辑的过程,之后就不断向开发者请教各种原理。从 GPT4 开始,LLM 达到了“足够好用”的水平,我便开始借助开源数据原语自学 Agentic Engineering。
如今,我在构建数据价值创造系统。我的主要方式是边做边学,尽可能多地搭建端到端的 pipeline。纽约市和纽约州拥有世界上最好的政府开放数据之一,因此它成了我的“练手场”,我围绕预算、交通、房地产、立法等多个领域构建 pipeline。现代开放数据原语加上编码 agent,让我能在数据增强和分析两个方向上深入挖掘,而不需要前期投入大量资金。除了纽约相关的核心 pipeline,我还基于 NFLverse 这类优质数据源和 Polymarket 这样的预测市场,做了一些体育领域的数据产品。
你是怎么发现 Dagster 的?
最初是在社交媒体上看到 Dagster,发现数据领域那些“信号很强”的人都喜欢它,所以等我的 pipeline 复杂度到了合适的程度,就尝试了一下。让我印象深刻的是,Dagster 把数据工程中很多通用的“后勤工作”(编排、调度、asset 检查、factory 创建等)都处理好了,同时仍然保持了原汁原味的 Python 风格。我喜欢它能脱离基础设施在本地快速开发,这样启动一个新 pipeline 的门槛很低,开发迭代也很快。
你用 Dagster 构建了什么项目?
我正在用开源数据原语构建一个垂直数据栈(vertical data stack)。你可以在 QueryStation.app 上浏览我公开的精选数据集。Dagster 让我能够构建详细的采集、清洗、连接和分析 pipeline。
我的数据 pipeline 主要围绕 Arrow、Parquet、DuckDB、Polars 和 DuckLake 构建。所有原始数据首先在 Parquet 层标准化,大部分数据整理工作由 Polars 的 LazyFrame 引擎完成。DuckDB 是我的引擎,负责把清洗后的中间 Parquet 文件写入 DuckLake,并运行下游的 SQL 分析 pipeline。DuckLake 是我的开放表格式提供者,让表的持续更新和数据应用层消费表都简单得多。Apache Arrow 保证各层之间可以互通,比如把 dataframe 传给 DuckDB,或者 TypeScript 数据应用层查询 DuckLake 表并接收 Arrow IPC。Arrow 让基于 SQL 的 API 构建变得容易。想要查询我已精选好的数据集(比如 NYC 311 或各种预算数据集)的用户,直接用 SQL 查询 QueryStation 就能拿回 Arrow 数据,读入 dataframe 后即可供 agent 分析。
Dagster 是整个数据工程系统的基石。它像一座桥梁,把一堆基于原语的脚本变成一个互联、连贯的系统。它通过自动处理 DAG 管理、调度、asset 检查、元数据存储等,把数据 pipeline 的生命周期标准化。Dagster 的设计理念就是把信任内嵌到数据 pipeline 中,让使用者真正愿意使用产出。这个仓库的设计刻意围绕 Dagster 形成固定风格,这样 agent 就能专注于特定场景的数据工程,而不必操心底层数据软件的实现。一个核心思路是降低构建和扩展新数据源的边际成本。我一直在开发 factory 风格的模式:给 AI agent 一个 URL 或 API 文档,agent 分析后就能基于模板生成正式的 Dagster asset,抓取文档站点并生成结构化的 markdown 文件。“Dagster expert” skill 对我把开发锚定在核心 Python 最佳实践上也很有帮助。
你解决过的最难或最有意思的问题是什么?
最有意思的挑战是让 factory 函数上线,并以可扩展的方式设计 asset。我不想为每个 asset 单独设计,而是希望有共享模式,让 asset 能继承检查、元数据、文档等通用行为。我朝着 asset registry 的概念努力:一种自动注册 asset、生成文档,并让 agent 在工作中沉淀知识的机制。
def create_socrata_pipeline( name: str, socrata_config: SocrataIngestConfig, schema: SchemaContract, *, # Organization (NEW) domain: str | None = None, geographic_scope: str | None = None, group: str | None = None, # Partitioning partitions_def: PartitionsDefinition | None = None, clean_partitions_def: PartitionsDefinition | None = None, partition_mapping: PartitionMapping | None = None, # Logic post_transform_fn: Callable[[pl.LazyFrame], pl.LazyFrame] | None = None, enrichments: StandardEnrichments | None = None, # ... more parameters ) -> PipelineResult: """Create a 2-stage Socrata pipeline: Landing (CSV) → Clean (Parquet).""" # Normalize schema: accept 2-tuple or 3-tuple, extract contract schema_contract, is_3tuple = normalize_schema(schema) # ... compose landing asset, clean asset, checks return PipelineResult( landing=landing_asset, clean=clean_asset, checks=checks, )
最难的部分是前期搭建。有些工作在 dg CLI 出现之前就完成了,所以很多结构只能手动组装。Dagster 里同一个问题往往有多种解法,找到最合理的思维模型需要反复试错。
我还在摸索在 Dagster 中处理开放表维护任务的最佳方式,比如 DuckLake 表所需的那类任务。一个思路是做一个维护 asset,把 checkpoint 或表维护行为打包进图里,同时支持在 Python 层面设置 DuckLake 配置。不过 DuckLake 还比较新,有些最佳实践只能靠时间来验证。
为什么 Dagster 适合这个项目?
Dagster 非常契合这个项目,因为其中包含大量领域上毫不相关的数据 pipeline,但它们需要在 asset 代码、检查、测试、元数据、调度等方面保持一致的模式。AssetFactory 模式让为数据工程中纯“后勤”部分设计稳定、可复用的代码 API 变得容易。
使用 Dagster asset factory 还能更容易地固化最佳实践,防止 agent 代码在 pipeline 之间产生漂移。举个例子,Polars 在做数据整理时,用 LazyFrame API 配合 scan 和 sink_parquet 比传统的 Eager API(.collect、read、write_parquet)性能好得多。我发现 agent 一旦有自由发挥的空间,就倾向于在 pipeline 里偷偷塞入手动 .collect 调用,造成过早物化、把惰性 pipeline 变成即时执行。Factory 函数提供了强大的护栏,自动保证转换全程惰性执行,agent 就可以专注于真正想做的数据整理逻辑,而不是实现细节。
Dagster 还有一个很棒的特性:它天然适合构建对 agent 友好、带有上下文领域知识系统的数据 pipeline。Dagster 把所有元数据(runs、asset 信息等)都存在 SQL 表里,可以用 Postgres,也能通过 API 查询。用户还可以创建自定义元数据并附加到 asset 上,这些元数据会显示在 Dagster UI 中。基于这些现成的原语,我能给 asset 附加各种类型的领域知识——数据来源、分析师发现的事实、数据质量或小毛病之类的注意事项——agent 在工作中发现的关键信息都可以用 Dagster 编码进 pipeline。把信息编码进 pipeline,能大大方便上下文在跨 agent 会话、跨模型提供商之间传递,也更容易在前人工作之上继续构建更深入、分层级的分析。
对刚上手 Dagster 的人有什么建议?
Dagster 最强的特质之一是高度可组合,但这恰恰也可能是新手的坑。在 Dagster 里,如何组织代码、逻辑放在哪里,往往有多种“正确”的实现方式,这容易导致重构瘫痪——总想找到那个用上“所有”特性的“最佳”方案。
最好的办法是先做出一条完整的端到端 pipeline,最好选一个和主工作代码独立的项目。去搭、去改、去弄坏,弄清楚各个部分如何协同。然后再逐步叠加复杂度,之后再加入调度、asset 检查、自定义元数据等高级特性。不断观察改动后什么会坏、什么会变,直到形成你满意的形态。我也建议定期让 agent 用 Dagster 团队提供的 Dagster-Expert 和 Dignified Pipeline skill 检查你的仓库,它们能帮你的代码锚定在最佳实践上,避免 agent 漂移。
不过,一旦跨过最初的学习曲线,就别怕在 Dagster 上玩得更大。它的可组合性让你能像一个代码原生的工作流构建器那样搭建非常强大的工作流。关键是先从简单的开始,做出能跑的东西,然后逐步增加复杂度,看看会出什么问题。
我个人推荐用我的仓库作为 Dagster 入门方式。这套固定风格的设计已经内置了很多核心基础设施和设计模式。任何使用 agent 的人都可以基于这个仓库构建自己的端到端数据 pipeline,把现有代码作为示例基础——我相信它能覆盖大多数人想构建的各类数据 pipeline。
有反馈或问题?欢迎到 Slack 或 Github 发起讨论。想和我们一起工作?查看我们的在招职位。想看更多类似内容?关注我们的 LinkedIn。