Răsfoiți Sursa

perf(ods/trd): ful_d merge 改反连接,避免每天全量 shuffle

原用 ROW_NUMBER() OVER (PARTITION BY id) 套在 full∪inc 上,等于每天把
6200 万行全量按 id 全 shuffle + 排序,只为合入 20 万行增量。改为
「昨日全量 LEFT ANTI JOIN 今日增量 id」+「今日增量」:右侧万级 id 走
broadcast,全量表只扫不 shuffle、不排序。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
tianyu.chu 2 săptămâni în urmă
părinte
comite
aef6157ae0

+ 48 - 11
jobs/ods/trd/ods_trd_card_group_info_ful_d.sql

@@ -1,20 +1,32 @@
 -- 作者:tianyu.chu
 -- 日期:2026-07-17
 -- 工单:(无)
--- 目的:拼团全量最新态日常 merge —— full dt=${dt} ← full dt=${pdt} ⨝ inc dt=${dt},按 id 取最新一版
+-- 目的:拼团全量最新态日常 merge —— full dt=${dt} = 昨日全量中今天没变的 + 今日增量
 -- 状态:[草案]
 -- 备注:sched=T,${dt}=业务日 T-1、${pdt}=T-2(即前一天的 full)。
+--
+--       **反连接而非窗口函数**:今天的全量 = 「昨日全量 LEFT ANTI JOIN 今日增量 id」+「今日增量」。
+--       右侧只有当日变更的 id(千级),Spark 会 broadcast,全量表只扫不 shuffle、也不排序。
+--       若改用 ROW_NUMBER() OVER (PARTITION BY id) 套在 union 上,等于每天把全量按 id 全shuffle+排序一遍,
+--       增量才几千行却要搬动全表,merge 的意义就没了。
+--       正确性:ODS inc 单个 dt 分区内一个 id 只有一行(已按 (id, ods_dt) 去重),且 inc dt=${dt}
+--       意味着该记录最后变更就在当天、必然比昨日全量那版新 —— 故「变了的取 inc、没变的留昨日」
+--       与「取 max(update_time)」等价。
+--       注:LEFT ANTI JOIN 是 Spark SQL 语法(Hive 2.1 不支持),本 SQL 走 spark-sql-starter,无碍。
+--
 --       **静态分区写入**(PARTITION(dt='${dt}'))是硬要求:本 SQL 读 full dt=${pdt} 又写 full,
 --       若用动态分区(PARTITION(dt)),输出路径=表根目录、与读的分区路径重叠,Spark 会报
 --       "Cannot overwrite a path that is also being read from"。静态分区的输出是具体分区目录,不重叠。
 --       起点由 manual/backfill/20260717_..._ful_d_seed.sql 一次性 seed;断链需重新 seed。
 --       历史多版本在 inc 表(拉链底座),本表只保当前态。lock 是保留字,SELECT 里必须带反引号。
+--
 --       **分区保留 2 天**(只留 ${dt} 与 ${pdt}):文末 ALTER 精确 drop ${rdt}(=${dt}-2,即 T-3),
 --       每天正好掉一个。用精确分区值而非 dt < 'x'——Spark SQL 的 DROP PARTITION 只接受 col=value,
 --       比较式是 Hive 语法、Spark 解析不了。DS 需额外传参 rdt=$[yyyyMMdd-3];手动跑加 -p rdt=<T-3>。
 --       代价:漏跑一天则那天本该掉的分区不会被补 drop,需手工清;断链重新 seed 即可。
 
 INSERT OVERWRITE TABLE ods.ods_trd_card_group_info_ful_d PARTITION (dt='${dt}')
+-- 昨日全量里今天没变的
 SELECT
     id, merchant_id, appid, name, code,
     status, specs, type, random_type, total_price,
@@ -40,16 +52,41 @@ SELECT
     act_type, waring_type, compensation_status, point_type, first_act_config,
     gift_config, version, extra_prop, use_member_discount, merchant_open,
     is_deleted
-FROM (
-    SELECT *,
-        ROW_NUMBER() OVER (PARTITION BY id ORDER BY COALESCE(update_time, create_time) DESC) AS rn
-    FROM (
-        SELECT * FROM ods.ods_trd_card_group_info_ful_d WHERE dt = '${pdt}'
-        UNION ALL
-        SELECT * FROM ods.ods_trd_card_group_info_inc_d WHERE dt = '${dt}'
-    ) u
-) t
-WHERE t.rn = 1;
+FROM (SELECT * FROM ods.ods_trd_card_group_info_ful_d WHERE dt = '${pdt}') f
+LEFT ANTI JOIN (
+    SELECT id FROM ods.ods_trd_card_group_info_inc_d WHERE dt = '${dt}'
+) i ON f.id = i.id
+
+UNION ALL
+
+-- 今日变了的(含今日新建)
+SELECT
+    id, merchant_id, appid, name, code,
+    status, specs, type, random_type, total_price,
+    copies, unit_price, sold_copies, release_time, cycle,
+    show_applet, title, msg, remark, create_time,
+    update_by, update_time, order_quota_min, order_quota_max, user_quota_max,
+    start_time, marketing_info, reviewmsg, `lock`, commission_rate,
+    year, sport, manufacturer, sets, act,
+    config, info_config, total_num, banner_end_time, add_banner,
+    finished_time, display_name, group_sets_no, close_payment_time, confirm_send_time,
+    close_payment_status, open_card, close_payment_record, group_full_time, live_create_time,
+    live_start_time, live_end_time, report_start_time, report_end_time, report_review_num,
+    report_review_first_time, report_review_end_time, review_hold_time, review_approval_time, review_num,
+    config_json, free_flag, mer_name, change_type, act_price,
+    act_config_json, real_sold_num, weight, hot_type, team_first,
+    prop1, prop2, prop3, point_rate, point_max,
+    point_min, list_id, list_code, mix_copies, sub_type,
+    act_point_type, payment_method, payment_total_price, payment_commission, payment_finished_price,
+    payment_remain_price, payment_online_price, exclusive, has_bg, merchant_sort,
+    del_flg, del_time, review_account, act_id, sold_end_time,
+    panini_list_id, hot_type_config, goods_type, report_flag, use_coupon,
+    user_level, custom, gift_card_id, group_show_name, min_card_num,
+    act_type, waring_type, compensation_status, point_type, first_act_config,
+    gift_config, version, extra_prop, use_member_discount, merchant_open,
+    is_deleted
+FROM ods.ods_trd_card_group_info_inc_d
+WHERE dt = '${dt}';
 
 -- 保留 2 天:drop 掉 ${dt}-2(T-3)这一个分区
 ALTER TABLE ods.ods_trd_card_group_info_ful_d DROP IF EXISTS PARTITION (dt='${rdt}');

+ 42 - 11
jobs/ods/trd/ods_trd_card_group_order_info_ful_d.sql

@@ -1,20 +1,31 @@
 -- 作者:tianyu.chu
 -- 日期:2026-07-17
 -- 工单:(无)
--- 目的:订单全量最新态日常 merge —— full dt=${dt} ← full dt=${pdt} ⨝ inc dt=${dt},按 id 取最新一版
+-- 目的:订单全量最新态日常 merge —— full dt=${dt} = 昨日全量中今天没变的 + 今日增量
 -- 状态:[草案]
 -- 备注:sched=T,${dt}=业务日 T-1、${pdt}=T-2(即前一天的 full)。
+--
+--       **反连接而非窗口函数**:今天的全量 = 「昨日全量 LEFT ANTI JOIN 今日增量 id」+「今日增量」。
+--       右侧只有当日变更的 id(万级),Spark 会 broadcast,6200 万行的 full 只扫不 shuffle、也不排序。
+--       若改用 ROW_NUMBER() OVER (PARTITION BY id) 套在 union 上,等于每天把全量按 id 全shuffle+排序一遍,
+--       增量才 20 万行却要搬动 6200 万行,merge 的意义就没了。
+--       正确性:ODS inc 单个 dt 分区内一个 id 只有一行(已按 (id, ods_dt) 去重),且 inc dt=${dt}
+--       意味着该单最后变更就在当天、必然比昨日全量那版新 —— 故「变了的取 inc、没变的留昨日」
+--       与「取 max(update_time)」等价。
+--
 --       **静态分区写入**(PARTITION(dt='${dt}'))是硬要求:本 SQL 读 full dt=${pdt} 又写 full,
 --       若用动态分区(PARTITION(dt)),输出路径=表根目录、与读的分区路径重叠,Spark 会报
 --       "Cannot overwrite a path that is also being read from"。静态分区的输出是具体分区目录,不重叠。
 --       起点由 manual/backfill/20260717_..._ful_d_seed.sql 一次性 seed;断链需重新 seed。
 --       历史多版本在 inc 表(拉链底座),本表只保当前态。
+--
 --       **分区保留 2 天**(只留 ${dt} 与 ${pdt}):文末 ALTER 精确 drop ${rdt}(=${dt}-2,即 T-3),
 --       每天正好掉一个。用精确分区值而非 dt < 'x'——Spark SQL 的 DROP PARTITION 只接受 col=value,
 --       比较式是 Hive 语法、Spark 解析不了。DS 需额外传参 rdt=$[yyyyMMdd-3];手动跑加 -p rdt=<T-3>。
 --       代价:漏跑一天则那天本该掉的分区不会被补 drop,需手工清;断链重新 seed 即可。
 
 INSERT OVERWRITE TABLE ods.ods_trd_card_group_order_info_ful_d PARTITION (dt='${dt}')
+-- 昨日全量里今天没变的
 SELECT
     id, group_info_id, merchant_id, user_id, shipping_address_id,
     purchase_count, order_no, accounts_payable, actual_payment, payment_type,
@@ -35,16 +46,36 @@ SELECT
     pre_wait_shipped_num, refuse_time, refuse_notice, pickup_time, waring_type,
     waring_status, point_type, delivery_end_time, serve_status, self_pickup_time,
     act_discount, is_deleted
-FROM (
-    SELECT *,
-        ROW_NUMBER() OVER (PARTITION BY id ORDER BY COALESCE(update_time, create_time) DESC) AS rn
-    FROM (
-        SELECT * FROM ods.ods_trd_card_group_order_info_ful_d WHERE dt = '${pdt}'
-        UNION ALL
-        SELECT * FROM ods.ods_trd_card_group_order_info_inc_d WHERE dt = '${dt}'
-    ) u
-) t
-WHERE t.rn = 1;
+FROM (SELECT * FROM ods.ods_trd_card_group_order_info_ful_d WHERE dt = '${pdt}') f
+LEFT ANTI JOIN (
+    SELECT id FROM ods.ods_trd_card_group_order_info_inc_d WHERE dt = '${dt}'
+) i ON f.id = i.id
+
+UNION ALL
+
+-- 今日变了的(含今日新建)
+SELECT
+    id, group_info_id, merchant_id, user_id, shipping_address_id,
+    purchase_count, order_no, accounts_payable, actual_payment, payment_type,
+    payment_time, coupon, discount, status, remark,
+    create_time, create_by, update_time, update_by, payment_status,
+    payment_status_desc, payment_success_time, del_flg, curier_company, refund_fee,
+    refund_time, anonymous, pick_up_type, ship_time, refund_success_time,
+    refund_recv_accout, refund_account, refund_request_source, card_price, act_price,
+    goods_price_json, payment_sub_type, team_first, refuse_status, prop1,
+    prop2, prop3, point, order_type, trade_amount,
+    refund_type, refund_reason, evaluation, user_refund_time, refund_status,
+    merchant_refund_reason, point_deduct, shipping_cost, merchant_remark, pay_record,
+    order_sub_type, give_user_code, give_order_id, read_flag, give_num,
+    invoice_id, combination_no, open_self, refund_desc, goods_allocate,
+    close_payment_status, close_payment_time, finished_time, expire_time, settlement_amount,
+    platform_coupon, platform_discount, discount_amount, member_discount, shipping_free_id,
+    shipping_free_amount, discount_point, un_shipped_num, pre_un_shipped_num, wait_shipped_num,
+    pre_wait_shipped_num, refuse_time, refuse_notice, pickup_time, waring_type,
+    waring_status, point_type, delivery_end_time, serve_status, self_pickup_time,
+    act_discount, is_deleted
+FROM ods.ods_trd_card_group_order_info_inc_d
+WHERE dt = '${dt}';
 
 -- 保留 2 天:drop 掉 ${dt}-2(T-3)这一个分区
 ALTER TABLE ods.ods_trd_card_group_order_info_ful_d DROP IF EXISTS PARTITION (dt='${rdt}');