1、全量覆盖
— 全量覆盖
/*
每日加载上游提供的全量数据。
该load方式适用于数据量不大的参数表、状态信息表、编码表等,业务接受数据完全覆盖,不追溯历史。
*/
INSERT OVERWRITE TABLE dwd.dwd_aaa_ywgc_ql PARTITION(pt_date='$load_time')
SELECT
业务字段
,CURRENT_TIMESTAMP() AS data_create_time
FROM
ods.ods_aaa_ywgc_ql
WHERE
pt_date = '$load_time';
2、增量追加(append)
— 增量追加(append)
/*
每日加载上游提供的增量数据。
该load方式适用于事件类表如流水表、埋点日志表等,无更新与删除操作
*/
INSERT OVERWRITE TABLE dwd.aaa_ywgc_zl PARTITION(pt_date='$load_date')
SELECT
业务字段
,CURRENT_TIMESTAMP() AS data_load_time
FROM ods.ods_aaa_ywgc_zl t1
WHERE pt_date='$load_date'
;
3、增量累全(增转全)
–增量累全(增转全)
/*
每日获取上游新增和更新数据,针对更新数据对历史数据进行更新,通过主键关联最新数据再回插至目标表。
*/
–获取当日增量数据
INSERT OVERWRITE TABLE dwd.dwd_aaa_ywgc_zl PARTITION (pt_date='$load_date')
SELECT
业务字段
,CURRENT_TIMESTAMP() AS data_create_time
FROM
ods.ods_aaa_ywgc_zl
WHERE
pt_date='$load_date'
;
— 将前一天全量数据加载到临时表
DROP TABLE IF EXISTS temp.dwd_aaa_ywgc_tmp1_ql;
CREATE TABLE IF NOT EXISTS temp.dwd_aaa_ywgc_tmp1_ql STORED AS ORC
AS
SELECT
业务字段
FROM dwd.dwd_aaa_ywgc_zzq –目标表
WHERE pt_date='$last_load_date' –last_load_date表示当前调度日期的前一天日期,通常来说真实的开发环境都会封装类似的日期参数
–pt_date=DATE_SUB('$load_date', 1); –没有可以这样处理
— 合并全量与增量数据(针对更新数据对历史数据进行更新,针对新增数据直接加载入历史表)
INSERT OVERWRITE TABLE dwd.dwd_aaa_ywgc_zzq PARTITION (pt_date='$load_date')
SELECT
t1.业务字段
,CURRENT_TIMESTAMP() AS data_create_time
FROM temp.dwd_aaa_ywgc_tmp1_ql t1
WHERE NOT EXISTS (
SELECT 1
FROM dwd.dwd_aaa_ywgc_zl t2
WHERE TRIM(COALESCE(t1.PKID,' '))=TRIM(COALESCE(t2.PKID,' '))
AND t2.pt_date='$load_date'
)
— 取出主键不等的数据(主键相等的数据:可能包含了重复数据与更新数据),相当于在全量数据中剔除了重复数据与要更新数据
UNION ALL
SELECT
业务字段
,CURRENT_TIMESTAMP() AS data_create_time
FROM dwd.dwd_aaa_ywgc_zl
WHERE pt_date='$load_date';
DROP TABLE IF EXISTS temp.dwd_aaa_ywgc_tmp1_ql;
4、增量拉链
–增量拉链
/*
拉链表记录数据历史变化,适用于用户信息变更追踪或者缓慢变化维场景(scd),增量拉链需每日加载上游提供的增量数据,全量拉链于此实现方式类似。
*/
— 今日增量数据
DROP TABLE IF EXISTS temp.aaa_yygc_h_zl_tmp1;
CREATE TABLE IF NOT EXISTS temp.aaa_yygc_h_zl_tmp1 STORED AS ORC
AS
SELECT
业务字段
,SHA256(CONCAT_WS('&',业务字段1,业务字段2,…)) AS all_col_flag –将全字段拼接,作为比较标识,外层可套个sha256h或MD5或其他自定义udf,能减少数据长度,也可加密数据
FROM ods.ods_aaa_ywgc_zl
WHERE pt_date='$load_date'
;
DROP TABLE IF EXISTS temp.aaa_yygc_h_ql_tmp2;
CREATE TABLE IF NOT EXISTS temp.aaa_yygc_h_ql_tmp2 STORED AS ORC
AS
— 将历史拉链表中开链数据导入到拉链临时表aaa_yygc_h_ql_tmp2中
INSERT OVERWRITE TABLE temp.aaa_yygc_h_ql_tmp2
SELECT
业务字段,
all_col_flag, –全字段拼接标识
st_dt, –拉链开链日期
'3000-12-31' AS end_dt –拉链闭链日期
— 将更新前数据还原
FROM dwd.aaa_yygc_h –目标拉链表
WHERE (pt_date='3000-12-31' AND st_dt < '$load_date') OR pt_date='$load_date'; — 支持重跑(将更新前数据还原)
–今日增量数据切片中的新增数据与更新数据
DROP TABLE IF EXISTS temp.aaa_yygc_h_ql_tmp3;
CREATE TABLE IF NOT EXISTS temp.aaa_yygc_h_ql_tmp3 STORED AS ORC
AS
SELECT
t1.业务字段,
t1.all_col_flag
FROM temp.aaa_yygc_h_zl_tmp1 t1 — 今日增量数据切片
LEFT JOIN temp.aaa_yygc_h_ql_tmp2 t2 — 历史开链数据
ON t1.all_col_flag=COALESCE(t2.all_col_flag, '')
WHERE t2.all_col_flag IS NULL; — 新增数据、更新数据(相当于在今日增量数据切片中剔除了无变动/重复数据)
–有效数据入end_dt='3000-12-31',无效数据入end_dt='$load_date'
DROP TABLE IF EXISTS temp.aaa_yygc_h_tmp4;
CREATE TABLE IF NOT EXISTS temp.aaa_yygc_h_tmp4 STORED AS ORC
AS
SELECT
t1.业务字段,
t1.all_col_flag,
t1.st_dt,
CASE
WHEN (t2.`pk` IS NOT NULL AND) THEN '$load_date'
ELSE t1.end_dt
END AS end_dt
— IS NOT NULL 更新前数据 (end_dt='$load_date')
— IS NULL 无变动数据 (end_dt='3000-12-31')
FROM temp.aaa_yygc_h_ql_tmp2 t1 — 历史开链数据
LEFT JOIN temp.aaa_yygc_h_ql_tmp3 t2 — 今日增量数据切片中的新增数据与更新数据
ON COALESCE(t1.`pk`, '') = COALESCE(t2.`pk`, '') — 能关联上的即为更新数据
UNION ALL
SELECT
业务字段,
all_col_flag,
'$load_date' AS st_dt,
'3000-12-31' AS end_dt
FROM temp.aaa_yygc_h_ql_tmp3; — 今日增量数据切片中的新增数据与更新后数据 (end_dt='3000-12-31')
— 数据分流、落表
set hive.exec.dynamic.partition=true;
set hive.exec.dynamic.partition.mode=nonstrict;
INSERT OVERWRITE TABLE dwd.aaa_yygc_h PARTITION (pt_date)
SELECT
业务字段,
all_col_flag,
CURRENT_TIMESTAMP() AS data_create_time
st_dt,
end_dt,
end_dt AS pt_date — pt_date=end_dt='3000-12-31' 或 '$load_date'
FROM temp.aaa_yygc_h_tmp4;
— 将数据进行分流,开链数据(无变动、新增、更新后)进入到'3000-12-31'分区,闭链数据(更新前数据)进入到'$load_date'分区
— '3000-12-31'分区下为最新的数据,'$load_date'分区下为老数据
DROP TABLE IF EXISTS temp.aaa_yygc_h_zl_tmp1;
DROP TABLE IF EXISTS temp.aaa_yygc_h_ql_tmp2;
DROP TABLE IF EXISTS temp.aaa_yygc_h_ql_tmp3;
DROP TABLE IF EXISTS temp.aaa_yygc_h_tmp4;


