、CTE 與 dbt 實戰(zhàn)應用)
Data Engineering Zoomcamp 之 SQL 復習指南窗口函數(shù)、CTE 與 dbt 實戰(zhàn)應用【免費下載鏈接】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本篇技術(shù)指南是 Data Engineering Zoomcamp 課程第四模塊Analytics Engineering的 SQL 復習資料圍繞窗口函數(shù)Window Functions、公共表表達式CTE兩大核心主題展開并結(jié)合課程倉庫中的 dbt 項目taxi_rides_ny展示它們在真實數(shù)倉模型中的實戰(zhàn)用法。讀完本文你將掌握ROW_NUMBER()、RANK()、DENSE_RANK()、LAG()、LEAD()、PERCENTILE_CONT()的語法與語義差異理解如何用 CTE 組織多步驟分析邏輯并能在 dbt 模型中熟練組合這些 SQL 能力為后續(xù)構(gòu)建可維護的分析工程打下堅實基礎。一、為什么在進入 dbt 之前要先復習 SQL在開始第四模塊Analytics Engineering / dbt之前先系統(tǒng)回顧 SQL 中兩類進階能力窗口函數(shù)與 CTE。原因很直接dbt 模型本質(zhì)上是 SQLdbt 中的每個 model 都是一個.sql文件最終會被編譯成普通 SQL 在目標數(shù)倉中執(zhí)行窗口函數(shù)是清洗與去重的主力例如用ROW_NUMBER() OVER (PARTITION BY ...)實現(xiàn)按業(yè)務鍵去重、取每組最新記錄這類模式在 dbt 的 staging/intermediate 層極為常見CTE 是 dbt 模型的標準組織方式一個模型通常由多個 CTE 串成聲明式流水線可讀性、可測試性都優(yōu)于多層嵌套子查詢。也就是說本文復習的 SQL 能力將直接轉(zhuǎn)化為你在 dbt 中編寫生產(chǎn)級模型的日常工具。二、窗口函數(shù)Window Functions2.1 什么是窗口函數(shù)窗口函數(shù)Window Function在一組與當前行相關(guān)的行上執(zhí)行計算這個行集合稱為窗口window。從計算類型上看它和聚合函數(shù)SUM()、AVG()、COUNT()等很相似但關(guān)鍵區(qū)別在于普通聚合函數(shù)會把多行折疊成一行輸出而窗口函數(shù)不會——每一行都保留自己的獨立身份同時附帶上窗口計算的結(jié)果。基本語法FUNCTION() OVER (PARTITION BY column_name ORDER BY column_name)窗口函數(shù)由兩部分組成其中OVER (...)這一半定義了你的窗口OVER (PARTITION BY column_name ORDER BY column_name)PARTITION BY把結(jié)果集按列值劃分成若干組可選。函數(shù)在每個分區(qū)內(nèi)獨立計算ORDER BY定義分區(qū)內(nèi)各行的處理順序很多窗口函數(shù)的語義依賴這個順序如累計求和、排名、取前一行等。常見窗口函數(shù)分類類別函數(shù)說明排名類ROW_NUMBER()在分區(qū)內(nèi)為每一行分配唯一的行號排名類RANK()類似ROW_NUMBER()但并列值取相同名次且后續(xù)名次會跳過排名類DENSE_RANK()類似RANK()并列取相同名次但名次不跳過、連續(xù)編號聚合類SUM() OVER()計算運行總計running total聚合類AVG() OVER()計算移動平均moving average前后行類LAG()取分區(qū)內(nèi)前一行或往前 N 行的值前后行類LEAD()取分區(qū)內(nèi)后一行或往后 N 行的值2.2 ROW_NUMBER()為每一行編號ROW_NUMBER()為每一行分配一個從 1 開始的連續(xù)編號編號順序由窗口中的ORDER BY決定如果使用了PARTITION BY則每個分區(qū)內(nèi)都從 1 重新開始計數(shù)。語法ROW_NUMBER() OVER (PARTITION BY column_name ORDER BY column_name)常見用途去重先用ROW_NUMBER()給每組按業(yè)務鍵分區(qū)的行編號再只保留編號為 1 的行排名需要唯一名次、不允許并列時使用取每組最新記錄按實體分組、按時間倒序編號后取第 1 行。示例 1不分區(qū)直接按金額排名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()生成的是臨時計算列只存在于查詢結(jié)果中不會修改原始表。示例 2按上車地點分區(qū)后組內(nèi)排名SELECT total_amount, PULocationID, ROW_NUMBER() OVER (PARTITION BY PULocationID ORDER BY total_amount DESC) AS ranking FROM greentaxi_trips LIMIT 10;該查詢在每個PULocationID分組內(nèi)按total_amount降序重新從 1 開始編號total_amountPULocationIDranking8.512244328.32244338.32244347.32244353.322443686.42234173.5234262.7234361.94234461.942345注意上表第一組因為LIMIT 10截斷了結(jié)果位置 224 的組內(nèi)編號從 432 開始——這恰恰說明了PARTITION BY讓每個分區(qū)獨立計數(shù)的行為如果去掉分區(qū)ROW_NUMBER()會在全表范圍內(nèi)連續(xù)編號。2.3 RANK() 與 DENSE_RANK()并列值如何處理ROW_NUMBER()、RANK()、DENSE_RANK()都按指定順序給行分配名次但在出現(xiàn)并列值時行為不同RANK()并列值取相同名次后續(xù)名次跳過跳過的數(shù)量取決于并列行數(shù)DENSE_RANK()并列值取相同名次后續(xù)名次連續(xù)不跳過。對比示例ScoreROW_NUMBER()RANK()DENSE_RANK()95111902229032285443可以看到兩個 90 分并列第 2 名RANK()的下一個名次跳到 4跳過 3而DENSE_RANK()則給 85 分排第 3 名。選擇哪個取決于業(yè)務語義需要名次連續(xù)如第 1、2、3 名用DENSE_RANK()需要標準競賽式名次如體育賽事有并列時自動空缺名次用RANK()需要逐行唯一編號用ROW_NUMBER()。2.4 LAG() 與 LEAD()訪問前一行/后一行業(yè)務中經(jīng)常需要把當前行與相鄰行做比較例如上一次乘車金額是多少下一次乘車金額是多少。LAG()和LEAD()可以在不進行自連接self-join的情況下直接把其他行的值拉到當前行。LAG()取之前的行LEAD()取之后的行。語法LAG(expression) OVER (PARTITION BY partition_expression ORDER BY order_expression)expression要取值的目標列offset可選往回或往后多少行默認 1即緊鄰的上一行/下一行PARTITION BY可選把結(jié)果集劃分成多個分區(qū)在每個分區(qū)內(nèi)獨立計算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這類相鄰行對照的寫法是環(huán)比/同比分析、事件序列分析的基石且避免了昂貴的自連接。2.5 PERCENTILE_CONT()線性插值計算百分位PERCENTILE_CONT()對指定的value_expression計算指定百分位數(shù)值采用線性插值linear interpolation方式因此即使百分位不恰好落在某個數(shù)據(jù)點上也能給出連續(xù)、平滑的分位數(shù)估計。語法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 分位。結(jié)果示意分區(qū)內(nèi)每行都會帶上該分區(qū)的 p90PULocationIDtotal_amountp9022417.351.922420.6751.92242151.922426.0651.922427.1351.922440.1451.922455.4651.922425.7451.922427.0251.92243751.9p90 的含義是90% 的數(shù)據(jù)都落在此值之下。上表中位置 224 的 p90 恒為 51.9說明該地點 90% 的行程金額低于 51.9。這類分位數(shù)指標在異常檢測、定價分析、服務質(zhì)量監(jiān)控中非常常用。三、公共表表達式Common Table Expression, CTE3.1 什么是 CTECTE 可以理解為查詢中的查詢。通過WITH語句你可以先創(chuàng)建若干臨時結(jié)果表再在后續(xù)查詢中引用它們讓復雜查詢變得更可讀、更易維護。這些臨時表只存在于當前這條主查詢的執(zhí)行期間。CTE 與子查詢subquery都能實現(xiàn)類似目標但各有側(cè)重可復用性CTE 可以在一條查詢內(nèi)被多次引用若多個查詢都需要同一邏輯還可以把該邏輯沉淀到視圖中可讀性CTE 把多步計算攤平成自上而下的命名步驟閱讀者能清晰把握分析脈絡。把 CTE 聲明在查詢開頭代碼的可讀性會大幅提升也讓分析邏輯更容易被別人以及未來的你理解。基本語法WITH cte_name AS ( SELECT column1, column2 FROM some_table WHERE condition ) SELECT * FROM cte_name;3.2 CTE 實戰(zhàn)找出金額第二大的行程示例找到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 與窗口函數(shù)的經(jīng)典組合在 CTEcte中用RANK() OVER (ORDER BY total_amount DESC)按金額從高到低給每行分配名次主查詢從cte中篩選rank 2的行即金額第二大的行程。結(jié)果lpep_pickup_datetimetotal_amountrank2019-10-10 15:22:49 UTC2878.32值得強調(diào)的是這里用RANK()而非ROW_NUMBER()的語義差異若存在并列第一RANK()下第二名是真正意義上的第二高值而ROW_NUMBER()只會機械地取第 2 行。選擇哪種取決于業(yè)務定義。四、dbt 模型中的 CTE 與窗口函數(shù)從復習到實戰(zhàn)CTE 與窗口函數(shù)在第四模塊的 dbt 課程中會被大量使用。下面以倉庫中的真實代碼為例展示它們?nèi)绾伪唤M織進 dbt 模型。4.1 模型示例基于 FHV 數(shù)據(jù)計算行程時長與 90 分位假設從 FHV 數(shù)據(jù)集出發(fā)要創(chuàng)建一個 dbt 模型為數(shù)據(jù)補充行程時長和行程時長 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子句創(chuàng)建了名為trip_duration_calculated的 CTE它相當于一張臨時表包含fhv_trips的全部列并額外用timestamp_diff(...)計算出每趟行程的時長trip_duration單位秒。第二步主查詢中組合 CTE 與窗口函數(shù)。外層SELECT對 CTE 結(jié)果按PUlocationID分區(qū)用PERCENTILE_CONT(trip_duration, 0.90)計算每個上車地點的行程時長 90 分位。PARTITION BY PUlocationID保證分位數(shù)按地點獨立計算分位 90 表示 90% 的行程時長小于等于該值。結(jié)果示意PUlocationIDtrip_durationtrip_duration_p901904512170.019013732170.01908172170.01905892170.019016482170.0325461988.0321511988.03217521988.03224261988.0328881988.0對PUlocationID 19090% 的行程時長 ≤ 2170.0 秒對PUlocationID 3290% 的行程時長 ≤ 1988.0 秒。4.2 倉庫實戰(zhàn)印證一窗口函數(shù)在去重與生成主鍵中的應用上面的 CTE 寫法不是孤立示例倉庫中 taxi_rides_ny/models/intermediate/int_trips.sql 就是一個把CTE 組織 窗口函數(shù)去重用到極致的真實模型。該模型在多個 CTE 之上做清洗與富化后最后用QUALIFY結(jié)合窗口函數(shù)去重-- 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 節(jié)用ROW_NUMBER()識別重復并保留一行的工程化落地按vendor_id pickup_datetime pickup_location_id service_type分組編號每組只保留dropoff_datetime最早row_number() 1的一條記錄。注意QUALIFY是 BigQuery 等數(shù)倉對HAVING的補充專門用于在窗口函數(shù)計算之后過濾行比先包一層子查詢更簡潔。4.3 倉庫實戰(zhàn)印證二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 構(gòu)成了一條清晰的轉(zhuǎn)換鏈先聲明數(shù)據(jù)來源{{ source(...) }}再做類型轉(zhuǎn)換與重命名。同樣地stg_yellow_tripdata.sql 與 int_trips_unioned.sql 延續(xù)了同一模式——后者還通過union all把綠、黃兩套出租車數(shù)據(jù)合并并各自打上service_type標簽體現(xiàn)了 CTE 在跨源整合中的價值。4.4 倉庫實戰(zhàn)印證三窗口/聚合能力與宏的配合在 fct_trips.sql 中可以看到行程時長不再手寫timestamp_diff而是調(diào)用倉庫自定義宏get_trip_duration_minutes{{ get_trip_duration_minutes(trips.pickup_datetime, trips.dropoff_datetime) }} as trip_duration_minutes,該宏定義在 macros/get_trip_duration_minutes.sql內(nèi)部使用 dbt 內(nèi)置的跨數(shù)據(jù)庫datediff宏{% macro get_trip_duration_minutes(pickup_datetime, dropoff_datetime) %} {{ dbt.datediff(pickup_datetime, dropoff_datetime, minute) }} {% endmacro %}這個細節(jié)揭示了一個重要工程思想窗口函數(shù)、日期函數(shù)等 SQL 能力在 dbt 中被抽象成可復用、跨數(shù)倉DuckDB、BigQuery、Snowflake、Redshift、PostgreSQL 等的宏。同理macros/safe_cast.sql 用target.type判斷在 BigQuery 上用safe_cast、其他數(shù)倉用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用于月度營收報表——這些聚合正是窗口函數(shù)背后的非窗口版計算兩者互補使用。五、小結(jié)把 SQL 能力轉(zhuǎn)化為 dbt 生產(chǎn)力回顧全文四條主線貫穿始終窗口函數(shù)不折疊行的特性讓每行在保留自身的同時附帶組內(nèi)排名、累計值、前后行值或分位數(shù)是清洗去重、分析排名、百分位、序列比較LAG/LEAD的利器ROW_NUMBER()/RANK()/DENSE_RANK()的并列語義差異需要按業(yè)務規(guī)則審慎選擇CTE 通過WITH組織多步邏輯讓復雜查詢可讀、可復用是 dbt 模型的標準組織方式在 dbt 中窗口函數(shù)與 CTE 與宏、QUALIFY、跨數(shù)倉適配等技術(shù)結(jié)合構(gòu)成生產(chǎn)級模型的核心語法基座。對照倉庫中 staging、intermediate、marts 三層的模型文件你會發(fā)現(xiàn)每一層都在復用本文的兩種語法。掌握它們就等于掌握了進入 Analytics Engineering 世界的鑰匙——接下來的 dbt 模塊將圍繞這些模型展開更系統(tǒng)的工程化實踐。【免費下載鏈接】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創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考