Ohhnews

分类导航

$ cd ..
DZone Java原文

dbt 携手 Apache Flink:数据工程师的统一工作流

#dbt#apache flink#数据工程#流处理#sql

[LOADING...]

管理 Snowflake、BigQuery 以及越来越多 Databricks 上的批处理 SQL 管道,以及 Apache Flink 上的流处理管道的数据工程师面临一个熟悉的问题:两套工具链、两套技能组合、两套 CI/CD 管道。dbt 正在向流处理领域扩展。本文将解释这在实践中意味着什么,为什么它对数据工程团队很重要,以及在 Confluent Cloud 上使用 Apache Flink 的具体实现是什么样子。

数据流处理遇见湖仓一体

数据湖曾承诺解决企业数据问题。现实却更加混乱。批处理管道产生过时信息,分析工作负载在业务事件发生数小时后才运行。等到查询运行时,行动窗口往往已经关闭。

湖仓一体模式改善了情况。Apache Iceberg 已成为主流的开放表格式,受到 Snowflake、Databricks、BigQuery 以及越来越多的查询引擎支持。团队可以直接在对象存储中的数据上运行 SQL 分析,而无需将其复制到专有数据仓库中。但仅靠湖仓一体并不能解决实时问题。数据仍以批处理方式到达,在源事件发生后数分钟或数小时。这一差距反映了更深层的架构分裂。使用 Apache Kafka 和 Flink 的数据流处理是运营层:它处理关键 SLA,支持事件驱动应用,并让业务系统实时运行。湖仓一体是分析层:它存储历史数据用于报表、机器学习以及近实时或批处理分析。这是两种不同的工作负载,在正常运行时间、数据丢失、延迟和吞吐量方面有不同要求。它们需要共存,而不迫使工程师构建和维护两条独立的管道。

Kafka、Flink 和 Iceberg 如何协同工作

这正是 Apache Kafka、Apache Flink 和 Apache Iceberg 的组合所解决的问题。Kafka 在源头捕获每个事件,并作为实时系统的运营骨干。Flink 处理并丰富运动中的数据,既支持即时运营决策,也为下游分析准备数据。Iceberg 将结果存储为受治理、可查询的表,供任何分析引擎使用,无论是 Snowflake、BigQuery 还是 Databricks。关于此架构的完整论述,包括模式演进、压缩和目录集成,请参见此处:数据流处理遇见湖仓一体:Apache Iceberg 实现统一的实时和批处理分析。问题不再是流处理和湖仓架构能否共存。它们已经共存。问题在于数据工程团队如何在不维护单独工具链的情况下跨两者工作。这就是 dbt 进入画面的地方。

什么是 dbt?

dbt(data build tool)是一个用于基于 SQL 的数据转换的开源框架。一个 dbt 模型是一个保存为文件的 SQL SELECT 语句。dbt 通过模型如何使用 ref() 相互引用来推断执行顺序。标准命令覆盖整个工程工作流:dbt run 对目标平台执行 SQL,dbt test 验证数据质量,dbt docs generate 生成可浏览的文档目录。dbt 成功的原因在于它为 SQL 工作带来的纪律。在 dbt 之前,转换逻辑分散在脚本和专有 ETL 工具中。dbt 用代码优先、版本控制的工作流取代了这一点,并内置了血缘、测试和文档。Snowflake 和 BigQuery 是当今大多数 dbt 采用所在。两者都是 SQL 原生,并针对 dbt 所围绕的 ELT 模式进行了优化。Redshift 是 AWS 环境中的强大第三平台。Databricks 最近在 dbt 采用方面增长,受无服务器 SQL 仓库投资推动,但其根源在 Spark 和 Python,使其成为 dbt 生态中较新的进入者。dbt Labs 在 2025 年初 ARR 超过 1 亿美元,拥有超过 5,000 名付费客户。目前约有 90,000 个 dbt 项目在生产中运行。Fivetran 和 dbt Labs 于 2025 年 10 月宣布的合并,创建了一家合并年收入近 6 亿美元的数据基础设施公司——这明确表明 dbt 已远远超越一个流行的开源工具,成为基础性企业数据基础设施。今天管理批处理和流处理的数据工程团队在两个分开的现实中运作。一边是 Snowflake 或 BigQuery:dbt 模型、版本控制的 SQL、自动化测试、生成的文档。另一边是 Apache Flink:Terraform 脚本、自定义部署代码或 Flink 控制台。技能和实践不能在两者之间转移。这种分离有真实成本。流处理管道更难测试、更难记录、更难交接。许多团队通过保持流处理逻辑最小化并将转换工作推向下游数据仓库来补偿,这重新引入了延迟并削弱了流处理的意义。愿景很直接:一个 SQL 工作流用于两者。在 Snowflake 或 BigQuery 上构建 dbt 模型的工程师应该能够将同样的方法应用于 Apache Flink 流处理管道,而无需切换工具或从头重建 CI/CD。两套工具链意味着两种测试策略、两种文档系统,以及两套需要招聘和保留的技能组合。治理执行在两个环境中变得不一致。SQL 是使这一现实可行的共同基础。Flink SQL 成熟且经过生产验证。Snowflake 和 BigQuery 是 SQL 原生的。Apache Iceberg 表可通过 SQL 在多个引擎中查询。dbt 用工程纪律封装 SQL。模型文件看起来相同。ref() 依赖解析以相同方式工作。测试和文档生成通过相同命令工作。组织不需要招聘单独的 Flink 基础设施专家。现有的数据工程团队可以拥有两边。Apache Iceberg 在存储层连接两个世界。一个 Flink 管道将结构化、受治理的事件写入组织自己 S3 桶中的 Iceberg 表。同一张表可立即被 Snowflake、BigQuery 或 Databricks 读取,无需任何额外 ETL 步骤。dbt 可以跨整个管道建模数据:在数据流经 Flink 时塑造它,并在数据落入数据仓库进行分析时再次转换它。这也是 Shift Left Architecture 2.0 的直接推动因素。Shift Left 方法将数据集成逻辑移近源头,在数据落入湖仓一体之前,在流处理层应用质量检查、丰富和治理。到目前为止,这需要大多数 dbt 原生团队不具备的流处理特定技能。dbt for Flink 大幅降低了这一门槛。完整的架构细节请参见此处:Shift Left 架构 2.0:面向实时数据产品的运营、分析和 AI 接口

具体示例:在 Confluent Cloud 上使用 Apache Flink 的 dbt

今天最具体的实现是 dbt-confluent 适配器,由 Confluent 与 confluent-sql Python 驱动一同发布。两者都是开源的,可在 PyPI 和 GitHub 上获得。数据工程师将流处理管道定义为 dbt 模型,并使用标准 dbt run 命令将其部署到 Flink 计算池。入门只需一步:

pip install dbt-confluent

支持三种物化方式:view 用于基于 Kafka 主题的虚拟 Flink SQL 视图,streaming_table 用于持续始终最新的结果集,streaming_source 用于将 Kafka 主题定义为 dbt 源。测试是确定性的,使用 Confluent Cloud 的快照查询能力返回有界的时间点结果,而不是在超时时静默通过。文档生成通过 INFORMATION_SCHEMA 集成工作,生成与 Snowflake 和 BigQuery 项目相同的可浏览目录。底层 confluent-sql 驱动符合 DB-API v2,意味着任何兼容工具都可以直接连接到 Confluent Cloud Flink:Airflow 和 Dagster 用于编排,Pandas 用于快照查询,Streamlit 用于实时仪表板,LangChain 用于 AI 代理工作流。对于已经在 dbt 中工作的数据工程师,这意味着围绕 Snowflake 或 BigQuery 构建的技能和实践直接转移到架构的流处理侧。

数据工程师用 dbt 拥有批处理和流处理

批处理和流处理工程之间的分离一直是组织上的而非技术上的。两个世界都使用 SQL。两者都需要测试、文档和可靠部署。工具只是从未弥合差距,因此组织配备并运营两个不同的工程学科。dbt 扩展到 Apache Flink 改变了这一等式。今天在 Snowflake 或 BigQuery 上运行 dbt 的数据工程师可以将相同的心智模型、命令和 CI/CD 管道应用于 Flink 流处理管道。无需 Flink 基础设施专业化。他们编写 SQL 模型、定义测试、生成文档并部署,正如他们为批处理所做的那样。含义很直接。对 dbt 技能和工具的投资现在进一步扩展到架构中。流处理可以由已经受信任处理批处理的同一数据工程团队逐步采用。一个团队、一个工具、一个治理标准,跨越运营和分析工作负载。用于 dbt 的 Flink 适配器与 Snowflake 或 BigQuery 上的 dbt 相比成熟度更早,团队应预期与不断发展的生态系统合作。但基础是坚实的,方向是明确的,核心架构组件已经在多个行业的大规模生产中运行。来自数据工程团队的需求真实且在增长。