• 5.从ODS层抽取数据到DWD层(字段脱敏,字段名字统一)


    5.从ODS层抽取数据到DWD层(字段脱敏,字段名字统一)

    5.1 需要注意事项

    1. 脱敏时可以使用md5,正则表达式(将数据中的数字变成星号)

    2. 表名后缀 .msl-d : 其中msl表示已经脱敏, d表示按天进行分区

    3. 源码名字命名规则:表名单词首字母打写,驼峰命名 ,例如:

    dwd_fcj_nwrs_sellbargain_msl_d

    1. 直接整合hive就不需要配置schema浪费时间

    但是需要开启enableHivesupport() 开启hive的元数据支持

    1. 用dwd用户执行任务会报一个错误,就是没有权限去从ods层拉取数据,需要申请权限

    2. 可以编写一个脚本脚本如下:

    1. # 分区
    2. ds=$1
    3. # 执行任务
    4. spark-submit \
    5. --master yarn-client \
    6. --class com.wt.dwd.DwdFcjNwrsSellbargainMskDay \
    7. ../target/dwd-1.0-SNAPSHOT.jar \
    8. $ds
    9. # 增加分区
    10. hive -e "alter table dwd.dwd_fcj_nwrs_sellbargain_msl_d add IF NOT EXISTS partition (ds='$ds')"
    1. 因为是T+1模式,所以需要在代码中指定一个变量可以使用val ds: String = args.head向代码中指定一个变量。然后再.jar包后面指定 $ds

    2. 还可以在脚本中动态的指定分区

    5.2 执行脚本结果如下:

    5.3 目录结构如下:

    5.3.1 dwd_fcj_nwrs_sellbargain_msl_d.sh 内容为:

    1. # 分区
    2. ds=$1
    3. # 执行任务
    4. spark-submit \
    5. --master yarn-client \
    6. --class com.wt.dwd.DwdFcjNwrsSellbargainMskDay \
    7. ../target/dwd-1.0-SNAPSHOT.jar \
    8. $ds
    9. # 增加分区
    10. 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 内容为:

    1. -- hive建表语句
    2. -- hive建表语句
    3. CREATE external TABLE IF NOT EXISTS dwd.dwd_fcj_nwrs_sellbargain_msl_d(
    4. id STRING comment '身份证号码',
    5. r_fwzl STRING comment '房产地址',
    6. htydjzmj STRING comment '合同中约定房子面积',
    7. tntjzmj STRING comment '房子内建筑面积',
    8. ftmj STRING comment '房子分摊建筑面积',
    9. time_tjba STRING comment '商品房备案时间',
    10. htzj STRING comment '合同总价'
    11. )PARTITIONED BY
    12. (
    13. ds STRING
    14. )
    15. ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
    16. STORED AS textfile
    17. location '/daas/motl/dwd/dwd_fcj_nwrs_sellbargain_msl_d/';

    5.3.3 DwdFcjNwrsSellbargainMskDay.scala内容为:

    1. package com.wt.dwd
    2. import org.apache.spark.sql.{DataFrame, SaveMode, SparkSession}
    3. object DwdFcjNwrsSellbargainMskDay {
    4. def main(args: Array[String]): Unit = {
    5. /**
    6. * 1. 创建spark环境
    7. *
    8. */
    9. val spark: SparkSession = SparkSession
    10. .builder
    11. //.master("local")
    12. .enableHiveSupport() //开启hive元数据支持,开启之后在spark中可以直接读取hive中的表,但是开启之后就不能再本地云心的了
    13. .getOrCreate()
    14. import spark.implicits._
    15. import org.apache.spark.sql.functions._
    16. /**
    17. * 获取时间分区的字段
    18. *
    19. */
    20. val ds: String = args.head
    21. /**
    22. * 2. 读取获取购房合同中的表
    23. * 必须带上库名,否则读不到
    24. *
    25. * 不可能读取所有的数据,我们只需要读取每一天的数据
    26. *
    27. */
    28. val sellbargain: DataFrame = spark
    29. .table("ods.ods_t_fcj_nwrs_sellbargain")
    30. .where($"ds" === ds)
    31. //对原始的数据进行托名 对id进行脱敏,然后将r_fwzl中的数字变成 * 号(通过正则表达式替换)
    32. val resultDF: DataFrame = sellbargain.select(
    33. upper(md5($"id")) as "id",
    34. regexp_replace($"r_fwzl", "\\d", "*") as "r_fwzl",
    35. $"htydjzmj",
    36. $"tntjzmj",
    37. $"ftmj",
    38. $"time_tjba",
    39. $"htzj"
    40. )
    41. resultDF
    42. .write
    43. .format("csv")
    44. .mode(SaveMode.Overwrite)
    45. .option("sep","\t")
    46. .save(s"/daas/motl/dwd/dwd_fcj_nwrs_sellbargain_msl_d/ds=$ds")
    47. //提交到集群运行 spark-submit --master yarn-client --class com.wt.dwd.DwdFcjNwrsSellbargainMskDay dwd-1.0-SNAPSHOT.jar
    48. }
    49. }
  • 相关阅读:
    Python学习基础笔记六十七——格式化字符串
    公众号数据分析总结怎么做?教你玩转公众号后台数据
    【Matlab】数据统计分析
    spring-security-oauth2(授权模式入门简单使用)
    window 系统里 chrome 浏览器一些实用的调试技巧
    CP Autosar中的PNC说明
    计算机毕业设计(60)php小程序毕设作品之共享充电桩小程序系统
    JVM 优化技术
    linux应用之文件读取
    java面试(缓存Redis)
  • 原文地址:https://blog.csdn.net/weixin_48370579/article/details/126219653