5.从ODS层抽取数据到DWD层(字段脱敏,字段名字统一)
5.1 需要注意事项
脱敏时可以使用md5,正则表达式(将数据中的数字变成星号)
表名后缀 .msl-d : 其中msl表示已经脱敏, d表示按天进行分区
源码名字命名规则:表名单词首字母打写,驼峰命名 ,例如:
dwd_fcj_nwrs_sellbargain_msl_d
- 直接整合hive就不需要配置schema浪费时间
但是需要开启enableHivesupport() 开启hive的元数据支持
用dwd用户执行任务会报一个错误,就是没有权限去从ods层拉取数据,需要申请权限
可以编写一个脚本脚本如下:
# 分区 ds=$1 # 执行任务 spark-submit \ --master yarn-client \ --class com.wt.dwd.DwdFcjNwrsSellbargainMskDay \ ../target/dwd-1.0-SNAPSHOT.jar \ $ds # 增加分区 hive -e "alter table dwd.dwd_fcj_nwrs_sellbargain_msl_d add IF NOT EXISTS partition (ds='$ds')"
因为是T+1模式,所以需要在代码中指定一个变量可以使用val ds: String = args.head向代码中指定一个变量。然后再.jar包后面指定 $ds
还可以在脚本中动态的指定分区
5.2 执行脚本结果如下:

5.3 目录结构如下:

5.3.1 dwd_fcj_nwrs_sellbargain_msl_d.sh 内容为:
- # 分区
- ds=$1
-
- # 执行任务
- spark-submit \
- --master yarn-client \
- --class com.wt.dwd.DwdFcjNwrsSellbargainMskDay \
- ../target/dwd-1.0-SNAPSHOT.jar \
- $ds
- # 增加分区
- hive -e "alter table dwd.dwd_fcj_nwrs_sellbargain_msl_d add IF NOT EXISTS partition (ds='$ds')"
5.3.2 dwd_fcj_nwrs_sellbargain_msl_d.sql 内容为:
- -- hive建表语句
- -- hive建表语句
- CREATE external TABLE IF NOT EXISTS dwd.dwd_fcj_nwrs_sellbargain_msl_d(
- id STRING comment '身份证号码',
- r_fwzl STRING comment '房产地址',
- htydjzmj STRING comment '合同中约定房子面积',
- tntjzmj STRING comment '房子内建筑面积',
- ftmj STRING comment '房子分摊建筑面积',
- time_tjba STRING comment '商品房备案时间',
- htzj STRING comment '合同总价'
- )PARTITIONED BY
- (
- ds STRING
- )
- ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
- STORED AS textfile
- location '/daas/motl/dwd/dwd_fcj_nwrs_sellbargain_msl_d/';
5.3.3 DwdFcjNwrsSellbargainMskDay.scala内容为:
- package com.wt.dwd
- import org.apache.spark.sql.{DataFrame, SaveMode, SparkSession}
-
- object DwdFcjNwrsSellbargainMskDay {
- def main(args: Array[String]): Unit = {
-
- /**
- * 1. 创建spark环境
- *
- */
- val spark: SparkSession = SparkSession
- .builder
- //.master("local")
- .enableHiveSupport() //开启hive元数据支持,开启之后在spark中可以直接读取hive中的表,但是开启之后就不能再本地云心的了
- .getOrCreate()
-
- import spark.implicits._
- import org.apache.spark.sql.functions._
-
- /**
- * 获取时间分区的字段
- *
- */
- val ds: String = args.head
- /**
- * 2. 读取获取购房合同中的表
- * 必须带上库名,否则读不到
- *
- * 不可能读取所有的数据,我们只需要读取每一天的数据
- *
- */
- val sellbargain: DataFrame = spark
- .table("ods.ods_t_fcj_nwrs_sellbargain")
- .where($"ds" === ds)
-
- //对原始的数据进行托名 对id进行脱敏,然后将r_fwzl中的数字变成 * 号(通过正则表达式替换)
- val resultDF: DataFrame = sellbargain.select(
- upper(md5($"id")) as "id",
- regexp_replace($"r_fwzl", "\\d", "*") as "r_fwzl",
- $"htydjzmj",
- $"tntjzmj",
- $"ftmj",
- $"time_tjba",
- $"htzj"
- )
-
- resultDF
- .write
- .format("csv")
- .mode(SaveMode.Overwrite)
- .option("sep","\t")
- .save(s"/daas/motl/dwd/dwd_fcj_nwrs_sellbargain_msl_d/ds=$ds")
-
- //提交到集群运行 spark-submit --master yarn-client --class com.wt.dwd.DwdFcjNwrsSellbargainMskDay dwd-1.0-SNAPSHOT.jar
- }
- }
