ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

从零掌握dbt-core:用代码思维构建可测试、可维护的数据转换管道

从零掌握dbt-core:用代码思维构建可测试、可维护的数据转换管道 1. 项目概述为什么你需要关注 dbt-core如果你正在和数据仓库、数据分析或者数据工程打交道听到“dbt”这个词的频率应该越来越高。它不是什么新的数据库技术而是一个彻底改变了数据分析师和工程师工作方式的工具。简单来说dbt 让你能用写代码主要是 SQL的方式去定义、测试和文档化你的数据转换逻辑并且这个“代码”是可以被版本控制、测试和协作的。听起来是不是有点像给 SQL 加上了软件工程的最佳实践没错这就是它的核心价值。我最初接触 dbt 是因为受够了传统 ETL 脚本的混乱。一堆散落在各处的 SQL 文件依赖关系全靠文件名里的日期来猜业务逻辑变更后根本不敢动历史代码生怕“牵一发而动全身”。dbt 的出现像是一套严谨的乐高说明书让搭建数据模型这件事变得可预测、可维护。而dbt-core是这一切的基础它是开源、可本地运行的核心引擎。网上很多教程一上来就讲云平台集成反而让初学者忽略了最本质的东西——如何用代码思维组织你的 SQL。这篇内容我们就扎扎实实地从dbt-core开始手把手带你构建第一个可运行、可测试的数据转换项目让你真正理解 dbt 的工作流和哲学而不是仅仅点几个按钮。2. 环境准备与项目初始化2.1 理解 dbt-core 的运作环境在动手安装之前我们先搞清楚 dbt-core 是什么以及它不是什么。dbt-core 本身不是一个数据库也不是一个调度器。它是一个命令行工具运行在你的本地机器或服务器上。它的核心工作是读取你编写的 SQL 和 YAML 配置文件连接到你的数据仓库如 Snowflake, BigQuery, Redshift, Postgres 等并在数据仓库中执行这些 SQL从而完成数据转换。它把数据仓库变成了一个“转换引擎”。因此要运行 dbt-core你需要准备两样东西一个可用的数据仓库连接这是 dbt 工作的“靶场”。对于入门学习我强烈推荐使用PostgreSQL。它免费、轻量且 dbt 对其支持非常好。你可以在本地安装 PostgreSQL或者使用云服务商提供的免费实例。Python 环境dbt-core 是一个 Python 包通过 pip 安装。这意味着你需要一个 Python 环境建议 Python 3.7 及以上。注意很多新手会困惑 dbt 和 Airflow 的关系。简单区分dbt 负责“转换”T in ELT即定义数据如何从原始状态变成业务可用的模型Airflow 负责“编排”和“调度”即决定何时、以何种顺序去运行 dbt 任务以及其他任务如数据提取。你可以用 Airflow 来调度dbt run命令。但入门阶段我们暂时不需要 Airflow先用 dbt 命令行手动执行理解其核心概念。2.2 一步步安装与验证假设我们已经有了一个本地 PostgreSQL 数据库里面有一些原始数据。现在我们来安装 dbt-core 并创建项目。首先为 dbt 创建一个独立的 Python 虚拟环境是个好习惯可以避免包依赖冲突。# 创建并激活虚拟环境以 macOS/Linux 为例 python -m venv dbt-env source dbt-env/bin/activate # 安装 dbt-core 以及适配你数据仓库的插件 # 这里我们安装适配 PostgreSQL 的插件 pip install dbt-core dbt-postgres安装完成后验证一下dbt --version你应该能看到 dbt-core 及其插件的版本信息。接下来初始化你的第一个 dbt 项目。dbt 会通过交互式提问来帮你创建项目骨架。dbt init它会问你几个问题项目名称例如my_first_dbt_project。这会成为你项目文件夹的名字。数据库连接选择postgres。主机、端口、用户、密码、数据库名、模式schema根据你的 PostgreSQL 配置填写。其中“模式”可以理解为数据集dbt 默认会在你指定的模式下创建模型。完成初始化后你会得到一个结构清晰的项目目录my_first_dbt_project/ ├── dbt_project.yml # 项目核心配置文件 ├── models/ # 所有模型SQL文件的存放目录 │ ├── staging/ # 推荐存放对原始数据做初步清洗的模型 │ ├── marts/ # 推荐存放面向业务的核心数据模型 │ └── example/ # dbt 自带的示例模型 ├── seeds/ # 存放静态数据文件CSV可加载到数据仓库 ├── snapshots/ # 用于实现缓慢变化维SCD逻辑 ├── tests/ # 自定义测试文件 ├── macros/ # 可复用的 Jinja 宏 ├── analyses/ # 临时或探索性分析 SQL └── target/ # 编译和运行的输出目录自动生成这个结构是 dbt 社区的约定俗成遵循它能让你的项目更容易被他人理解。dbt_project.yml是这个项目的“总控中心”我们马上会深入查看。3. 核心概念与项目结构深度解析3.1 解剖 dbt_project.yml项目的控制中心让我们打开自动生成的dbt_project.yml文件这是理解 dbt 项目配置的起点。# dbt_project.yml name: my_first_dbt_project version: 1.0.0 config-version: 2 profile: my_first_dbt_project # 对应 ~/.dbt/profiles.yml 中的配置名 model-paths: [models] # 模型文件所在路径 seed-paths: [seeds] # 种子数据路径 test-paths: [tests] # 测试文件路径 macro-paths: [macros] # 宏文件路径 snapshot-paths: [snapshots] # 快照路径 models: my_first_dbt_project: # 项目名下的模型配置 # 在此处配置的选项会应用于本项目所有模型 materialized: view # 默认物化方式为视图关键配置解读name和version: 项目标识在复杂项目中用于依赖管理。profile: 这是最关键的链接之一。它指向~/.dbt/profiles.yml文件中的一个配置块那里存储了数据库连接的用户名、密码等敏感信息。项目代码本身不包含密码实现了代码与凭证的分离。materialized: 物化策略。这是 dbt 的核心概念之一决定了你的 SQL 模型在数据仓库中将以何种物理形式存在。常见的有view视图每次查询时动态计算。构建快节省存储但查询性能可能较慢。table表执行dbt run时创建为物理表。查询快但构建慢占用存储。incremental增量表只处理新增数据极大提升大表构建效率。ephemeral临时模型不作为对象存在于数据库仅被其他模型引用时内联展开。在入门阶段我们可以先全部使用view因为它能快速迭代。随着模型成熟和性能要求提高再逐步将核心模型改为table或incremental。3.2 模型Models你的核心 SQL 代码模型是 dbt 项目的基石本质上就是一个.sql文件。但和普通 SQL 文件不同dbt 模型使用Jinja 模板语言对 SQL 进行了增强使其具备可编程性。一个最简单的模型文件models/staging/stg_customers.sql可能长这样{{ config( materializedview ) }} select id as customer_id, first_name, last_name, email, created_at, updated_at from {{ source(raw, customers) }} -- 引用源数据 where email is not null -- 简单的数据清洗这里发生了什么{{ config(...) }}: 这是 Jinja 语法用于设置此模型的特定配置这里覆盖了项目默认的view显式声明为视图。你可以在模型级配置物化方式、是否启用/禁用、设置唯一键等。{{ source(raw, customers) }}: 这是source函数。它引用的是在models/sources.yml中定义的数据源。这样做的好处是如果底层表名或模式发生变化你只需要更新 YAML 文件而无需修改所有 SQL 模型实现了声明式的数据血缘管理。3.3 源Sources与引用Ref建立数据血缘定义源Sources在models/sources.yml中我们定义原始数据表version: 2 sources: - name: raw database: my_database schema: raw_schema tables: - name: customers - name: orders description: 原始订单表包含所有订单记录 # 可以添加描述这个文件告诉 dbt“在my_database.raw_schema模式下存在名为customers和orders的表它们是我的数据源头。”模型间引用Ref在模型 SQL 中永远不要使用硬编码的表名去引用另一个 dbt 模型。取而代之的是使用{{ ref() }}函数。 例如在models/marts/dim_customers.sql中{{ config( materializedtable ) }} select customer_id, first_name, last_name, email from {{ ref(stg_customers) }} -- 引用上游模型 stg_customers{{ ref(stg_customers) }}会被 dbt 在运行时解析为stg_customers模型在数据库中的实际名称通常是schema.table的形式。这样做有巨大优势自动依赖管理dbt 能自动分析出dim_customers依赖于stg_customers并以此决定运行顺序。环境隔离在开发、测试、生产环境中你可以通过配置让ref指向不同模式下的表而代码无需改动。避免循环依赖dbt 会检查并阻止A引用B同时B又引用A的情况。3.4 测试Tests为数据质量上保险dbt 内置了一套简单而强大的数据测试框架。测试主要分两类通用测试在 YAML 文件中声明包括unique: 字段值是否唯一。not_null: 字段值是否非空。accepted_values: 字段值是否在指定列表中。relationships外键本表字段的值是否存在于另一表的指定字段中。自定义测试本质上是一个返回行数的 SQL 查询。如果查询返回 0 行则测试通过返回任何行则测试失败返回的行就是有问题的数据。如何为模型添加测试在模型所在的目录如models/staging/下创建一个.yml文件例如schema.ymlversion: 2 models: - name: stg_customers description: 清洗后的客户信息表 columns: - name: customer_id description: 客户唯一标识 tests: - unique - not_null - name: email tests: - not_null - name: first_name description: 客户名 sources: - name: raw tables: - name: customers columns: - name: id tests: - unique - not_null运行测试命令# 运行所有测试 dbt test # 运行特定模型的测试 dbt test --models stg_customers # 运行特定源的测试 dbt test --source raw当测试失败时dbt 会清晰告诉你哪个模型、哪个字段、违反了哪种测试并输出导致失败的具体数据行在target/目录下生成run_results.json和compiled/目录中可查看失败查询。这是数据管道可靠性的基石。4. 完整工作流实战构建一个迷你分析模型现在我们把所有概念串联起来完成一个从原始数据到分析模型的完整流程。假设我们在 PostgreSQL 的raw_schema中有两张原始表raw.customers和raw.orders。4.1 步骤一定义数据源创建models/staging/sources.ymlversion: 2 sources: - name: raw schema: raw_schema database: my_database # 如果与 profile 中一致可省略 tables: - name: customers description: 原始客户信息表 - name: orders description: 原始订单事实表4.2 步骤二创建基础模型1. 创建客户信息模型(models/staging/stg_customers.sql){{ config( materializedview, tags[staging] ) }} select id as customer_id, trim(first_name) as first_name, -- 去除首尾空格 trim(last_name) as last_name, lower(email) as email, -- 统一为小写 created_at, updated_at, -- 添加一个数据新鲜度检查日期 current_date as etl_date from {{ source(raw, customers) }} where email is not null and email like %%.% -- 简单的邮箱格式校验这里我们使用了tags配置可以为模型打上标签方便后续按标签分组运行或测试。2. 创建订单模型(models/staging/stg_orders.sql){{ config( materializedview, tags[staging] ) }} select id as order_id, user_id as customer_id, -- 注意原始表可能叫 user_id我们统一为 customer_id order_date, status, amount, current_date as etl_date from {{ source(raw, orders) }} where amount 0 -- 确保金额非负 and order_date current_date -- 订单日期不应在未来4.3 步骤三创建核心业务模型现在我们基于清洗后的数据创建一个面向分析的核心模型客户订单聚合表。创建models/marts/fct_customer_orders.sql{{ config( materializedtable, -- 核心业务模型物化为表以提升查询性能 tags[marts, core] ) }} with customer_orders as ( select customer_id, count(*) as number_of_orders, min(order_date) as first_order_date, max(order_date) as most_recent_order_date, sum(amount) as lifetime_value from {{ ref(stg_orders) }} where status completed -- 只计算已完成的订单 group by 1 ) select c.customer_id, c.first_name, c.last_name, c.email, co.number_of_orders, co.first_order_date, co.most_recent_order_date, co.lifetime_value, -- 计算客户活跃天数从首次订单到最近订单 (co.most_recent_order_date - co.first_order_date) as active_days, -- 判断是否为活跃客户最近30天内有订单 case when co.most_recent_order_date current_date - interval 30 days then true else false end as is_active_customer, c.etl_date from {{ ref(stg_customers) }} c left join customer_orders co on c.customer_id co.customer_id4.4 步骤四为模型添加测试与文档创建models/marts/schema.ymlversion: 2 models: - name: fct_customer_orders description: | 客户订单事实聚合表。 此模型连接客户与订单信息计算客户生命周期价值、订单数等核心指标。 是下游客户分析、报表的主要数据源。 columns: - name: customer_id description: 客户唯一标识关联 dim_customers tests: - unique - not_null - relationships: to: ref(stg_customers) field: customer_id - name: lifetime_value description: 客户历史累计消费总额 tests: - not_null - name: number_of_orders description: 客户历史总订单数 tests: - not_null - name: is_active_customer description: 是否为活跃客户最近30天有订单4.5 步骤五运行与测试现在让我们执行整个管道# 1. 运行所有模型dbt 会自动解析依赖关系按正确顺序执行 dbt run # 2. 运行所有测试 dbt test # 3. 生成项目文档网站一个可交互的静态站点 dbt docs generate dbt docs serve # 在本地启动一个服务器查看文档默认 http://localhost:8080执行dbt run后你会在数据库中看到创建了stg_customers(视图)、stg_orders(视图) 和fct_customer_orders(表)。dbt docs serve启动的站点会展示完整的 DAG有向无环图依赖关系、模型描述、列描述和测试结果数据血缘一目了然。5. 进阶技巧与避坑指南5.1 利用 Jinja 实现动态 SQLJinja 模板引擎是 dbt 强大灵活性的来源。除了config、ref、source你还可以使用控制结构{% set payment_methods [credit_card, coupon, bank_transfer, gift_card] %} select order_id, {% for payment_method in payment_methods %} sum(case when payment_method {{ payment_method }} then amount else 0 end) as {{ payment_method }}_amount {% if not loop.last %},{% endif %} {% endfor %} from {{ ref(stg_payments) }} group by 1这段代码会根据payment_methods列表动态生成列避免了手动编写重复的 CASE WHEN 语句。使用宏Macros封装可复用逻辑 在macros/目录下创建cents_to_dollars.sql{% macro cents_to_dollars(column_name, precision2) %} ({{ column_name }} / 100.0)::numeric(16, {{ precision }}) {% endmacro %}在模型中调用select order_id, {{ cents_to_dollars(amount_cents) }} as amount_dollars from {{ ref(stg_orders) }}5.2 物化策略选择与增量模型对于数据量大的表全量重建 (table) 成本高昂。这时需要使用增量模型 (incremental)。{{ config( materializedincremental, unique_keyid -- 增量去重的依据 ) }} select * from {{ source(raw, large_log_table) }} {% if is_incremental() %} -- 这是一个增量运行只处理新数据 where created_at (select max(created_at) from {{ this }}) {% endif %}关键点is_incremental(): dbt 提供的宏在增量运行时返回true。unique_key: 确保基于此键进行 upsert合并更新操作避免重复。{{ this }}: 引用当前模型本身在数据库中的表。实操心得增量模型是性能优化的利器但也是坑最多的地方。务必确保你的where条件能准确、可靠地过滤出新数据并且unique_key是真正唯一的。建议先在开发环境用小数据量充分测试增量逻辑。5.3 开发、测试与生产环境管理通过profiles.yml文件管理多环境配置。文件通常位于~/.dbt/下。# ~/.dbt/profiles.yml my_first_dbt_project: # 项目配置名 target: dev # 默认目标环境 outputs: dev: type: postgres host: localhost user: dev_user pass: dev_password port: 5432 dbname: analytics_dev schema: dbt_dev prod: type: postgres host: prod-db.company.com user: {{ env_var(DBT_PROD_USER) }} # 使用环境变量更安全 pass: {{ env_var(DBT_PROD_PASS) }} port: 5432 dbname: analytics schema: dbt_prod通过命令行切换环境dbt run --target prod # 在生产环境运行 dbt run --target dev # 在开发环境运行默认重要安全实践永远不要将密码等敏感信息硬编码在代码或配置文件中提交到版本库。使用环境变量如env_var或秘密管理工具。5.4 常见问题与排查实录问题1dbt run失败报错 “relation ‘xxx’ already exists”。原因通常是因为之前运行失败或中断导致表被部分创建而 dbt 尝试创建同名表时冲突。解决手动到数据库中删除该表或视图。或者在开发时可以先运行dbt run --full-refresh强制全量重建所有模型。检查模型配置的materialized是否正确是否与现有对象类型冲突。问题2测试失败但数据看起来没问题。原因可能是测试逻辑与业务逻辑不符或者数据中存在你未考虑到的边缘情况。排查运行dbt test --store-failures。这个命令会将测试失败的数据行插入到数据库的临时表中。直接查询这些临时表表名通常为test_项目名_测试名查看具体是哪些数据导致了测试失败。根据失败数据判断是测试条件过于严格需要调整还是数据本身确实存在质量问题。问题3模型运行顺序不符合预期。原因dbt 通过ref()函数自动解析依赖。如果运行顺序错乱很可能是因为某些模型间缺少ref()引用而是使用了硬编码的表名导致 dbt 无法识别依赖关系。解决检查所有 SQL 文件确保引用其他 dbt 模型时一律使用{{ ref(model_name) }}引用原始数据时使用{{ source(source_name, table_name) }}。问题4文档网站中的 DAG 图混乱或缺失。原因依赖关系解析错误或者dbt docs generate时使用的target/目录是旧的编译结果。解决运行dbt clean清理target/和dbt_packages/目录。重新运行dbt run或dbt compile。再执行dbt docs generate。6. 项目组织与团队协作最佳实践当项目从个人玩具发展为团队工具时良好的组织结构至关重要。1. 模块化目录结构不要把所有模型都扔在models/根目录下。建议按业务域或数据流阶段划分models/ ├── staging/ # 数据清洗层 │ ├── jaffle_shop/ # 来自 Jaffle Shop 数据源 │ └── stripe/ # 来自 Stripe 支付数据源 ├── intermediate/ # 中间模型复杂转换的中间步骤 ├── marts/ # 数据集市层面向业务 │ ├── finance/ │ ├── marketing/ │ └── product/ └── utils/ # 工具模型如日期维度表在dbt_project.yml中可以为不同目录配置不同的物化策略models: my_project: staging: materialized: view tags: staging marts: materialized: table tags: marts2. 使用 Tags 进行分组管理标签可以让你灵活地选择运行或测试一部分模型。# 只运行打上 finance 标签的模型 dbt run --select tag:finance # 运行所有标记为 staging 的模型及其下游依赖 dbt run --select tag:staging # 测试核心业务模型 dbt test --select tag:core3. 版本控制与代码审查将整个 dbt 项目除了target/和dbt_packages/纳入 Git 版本控制。使用 Pull Request 工作流进行代码变更。dbt 的模块化和测试框架使得代码审查可以聚焦于业务逻辑和数据质量。在 CI/CD 流水线中集成dbt compile和dbt test确保合并的代码不会破坏现有功能。4. 文档即代码把schema.yml文件中的描述当作重要的文档来写。清晰的列描述、业务定义和示例能极大降低团队的理解成本。生成的文档网站是数据资产最好的门户。从在命令行里敲下dbt init开始到构建起一个清晰、可测试、有文档的数据转换管道这个过程最深刻的体会是思维的转变。你不再是在编写一个个孤立的 SQL 脚本而是在构建一个相互关联、有生命力的数据产品。每一次dbt run都是对数据逻辑的一次部署每一次dbt test都是对数据质量的一次守护。刚开始可能会觉得 YAML 配置和 Jinja 语法有些繁琐但一旦习惯你会发现它带来的秩序感和可维护性是传统脚本方式无法比拟的。遇到问题时多利用dbt compile查看编译后的 SQL多使用dbt docs generate可视化依赖关系这两个命令是调试和理解项目的最佳助手。
返回列表