ARTICLE DETAIL

资讯详情

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

data-engineering-zoomcamp SQL 复习指南:窗口函数、CTE 与 dbt 模型实战

data-engineering-zoomcamp SQL 复习指南:窗口函数、CTE 与 dbt 模型实战 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本篇是>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()编号、LAG/LEAD取值方向以及累计类函数的计算方向。常见窗口函数排名函数Ranking FunctionsROW_NUMBER()在分区内为每一行分配唯一的行号。RANK()与ROW_NUMBER()类似但遇到并列值会给相同排名并跳过后续序号。DENSE_RANK()与RANK()类似但并列时不产生序号空洞。作为窗口函数使用的聚合函数SUM() OVER()计算运行总计running total。AVG() OVER()计算移动平均值。前后行取值函数Lag / LeadLAG()取当前行之前某行的值。LEAD()取当前行之后某行的值。Row Number 行号ROW_NUMBER()名副其实——为给定行显示序号。它从 1 开始按照窗口语句中的ORDER BY部分对行进行编号使用PARTITION BY子句时会在每个分区内重新从 1 开始计数。语法ROW_NUMBER() OVER (PARTITION BY column_name ORDER BY column_name)常见用途去重用ROW_NUMBER()识别重复行过滤掉行号大于 1 的行只保留一条。数据排名按特定条件排名但要求每行有唯一序号。取最新记录结合PARTITION BY选出每个类别中最新的记录。示例 1按 total_amount 降序取前 10SELECT total_amount, ROW_NUMBER() OVER (ORDER BY total_amount DESC) AS ranking FROM greentaxi_trips LIMIT 10;该查询返回表中最高的 10 个total_amount值并附带表示排名的行号total_amountranking4012.312878.322438.832156.342109.852017.361971.0571958.881762.891600.810用ROW_NUMBER()生成的列是临时计算列不会修改原表它只是应用于查询结果数据的一次计算。示例 2按 PULocationID 分区排名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降序为每行分配排名total_amountPULocationIDranking8.512244328.32244338.32244347.32244353.322443686.42234173.5234262.7234361.94234461.942345注意PULocationID 224的排名从 432 开始说明该分区内已有大量记录被先行编号而PULocationID 234的分区从 1 重新开始计数——这就是PARTITION BY让每个分区独立编号的直观体现。Rank 与 Dense RankROW_NUMBER()、RANK()和DENSE_RANK()都是按指定顺序为行分配排名的窗口函数但在排名列存在重复值时行为不同。RANK()遇到并列会跳过序号1、2、2、4。DENSE_RANK()与RANK()类似但并列时不跳过序号1、2、2、3。示例对比Score 相同得 90 的两行并列ScoreROW_NUMBER()RANK()DENSE_RANK()95111902229032285443实际选择建议需要唯一编号如后续用于过滤、去重用ROW_NUMBER()业务上希望“并列共享同一名次”时用RANK()或DENSE_RANK()其中希望后续排名连续如“Top 3”分组用DENSE_RANK()更合适。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在真实分析中这类“环比/环比”取值常用于计算相邻时段指标变化率、识别异常跳变等场景。Percentile Cont 百分位计算PERCENTILE_CONT对value_expression计算指定百分位的值采用线性插值方式因此结果可能不是数据集中实际存在的值。语法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。PARTITION BY PULocationID按上车地点分组使 90 分位在每个地点分别计算。查询结果示意PULocationIDtotal_amountp9022417.351.922420.6751.92242151.922426.0651.922427.1351.922440.1451.922455.4651.922425.7451.922427.0251.92243751.9p90 的含义是“90% 的值低于该数值”。上表中 p90 恒为 51.9说明对于地点22490% 的行程总费用低于 51.9。这类分位指标在计费、风控与服务水平分析中非常常用。Common Table Expression 公用表表达式CTECommon Table Expression公用表表达式可以理解为“查询中的查询”。借助WITH语句你可以创建临时表来暂存中间结果让复杂查询更易读、更易维护。这些临时表只在主查询执行期间存在。CTE 与子查询都是强大的工具能达到相似目的但各有适用场景。两者的主要区别在于CTE 在整条查询期间可复用且可读性更好。把 CTE 声明在查询开头能显著提升代码可读性让分析逻辑一目了然。语法WITH cte_name AS ( SELECT column1, column2 FROM some_table WHERE condition ) SELECT * FROM cte_name;示例找出金额第二大的行程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的 CTE通过RANK()窗口函数按total_amount降序为每行分配排名随后在主查询中引用该 CTE 并过滤rank 2。这里使用RANK()而非ROW_NUMBER()的意义在于若有多笔并列最高金额RANK()会把它们都排为 1而紧随其后的才是真正的“第二”名次。查询结果lpep_pickup_datetimetotal_amountrank2019-10-10 15:22:49 UTC2878.32dbt models 和 CTECTE 与窗口函数在第 4 模块dbt中会被大量使用。dbt 模型本身就是“以 CTE 组织逻辑的 SELECT 语句”因此本节先看原文档的入门示例再深入本仓库真实模型展示完整的实战形态。入门示例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第 1 步理解 CTEWITH子句创建名为trip_duration_calculated的 CTE它相当于一张临时表包含fhv_trips的全部列并额外用timestamp_diff(dropOff_datetime, pickup_datetime, second)以秒为单位计算每趟行程的时长。第 2 步主查询使用 CTE 与窗口函数主查询计算每个PUlocationID的行程时长 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 秒。仓库实战CTE 窗口函数 宏的完整数据链路在 04-analytics-engineering/taxi_rides_ny 这个 dbt 项目中上述技巧贯穿了从 staging 到 marts 的全过程可以从源码中逐一印证。1. staging 层CTE 组织重命名与清洗stg_yellow_tripdata.sql 与 stg_green_tripdata.sql 都用标准的两段式 CTEsource读取{{ source(raw, yellow_tripdata) }}等原始数据renamed统一列名与类型如cast(vendorid as integer) as vendor_id最后select * from renamed并额外做了vendorid is not null的数据质量过滤以及 dev 环境的日期采样{% if target.name dev %} where pickup_datetime 2019-01-01 and pickup_datetime 2019-02-01 {% endif %}2. intermediate 层CTE union all合并多源int_trips_unioned.sql 用green_trips/yellow_trips两个 CTE 将绿、黄出租车数据字段对齐例如给黄色出租车补上cast(0 as numeric) as ehail_fee、cast(1 as integer) as trip_type并添加Green/Yellow的service_type标签最后union all合并——这正是 CTE“可读、可复用”的典型应用。3. intermediate 层qualify row_number() over(...)去重int_trips.sql 演示了原文档提到的 ROW_NUMBER 头号用途——去重。它在 CTEcleaned_and_enriched中生成代理键并关联支付类型描述后用 BigQuery 的QUALIFY子句配合ROW_NUMBER()保留每组唯一记录select * from cleaned_and_enriched -- 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按vendor_id pickup_datetime pickup_location_id service_type分区并取行号为 1 的行即“同一秒内同一地点同一服务的重复行程只保留第一条”——这是生产级数据管道中典型的幂等去重写法。4. marts 层窗口函数式的宏 CTE 装配事实表fct_trips.sql 将行程事实表与dim_zones维度表做left join丰富上/下车行政区与区域名并调用自定义宏计算行程时长{{ get_trip_duration_minutes(trips.pickup_datetime, trips.dropoff_datetime) }} as trip_duration_minutes宏定义位于 get_trip_duration_minutes.sql内部使用 dbt 内置的跨数据库datediff{% macro get_trip_duration_minutes(pickup_datetime, dropoff_datetime) %} {{ dbt.datediff(pickup_datetime, dropoff_datetime, minute) }} {% endmacro %}它可以在 DuckDB、BigQuery、Snowflake、Redshift、PostgreSQL 等不同平台上无缝运行与入门示例中手写timestamp_diff的思路一致但通过宏做到了跨方言复用。对应的列类型、非空与取值约束在 marts/schema.yml 中以data_tests如unique、not_null、accepted_values、relationships声明便于在模型运行时自动校验。5. reporting 层窗口/聚合知识向报表聚合迁移最后fct_monthly_zone_revenue.sql 从fct_trips出发用sum()、count()、avg()配合group by pickup_zone, revenue_month, service_type产出“按上车区域 月份 服务类型”的月度营收报表并用coalesce(pickup_zone, Unknown Zone)兜底缺失区域同时针对bigquery与duckdb用if / elif适配了不同的date_trunc写法。这体现了掌握聚合函数与GROUP BY语义后如何从明细模型逐层聚合出面向报表与仪表盘的宽表。小结与进阶建议回顾本篇要点窗口函数OVER (PARTITION BY ... ORDER BY ...)在保留明细行的同时提供分组计算ROW_NUMBER()去重/唯一编号RANK()/DENSE_RANK()处理并列排名LAG()/LEAD()免自连接取前后行PERCENTILE_CONT()计算线性插值分位。CTE用WITH把复杂查询拆解为可读、可复用的临时逻辑块是子查询更清晰的组织方式。dbt 中的落地本仓库的 staging → intermediate → marts 分层模型是 CTE、窗口函数、宏与测试配置组合的完整范本。建议下一步直接打开 04-analytics-engineering/taxi_rides_ny 逐层阅读上述模型文件尝试将本篇的PERCENTILE_CONT与QUALIFY写法迁移到自己的报表查询中并结合 SQL.md 原文 反复对照直到能够不看示例独立写出多 CTE 窗口函数的查询。【免费下载链接】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),仅供参考
返回列表