Commit 143d61c1 by hejiangming

no message

parent a05de552
......@@ -17,6 +17,8 @@ import sys
sys.path.append(os.path.dirname(sys.path[0])) # 上级目录
from utils.templates import Templates
from utils.common_util import CommonUtil
from utils.hdfs_utils import HdfsUtils
# 分组排序的udf窗口函数
from pyspark.sql.window import Window
from pyspark.sql import functions as F
......@@ -62,6 +64,17 @@ class DwtAbaStAnalyticsReport(Templates):
self.partitions_by = ['site_name', 'date_type', 'date_info']
self.reset_partitions(65)
# 落表走 Templates.save_data 的 saveAsTable(mode='append'),不会覆盖旧数据
# 所以写之前先清掉本次分区目录下的文件,同一分区重跑不会追加成两份
partition_dict = {
"site_name": self.site_name,
"date_type": self.date_type,
"date_info": self.date_info,
}
hdfs_path = CommonUtil.build_hdfs_path(self.db_save, partition_dict=partition_dict)
print(f"清除hdfs目录中.....{hdfs_path}")
HdfsUtils.delete_file_in_folder(hdfs_path)
self.u_get_buy_box_num = self.spark.udf.register("u_get_buy_box_num", self.udf_get_buy_box_num, StringType())
self.u_get_buy_box = self.spark.udf.register("u_get_buy_box", self.udf_get_buy_box_type, StringType())
self.u_get_img_type = self.spark.udf.register("u_get_img_type", self.udf_get_img_type, StringType())
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment