Commit 0b5a5695 by chenyuanjie

店铺asin对应关系表-计算优化

parent 5706eb8c
...@@ -12,37 +12,31 @@ ...@@ -12,37 +12,31 @@
import os import os
import sys import sys
import re
from functools import reduce
sys.path.append(os.path.dirname(sys.path[0])) # 上级目录 sys.path.append(os.path.dirname(sys.path[0])) # 上级目录
from utils.templates import Templates
# 分组排序的udf窗口函数
from pyspark.sql.window import Window from pyspark.sql.window import Window
from pyspark.sql import functions as F from pyspark.sql import functions as F
from utils.spark_util import SparkUtil from utils.spark_util import SparkUtil
from pyspark.sql.types import StringType, IntegerType, DoubleType
from utils.hdfs_utils import HdfsUtils from utils.hdfs_utils import HdfsUtils
from utils.common_util import CommonUtil from utils.common_util import CommonUtil
class DwtFdAsinInfo(object): class DwtFdAsinInfo(object):
def __init__(self, site_name='us'): def __init__(self, site_name='us', date_type='month', date_info='2026-06'):
super().__init__() super().__init__()
self.hive_tb = "dim_fd_asin_info" self.hive_tb = "dim_fd_asin_info"
self.site_name = site_name self.site_name = site_name
self.date_type = date_type
self.date_info = date_info
self.partition_dict = { self.partition_dict = {
"site_name": site_name, "site_name": site_name,
} }
# 落表路径校验 # 落表路径校验
self.hdfs_path = CommonUtil.build_hdfs_path(self.hive_tb, partition_dict=self.partition_dict) self.hdfs_path = CommonUtil.build_hdfs_path(self.hive_tb, partition_dict=self.partition_dict)
app_name = f"{self.hive_tb}:{self.site_name}" app_name = f"{self.hive_tb}:{self.site_name}_{self.date_type}_{self.date_info}"
self.spark = SparkUtil.get_spark_session(app_name) self.spark = SparkUtil.get_spark_session(app_name)
self.partitions_num = CommonUtil.reset_partitions(self.site_name, 80) self.partitions_num = CommonUtil.reset_partitions(self.site_name, 80)
# 初始化全局变量df--ods获取数据的原始df # 初始化全局变量df--ods获取数据的原始df
self.df_seller_account_syn = self.spark.sql("select 1+1;") self.df_seller_account_syn = self.spark.sql("select 1+1;")
self.df_seller_account_feedback = self.spark.sql("select 1+1;") self.df_seller_account_feedback = self.spark.sql("select 1+1;")
...@@ -50,59 +44,48 @@ class DwtFdAsinInfo(object): ...@@ -50,59 +44,48 @@ class DwtFdAsinInfo(object):
# 初始化全局变量df--dwd层转换输出的df # 初始化全局变量df--dwd层转换输出的df
self.df_save = self.spark.sql(f"select 1+1;") self.df_save = self.spark.sql(f"select 1+1;")
# 注册自定义函数 (UDF)
self.u_handle_seller_unique = F.udf(self.udf_handle_seller_unique, StringType())
@staticmethod
def udf_handle_seller_unique(fd_url):
return fd_url.split("seller=")[-1].split('&')[0]
# 1.获取原始数据
def read_data(self): def read_data(self):
# 获取ods_seller相关原始 # 获取爬虫店铺记录
print("获取 ods_seller_account_syn") print("获取 ods_seller_account_syn")
sql = f"select id as fd_account_id,seller_id as unique_id, account_name as fd_account_name," \ sql = f"""
f"lower(account_name) as fd_account_name_lower, url as fd_url, created_at, updated_at from " \ select id as fd_account_id, seller_id as unique_id, account_name as fd_account_name,
f" ods_seller_account_syn where site_name='{self.site_name}' " \ lower(account_name) as fd_account_name_lower, url as fd_url
# f"group by id,account_name" from ods_seller_account_syn where site_name = '{self.site_name}'
"""
# seller_id 本身在 ods_seller_account_syn 里是唯一的,不需要再去重
self.df_seller_account_syn = self.spark.sql(sqlQuery=sql) self.df_seller_account_syn = self.spark.sql(sqlQuery=sql)
print(sql) print(sql)
# 避免重复数据
self.df_seller_account_syn = self.df_seller_account_syn.orderBy(self.df_seller_account_syn.fd_account_id.desc())
self.df_seller_account_syn = self.df_seller_account_syn.drop_duplicates(['unique_id'])
# print("self.df_seller_account_syn:", self.df_seller_account_syn.show(10, truncate=False))
# 获取店铺详情表:只读传参指定的这一个分区(date_type+date_info 即调度传入的最新分区)
print("获取 ods_seller_account_feedback") print("获取 ods_seller_account_feedback")
# 抓取会有重复,因此需要全局去重,count取最大即是最近的一个数据 sql = f"""
sql = f"select unique_id, " \ select seller_id as unique_id, country_name as fd_country_name, created_at
f" country_name as fd_country_name " \ from ods_seller_account_feedback
f"from (select seller_id as unique_id,country_name, " \ where site_name = '{self.site_name}' and date_type = '{self.date_type}' and date_info = '{self.date_info}'
f" row_number() over (partition by seller_id order by created_at desc) sort_flag " \ """
f" from ods_seller_account_feedback " \
f" where site_name = '{self.site_name}'" \
f" ) t1 " \
f"where sort_flag = 1; "
self.df_seller_account_feedback = self.spark.sql(sqlQuery=sql) self.df_seller_account_feedback = self.spark.sql(sqlQuery=sql)
print(sql) print(sql)
# print("self.df_seller_account_feedback", self.df_seller_account_feedback.show(10, truncate=False)) window = Window.partitionBy('unique_id').orderBy(F.col('created_at').desc())
self.df_seller_account_feedback = self.df_seller_account_feedback.withColumn(
# 获取我们内部的店铺与asin的数据库(从搜索词抓下来,店铺与asin的关系表) 'rank', F.row_number().over(window)
).filter('rank = 1').withColumn(
'fb_crawl_date', F.date_format(F.col('created_at'), 'yyyy-MM-dd HH:mm:ss')
).drop('rank', 'created_at')
# 获取店铺与asin的对应关系库(店铺与asin所有历史对应关系表)
print("获取 ods_seller_asin_account") print("获取 ods_seller_asin_account")
sql = f""" sql = f"""
select seller_id as unique_id,asin from ods_seller_asin_account select seller_id as unique_id, asin from ods_seller_asin_account where site_name='{self.site_name}'
where site_name='{self.site_name}' """
"""
self.df_fd_asin = self.spark.sql(sqlQuery=sql) self.df_fd_asin = self.spark.sql(sqlQuery=sql)
self.df_fd_asin = self.df_fd_asin.drop_duplicates(['unique_id', 'asin']) self.df_fd_asin = self.df_fd_asin.drop_duplicates(['unique_id', 'asin'])
print(sql) print(sql)
def save_data(self):
def sava_data(self): df_save = self.df_seller_account_syn.join(
# 店铺+asin 通过fd_account_name 拼接获取到fd_unique self.df_seller_account_feedback, on='unique_id', how='left'
df_save = self.df_seller_account_syn.join(self.df_seller_account_feedback, on='unique_id', how='left') ).join(
df_save = df_save.join(
self.df_fd_asin, on='unique_id', how='left' self.df_fd_asin, on='unique_id', how='left'
) )
...@@ -114,8 +97,9 @@ class DwtFdAsinInfo(object): ...@@ -114,8 +97,9 @@ class DwtFdAsinInfo(object):
F.col('fd_country_name'), F.col('fd_country_name'),
F.col('fd_url'), F.col('fd_url'),
F.col('asin'), F.col('asin'),
F.date_format(F.current_timestamp(), 'yyyy-MM-dd HH:mm:SS').alias('created_at'), F.date_format(F.current_timestamp(), 'yyyy-MM-dd HH:mm:ss').alias('created_at'),
F.date_format(F.current_timestamp(), 'yyyy-MM-dd HH:mm:SS').alias('updated_at'), F.date_format(F.current_timestamp(), 'yyyy-MM-dd HH:mm:ss').alias('updated_at'),
F.col('fb_crawl_date'),
F.lit(self.site_name).alias('site_name'), F.lit(self.site_name).alias('site_name'),
) )
...@@ -129,11 +113,11 @@ class DwtFdAsinInfo(object): ...@@ -129,11 +113,11 @@ class DwtFdAsinInfo(object):
def run(self): def run(self):
self.read_data() self.read_data()
self.sava_data() self.save_data()
if __name__ == '__main__': if __name__ == '__main__':
site_name = sys.argv[1] # 参数1:站点 site_name = sys.argv[1] # 参数1:站点
# date_type = sys.argv[2] # 参数2:类型:week/4_week/month/quarter date_type = sys.argv[2] # 参数2:类型:week/4_week/month/quarter
# date_info = sys.argv[3] # 参数3:年-周/年-月/年-季, 比如: 2022-1 date_info = sys.argv[3] # 参数3:年-周/年-月/年-季, 比如: 2022-1
handle_obj = DwtFdAsinInfo(site_name=site_name) handle_obj = DwtFdAsinInfo(site_name=site_name, date_type=date_type, date_info=date_info)
handle_obj.run() handle_obj.run()
\ No newline at end of file
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