
Data Engineering Zoomcamp 之 SQL 复习指南窗口函数、CTE 与 dbt 实战应用【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp本篇技术指南是 Data Engineering Zoomcamp 课程第四模块Analytics Engineering的 SQL 复习资料围绕窗口函数Window Functions、公共表表达式CTE两大核心主题展开并结合课程仓库中的 dbt 项目taxi_rides_ny展示它们在真实数仓模型中的实战用法。读完本文你将掌握ROW_NUMBER()、RANK()、DENSE_RANK()、LAG()、LEAD()、PERCENTILE_CONT()的语法与语义差异理解如何用 CTE 组织多步骤分析逻辑并能在 dbt 模型中熟练组合这些 SQL 能力为后续构建可维护的分析工程打下坚实基础。一、为什么在进入 dbt 之前要先复习 SQL在开始第四模块Analytics Engineering / dbt之前先系统回顾 SQL 中两类进阶能力窗口函数与 CTE。原因很直接dbt 模型本质上是 SQLdbt 中的每个 model 都是一个.sql文件最终会被编译成普通 SQL 在目标数仓中执行窗口函数是清洗与去重的主力例如用ROW_NUMBER() OVER (PARTITION BY ...)实现按业务键去重、取每组最新记录这类模式在 dbt 的 staging/intermediate 层极为常见CTE 是 dbt 模型的标准组织方式一个模型通常由多个 CTE 串成声明式流水线可读性、可测试性都优于多层嵌套子查询。也就是说本文复习的 SQL 能力将直接转化为你在 dbt 中编写生产级模型的日常工具。二、窗口函数Window Functions2.1 什么是窗口函数窗口函数Window Function在一组与当前行相关的行上执行计算这个行集合称为窗口window。从计算类型上看它和聚合函数SUM()、AVG()、COUNT()等很相似但关键区别在于普通聚合函数会把多行折叠成一行输出而窗口函数不会——每一行都保留自己的独立身份同时附带上窗口计算的结果。基本语法FUNCTION() OVER (PARTITION BY column_name ORDER BY column_name)窗口函数由两部分组成其中OVER (...)这一半定义了你的窗口OVER (PARTITION BY column_name ORDER BY column_name)PARTITION BY把结果集按列值划分成若干组可选。函数在每个分区内独立计算ORDER BY定义分区内各行的处理顺序很多窗口函数的语义依赖这个顺序如累计求和、排名、取前一行等。常见窗口函数分类类别函数说明排名类ROW_NUMBER()在分区内为每一行分配唯一的行号排名类RANK()类似ROW_NUMBER()但并列值取相同名次且后续名次会跳过排名类DENSE_RANK()类似RANK()并列取相同名次但名次不跳过、连续编号聚合类SUM() OVER()计算运行总计running total聚合类AVG() OVER()计算移动平均moving average前后行类LAG()取分区内前一行或往前 N 行的值前后行类LEAD()取分区内后一行或往后 N 行的值2.2 ROW_NUMBER()为每一行编号ROW_NUMBER()为每一行分配一个从 1 开始的连续编号编号顺序由窗口中的ORDER BY决定如果使用了PARTITION BY则每个分区内都从 1 重新开始计数。语法ROW_NUMBER() OVER (PARTITION BY column_name ORDER BY column_name)常见用途去重先用ROW_NUMBER()给每组按业务键分区的行编号再只保留编号为 1 的行排名需要唯一名次、不允许并列时使用取每组最新记录按实体分组、按时间倒序编号后取第 1 行。示例 1不分区直接按金额排名SELECT total_amount, ROW_NUMBER() OVER (ORDER BY total_amount DESC) AS ranking FROM greentaxi_trips LIMIT 10;该查询返回表中total_amount最高的 10 行并为每行附带一个表示排名的行号total_amountranking4012.312878.322438.832156.342109.852017.361971.0571958.881762.891600.810需要注意的是ROW_NUMBER()生成的是临时计算列只存在于查询结果中不会修改原始表。示例 2按上车地点分区后组内排名SELECT total_amount, PULocationID, ROW_NUMBER() OVER (PARTITION BY PULocationID ORDER BY total_amount DESC) AS ranking FROM greentaxi_trips LIMIT 10;该查询在每个PULocationID分组内按total_amount降序重新从 1 开始编号total_amountPULocationIDranking8.512244328.32244338.32244347.32244353.322443686.42234173.5234262.7234361.94234461.942345注意上表第一组因为LIMIT 10截断了结果位置 224 的组内编号从 432 开始——这恰恰说明了PARTITION BY让每个分区独立计数的行为如果去掉分区ROW_NUMBER()会在全表范围内连续编号。2.3 RANK() 与 DENSE_RANK()并列值如何处理ROW_NUMBER()、RANK()、DENSE_RANK()都按指定顺序给行分配名次但在出现并列值时行为不同RANK()并列值取相同名次后续名次跳过跳过的数量取决于并列行数DENSE_RANK()并列值取相同名次后续名次连续不跳过。对比示例ScoreROW_NUMBER()RANK()DENSE_RANK()95111902229032285443可以看到两个 90 分并列第 2 名RANK()的下一个名次跳到 4跳过 3而DENSE_RANK()则给 85 分排第 3 名。选择哪个取决于业务语义需要名次连续如第 1、2、3 名用DENSE_RANK()需要标准竞赛式名次如体育赛事有并列时自动空缺名次用RANK()需要逐行唯一编号用ROW_NUMBER()。2.4 LAG() 与 LEAD()访问前一行/后一行业务中经常需要把当前行与相邻行做比较例如上一次乘车金额是多少下一次乘车金额是多少。LAG()和LEAD()可以在不进行自连接self-join的情况下直接把其他行的值拉到当前行。LAG()取之前的行LEAD()取之后的行。语法LAG(expression) OVER (PARTITION BY partition_expression ORDER BY order_expression)expression要取值的目标列offset可选往回或往后多少行默认 1即紧邻的上一行/下一行PARTITION BY可选把结果集划分成多个分区在每个分区内独立计算ORDER BY定义行处理顺序决定前/后的参照系。示例把相邻行程的金额串起来SELECT lpep_pickup_datetime, total_amount, LAG(total_amount) OVER (ORDER BY lpep_pickup_datetime) as prev_total_amount, LEAD(total_amount) OVER (ORDER BY lpep_pickup_datetime) as next_total_amount FROM greentaxi_trips ORDER BY lpep_pickup_datetime查询返回每一趟行程的上车时间、金额以及按时间排序后的上一趟金额和下一趟金额lpep_pickup_datetimetotal_amountprev_total_amountnext_total_amount2008-12-31 23:33:38 UTC7.36.35.32008-12-31 23:42:31 UTC5.37.314.552008-12-31 23:47:51 UTC14.555.319.552008-12-31 23:57:46 UTC19.5514.559.82009-01-01 00:00:00 UTC9.819.5581.3这类相邻行对照的写法是环比/同比分析、事件序列分析的基石且避免了昂贵的自连接。2.5 PERCENTILE_CONT()线性插值计算百分位PERCENTILE_CONT()对指定的value_expression计算指定百分位数值采用线性插值linear interpolation方式因此即使百分位不恰好落在某个数据点上也能给出连续、平滑的分位数估计。语法PERCENTILE_CONT(value_expression, percentile) OVER (PARTITION BY partition_expression)示例计算每个上车地点的 90 分位金额SELECT PULocationID, total_amount, PERCENTILE_CONT(total_amount, 0.9) OVER (PARTITION BY PULocationID) AS p90 FROM greentaxi_tripsPERCENTILE_CONT(total_amount, 0.9)计算total_amount的 90 分位p90即 90% 的金额低于该值PARTITION BY PULocationID按上车地点分组每个地点独立计算各自的 90 分位。结果示意分区内每行都会带上该分区的 p90PULocationIDtotal_amountp9022417.351.922420.6751.92242151.922426.0651.922427.1351.922440.1451.922455.4651.922425.7451.922427.0251.92243751.9p90 的含义是90% 的数据都落在此值之下。上表中位置 224 的 p90 恒为 51.9说明该地点 90% 的行程金额低于 51.9。这类分位数指标在异常检测、定价分析、服务质量监控中非常常用。三、公共表表达式Common Table Expression, CTE3.1 什么是 CTECTE 可以理解为查询中的查询。通过WITH语句你可以先创建若干临时结果表再在后续查询中引用它们让复杂查询变得更可读、更易维护。这些临时表只存在于当前这条主查询的执行期间。CTE 与子查询subquery都能实现类似目标但各有侧重可复用性CTE 可以在一条查询内被多次引用若多个查询都需要同一逻辑还可以把该逻辑沉淀到视图中可读性CTE 把多步计算摊平成自上而下的命名步骤阅读者能清晰把握分析脉络。把 CTE 声明在查询开头代码的可读性会大幅提升也让分析逻辑更容易被别人以及未来的你理解。基本语法WITH cte_name AS ( SELECT column1, column2 FROM some_table WHERE condition ) SELECT * FROM cte_name;3.2 CTE 实战找出金额第二大的行程示例找到total_amount第二大的行程WITH cte AS ( SELECT lpep_pickup_datetime, total_amount, RANK() OVER (ORDER BY total_amount DESC) AS rank FROM greentaxi_trips ) SELECT * FROM cte WHERE rank 2;这个查询展示了 CTE 与窗口函数的经典组合在 CTEcte中用RANK() OVER (ORDER BY total_amount DESC)按金额从高到低给每行分配名次主查询从cte中筛选rank 2的行即金额第二大的行程。结果lpep_pickup_datetimetotal_amountrank2019-10-10 15:22:49 UTC2878.32值得强调的是这里用RANK()而非ROW_NUMBER()的语义差异若存在并列第一RANK()下第二名是真正意义上的第二高值而ROW_NUMBER()只会机械地取第 2 行。选择哪种取决于业务定义。四、dbt 模型中的 CTE 与窗口函数从复习到实战CTE 与窗口函数在第四模块的 dbt 课程中会被大量使用。下面以仓库中的真实代码为例展示它们如何被组织进 dbt 模型。4.1 模型示例基于 FHV 数据计算行程时长与 90 分位假设从 FHV 数据集出发要创建一个 dbt 模型为数据补充行程时长和行程时长 90 分位两个字段WITH trip_duration_calculated AS ( SELECT *, timestamp_diff(dropOff_datetime, pickup_datetime, second) as trip_duration FROM fhv_trips ) SELECT PUlocationID, trip_duration, PERCENTILE_CONT(trip_duration, 0.90) OVER (PARTITION BY PUlocationID) AS trip_duration_p90 FROM trip_duration_calculated第一步理解 CTE。WITH子句创建了名为trip_duration_calculated的 CTE它相当于一张临时表包含fhv_trips的全部列并额外用timestamp_diff(...)计算出每趟行程的时长trip_duration单位秒。第二步主查询中组合 CTE 与窗口函数。外层SELECT对 CTE 结果按PUlocationID分区用PERCENTILE_CONT(trip_duration, 0.90)计算每个上车地点的行程时长 90 分位。PARTITION BY PUlocationID保证分位数按地点独立计算分位 90 表示 90% 的行程时长小于等于该值。结果示意PUlocationIDtrip_durationtrip_duration_p901904512170.019013732170.01908172170.01905892170.019016482170.0325461988.0321511988.03217521988.03224261988.0328881988.0对PUlocationID 19090% 的行程时长 ≤ 2170.0 秒对PUlocationID 3290% 的行程时长 ≤ 1988.0 秒。4.2 仓库实战印证一窗口函数在去重与生成主键中的应用上面的 CTE 写法不是孤立示例仓库中 taxi_rides_ny/models/intermediate/int_trips.sql 就是一个把CTE 组织 窗口函数去重用到极致的真实模型。该模型在多个 CTE 之上做清洗与富化后最后用QUALIFY结合窗口函数去重-- Deduplicate: if multiple trips match (same vendor, second, location, service), keep first qualify row_number() over( partition by vendor_id, pickup_datetime, pickup_location_id, service_type order by dropoff_datetime ) 1这里的逻辑正是本文 2.2 节用ROW_NUMBER()识别重复并保留一行的工程化落地按vendor_id pickup_datetime pickup_location_id service_type分组编号每组只保留dropoff_datetime最早row_number() 1的一条记录。注意QUALIFY是 BigQuery 等数仓对HAVING的补充专门用于在窗口函数计算之后过滤行比先包一层子查询更简洁。4.3 仓库实战印证二CTE 作为 dbt 模型的标准组织方式dbt 模型普遍采用多个 CTE 串行的声明式写法。以 stg_green_tripdata.sql 为例with source as ( select * from {{ source(raw, green_tripdata) }} ), renamed as ( select cast(vendorid as integer) as vendor_id, cast(pulocationid as integer) as pickup_location_id, cast(lpep_pickup_datetime as timestamp) as pickup_datetime, ... from source where vendorid is not null ) select * from renamedsource→renamed两个 CTE 构成了一条清晰的转换链先声明数据来源{{ source(...) }}再做类型转换与重命名。同样地stg_yellow_tripdata.sql 与 int_trips_unioned.sql 延续了同一模式——后者还通过union all把绿、黄两套出租车数据合并并各自打上service_type标签体现了 CTE 在跨源整合中的价值。4.4 仓库实战印证三窗口/聚合能力与宏的配合在 fct_trips.sql 中可以看到行程时长不再手写timestamp_diff而是调用仓库自定义宏get_trip_duration_minutes{{ get_trip_duration_minutes(trips.pickup_datetime, trips.dropoff_datetime) }} as trip_duration_minutes,该宏定义在 macros/get_trip_duration_minutes.sql内部使用 dbt 内置的跨数据库datediff宏{% macro get_trip_duration_minutes(pickup_datetime, dropoff_datetime) %} {{ dbt.datediff(pickup_datetime, dropoff_datetime, minute) }} {% endmacro %}这个细节揭示了一个重要工程思想窗口函数、日期函数等 SQL 能力在 dbt 中被抽象成可复用、跨数仓DuckDB、BigQuery、Snowflake、Redshift、PostgreSQL 等的宏。同理macros/safe_cast.sql 用target.type判断在 BigQuery 上用safe_cast、其他数仓用cast保证同一份模型在不同平台上行为一致。如果你在本地用 DuckDB 跑这套项目可以按 setup/local_setup.md 的指引安装dbt-duckdb、配置profiles.yml并运行dbt debug验证连接。此外macros/get_vendor_data.sql 展示了宏如何在编译期生成CASE表达式把vendor_id映射为厂商名称而 models/marts/reporting/fct_monthly_zone_revenue.sql 则把group by聚合sum、count、avg用于月度营收报表——这些聚合正是窗口函数背后的非窗口版计算两者互补使用。五、小结把 SQL 能力转化为 dbt 生产力回顾全文四条主线贯穿始终窗口函数不折叠行的特性让每行在保留自身的同时附带组内排名、累计值、前后行值或分位数是清洗去重、分析排名、百分位、序列比较LAG/LEAD的利器ROW_NUMBER()/RANK()/DENSE_RANK()的并列语义差异需要按业务规则审慎选择CTE 通过WITH组织多步逻辑让复杂查询可读、可复用是 dbt 模型的标准组织方式在 dbt 中窗口函数与 CTE 与宏、QUALIFY、跨数仓适配等技术结合构成生产级模型的核心语法基座。对照仓库中 staging、intermediate、marts 三层的模型文件你会发现每一层都在复用本文的两种语法。掌握它们就等于掌握了进入 Analytics Engineering 世界的钥匙——接下来的 dbt 模块将围绕这些模型展开更系统的工程化实践。【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考