Commit 6a615706 by hejiangming

店铺与选品数据对齐+新增字段

parent 102f3809
...@@ -7,14 +7,14 @@ ...@@ -7,14 +7,14 @@
本次只换数据源(ODS -> DIM),聚合与三个占比公式沿用原逻辑不动。 本次只换数据源(ODS -> DIM),聚合与三个占比公式沿用原逻辑不动。
@SourceTable : @SourceTable :
1.dim_fb_detail / ods_seller_account_feedback (店铺集合,按月门控) 1.dim_fb_detail / ods_seller_account_feedback (店铺集合,按月门控)
2.dim_fb_asin_info (店铺-ASIN 关系) 2.dwt_flow_asin (当月:account_id 店铺归属 + category_first_id 一级分类 + asin_is_new 新品 + 市场分母)
3.dwt_flow_asin (asin_is_new,按 asin 取值) 3.dim_bsr_category_tree (一级分类名称)
4.dim_cal_asin_history_detail (一级分类 + 市场总量分母)
5.dim_bsr_category_tree (一级分类名称)
@SinkTable : @SinkTable :
1.dwt_fb_category_report 1.dwt_fb_category_report
@CreateTime : 2023/07/18 17:33 @CreateTime : 2023/07/18 17:33
@UpdateTime : 2026/08/17 17:59 @UpdateTime : 2026/09/10 00:00
2026-09 按月严格同源改口径:店铺-asin 归属/一级分类/新品/市场分母全部改用 dwt_flow_asin 当月,
与主表 dwt_fb_base_report(已改 flow 当月)对齐;原 dim_fb_asin_info(全历史名单)+dim_cal_asin_history_detail(全历史)弃用。
""" """
import os import os
...@@ -48,9 +48,8 @@ class DwtFbCategoryReport(object): ...@@ -48,9 +48,8 @@ class DwtFbCategoryReport(object):
self.spark = SparkUtil.get_spark_session(app_name) self.spark = SparkUtil.get_spark_session(app_name)
# 初始化全局df # 初始化全局df
self.df_fb_asin_info = self.spark.sql(f"select 1+1;") # 主表店铺 x 店铺-ASIN 关系 self.df_fb_asin_info = self.spark.sql(f"select 1+1;") # 主表店铺 x 当月 flow 归属 asin
self.df_flow_asin = self.spark.sql(f"select 1+1;") # asin -> is_asin_new self.df_flow_all = self.spark.sql(f"select 1+1;") # 全站当月 flow asin(市场分母用)
self.df_asin_history = self.spark.sql(f"select 1+1;")
self.df_cate_name = self.spark.sql(f"select 1+1;") self.df_cate_name = self.spark.sql(f"select 1+1;")
self.df_fb_cate_asin_cal = self.spark.sql(f"select 1+1;") self.df_fb_cate_asin_cal = self.spark.sql(f"select 1+1;")
self.df_fb_asin_cal = self.spark.sql(f"select 1+1;") self.df_fb_asin_cal = self.spark.sql(f"select 1+1;")
...@@ -94,49 +93,36 @@ class DwtFbCategoryReport(object): ...@@ -94,49 +93,36 @@ class DwtFbCategoryReport(object):
df_fb_seller = self.spark.sql(sqlQuery=sql) df_fb_seller = self.spark.sql(sqlQuery=sql)
print(sql) print(sql)
# 店铺-ASIN 关系:dim_fb_asin_info 从 2026-06 起有分区,历史月读 2026-06 # 换口径(2026-09,按月严格同源):店铺-asin 归属、一级分类、新品标记、市场分母全部改用 dwt_flow_asin 当月
# 替代原来的 ods_seller_asin_account —— 那张表读取时带 created_at <= cal_date 过滤, # 起因:主表 fb_asin_total 等已改 flow 当月口径,本表原来用 dim_fb_asin_info(全历史名单)+dim_cal_asin_history_detail(全历史分类/市场)对不上
# 爬虫先删后增会刷新 created_at,跨月重跑时活跃店铺被误筛成 0 商品;dim 层已去掉该过滤 # 旧源(注释保留备查):
asin_info_date = self.date_info if self.date_info >= '2026-06' else '2026-06' # 店铺-asin: dim_fb_asin_info(全历史名单,2026-06 起分区,历史月读 2026-06)
print(f"获取 dim_fb_asin_info(店铺-ASIN 关系,date_info={asin_info_date})") # asin_info_date = self.date_info if self.date_info >= '2026-06' else '2026-06'
# df_fb_asin = spark.sql("select seller_id, asin from dim_fb_asin_info where ... date_info='{asin_info_date}'")
# self.df_fb_asin_info = df_fb_seller.join(df_fb_asin, on='seller_id', how='left').drop_duplicates(['seller_id','asin']).cache()
# 新品标记: self.df_flow_asin = spark.sql("select asin, asin_is_new as is_asin_new from dwt_flow_asin where ...").drop_duplicates(['asin'])
# 分类+市场分母: self.df_asin_history = spark.sql("select asin, category_first_id as bsr_cate_1_id from dim_cal_asin_history_detail where site_name=...").cache()
# flow 当月:account_id 归属店铺、category_first_id 当月一级分类、asin_is_new 新品标记,一次取全
print("获取 dwt_flow_asin(account_id/一级分类/新品,当月)")
sql = f""" sql = f"""
select seller_id, asin select account_id as seller_id, asin,
from dim_fb_asin_info asin_is_new as is_asin_new,
where site_name = '{self.site_name}' category_first_id as bsr_cate_1_id
and date_type = '{self.date_type}'
and date_info = '{asin_info_date}'
"""
df_fb_asin = self.spark.sql(sqlQuery=sql)
print(sql)
# 主表店铺 left join 店铺-ASIN 关系:店铺一个不丢,没有 asin 的店铺 asin 为 null
self.df_fb_asin_info = df_fb_seller.join(df_fb_asin, on='seller_id', how='left')
self.df_fb_asin_info = self.df_fb_asin_info.drop_duplicates(['seller_id', 'asin']).cache()
# asin_is_new 从 flow 按 asin 取值(不是按店铺,不影响店铺集合)
# 用 flow 已算好的值而不是 UDF 重算,口径与 dwt_fb_base_report.fb_new_asin_num 对齐
print("获取 dwt_flow_asin(asin_is_new)")
sql = f"""
select asin, asin_is_new as is_asin_new
from dwt_flow_asin from dwt_flow_asin
where site_name = '{self.site_name}' where site_name = '{self.site_name}'
and date_type = '{self.date_type}' and date_type = '{self.date_type}'
and date_info = '{self.date_info}' and date_info = '{self.date_info}'
""" """
self.df_flow_asin = self.spark.sql(sqlQuery=sql).drop_duplicates(['asin']) df_flow = self.spark.sql(sqlQuery=sql).drop_duplicates(['asin'])
print(sql) print(sql)
# 市场分母用:全站当月 flow asin(不筛 account_id)
self.df_flow_all = df_flow.cache()
# 店铺归属用:account_id 非空 = 当月归属到某店的 asin
df_flow_store = df_flow.filter(F.col('seller_id').isNotNull())
# dim_cal_asin_history_detail: 两个用途 # 主表店铺 left join 当月 flow 归属 asin:feedback 店一个不丢,当月没 flow asin 的店 asin 为 null(占比落 null,前端显示 -)
# 1. 给上面的 seller-asin 补一级分类(asin 全历史维表,覆盖率高于 flow 当月快照) self.df_fb_asin_info = df_fb_seller.join(df_flow_store, on='seller_id', how='left')
# 2. 按一级分类统计全量 asin 数,作为市场占比 fb_market_rate 的分母 self.df_fb_asin_info = self.df_fb_asin_info.drop_duplicates(['seller_id', 'asin']).cache()
print("获取 dim_cal_asin_history_detail (bsr_cate_1_id + bsr_asin_num)")
sql = f"""
select asin, category_first_id as bsr_cate_1_id
from dim_cal_asin_history_detail
where site_name = '{self.site_name}'
"""
self.df_asin_history = self.spark.sql(sqlQuery=sql).cache()
print(sql)
# 一级分类名称 # 一级分类名称
print("获取 dim_bsr_category_tree") print("获取 dim_bsr_category_tree")
...@@ -158,15 +144,14 @@ class DwtFbCategoryReport(object): ...@@ -158,15 +144,14 @@ class DwtFbCategoryReport(object):
self.sava_data() self.sava_data()
def handle_fb_agg(self): def handle_fb_agg(self):
# 旧代码:df_fb_asin_info 直接来自 flow,已自带 is_asin_new 和 bsr_cate_1_id,不用再 join # df_fb_asin_info 来自 flow 当月,已自带 is_asin_new 和 bsr_cate_1_id,不用再 join
# self.df_fb_cate_asin_cal = self.df_fb_asin_info # 旧代码(名单口径下 df_fb_asin_info 只有 seller_id+asin,需 join 补分类/新品,注释保留):
# self.df_fb_cate_asin_cal = self.df_fb_asin_info \
# 现在 df_fb_asin_info 只有 seller_id + asin,按 asin 补一级分类和新品标记 # .join(self.df_asin_history, on='asin', how='left') \
self.df_fb_cate_asin_cal = self.df_fb_asin_info \ # .join(self.df_flow_asin, on='asin', how='left')
.join(self.df_asin_history, on='asin', how='left') \ self.df_fb_cate_asin_cal = self.df_fb_asin_info
.join(self.df_flow_asin, on='asin', how='left')
# 当月无 BSR 分类的 asin 归到 '无'(按月口径:当月没上榜就不硬塞历史分类)
# 分类取不到的 asin 归到 '无',与原逻辑一致
self.df_fb_cate_asin_cal = self.df_fb_cate_asin_cal.na.fill({'bsr_cate_1_id': '无'}) self.df_fb_cate_asin_cal = self.df_fb_cate_asin_cal.na.fill({'bsr_cate_1_id': '无'})
# 按 seller_id + 一级分类聚合: 分类下 asin 数 + 新品数 # 按 seller_id + 一级分类聚合: 分类下 asin 数 + 新品数
...@@ -175,13 +160,14 @@ class DwtFbCategoryReport(object): ...@@ -175,13 +160,14 @@ class DwtFbCategoryReport(object):
F.sum("is_asin_new").alias("fb_cate_new_asin_num"), F.sum("is_asin_new").alias("fb_cate_new_asin_num"),
) )
# 店铺总 asin 数:数 dim_fb_asin_info 的关系条数,与 dwt_fb_base_report.fb_asin_total 同口径 # 店铺总 asin 数:数当月 flow 归属该店的 asin 条数,与 dwt_fb_base_report.fb_asin_total(现 flow 当月口径) 同口径
# count 会跳过 null,所以没有任何 asin 的店铺这里是 0 # count 会跳过 null,所以当月没 flow asin 的店这里是 0(占比落 null,前端显示 -)
self.df_fb_asin_cal = self.df_fb_asin_info.groupby(['seller_id']).agg( self.df_fb_asin_cal = self.df_fb_asin_info.groupby(['seller_id']).agg(
F.count("asin").alias("fb_asin_num")) F.count("asin").alias("fb_asin_num"))
# BSR 分类市场总量(仍用 dim_cal_asin_history_detail 全量 asin, 作为市场占比分母) # BSR 分类市场总量(市场占比分母):改用当月 flow 全站 asin 按一级分类数
self.df_bsr_asin_cal = self.df_asin_history.groupby(['bsr_cate_1_id']).agg( # 旧(全历史累计市场): self.df_bsr_asin_cal = self.df_asin_history.groupby(['bsr_cate_1_id']).agg(F.count("asin").alias("bsr_asin_num"))
self.df_bsr_asin_cal = self.df_flow_all.na.fill({'bsr_cate_1_id': '无'}).groupby(['bsr_cate_1_id']).agg(
F.count("asin").alias("bsr_asin_num")) F.count("asin").alias("bsr_asin_num"))
# 合并取到卖家分类的asin数量 和 bsr分类asin数量 # 合并取到卖家分类的asin数量 和 bsr分类asin数量
......
...@@ -142,7 +142,17 @@ if __name__ == '__main__': ...@@ -142,7 +142,17 @@ if __name__ == '__main__':
# Feedback 同比变化率(对比去年同月) # Feedback 同比变化率(对比去年同月)
"count_30_day_yoy_rate", "count_30_day_yoy_rate",
"count_1_year_yoy_rate", "count_1_year_yoy_rate",
"count_life_time_yoy_rate" "count_life_time_yoy_rate",
# 店铺品牌 Top3(name + 占比,占比空位/无数据 -1.0,name 空位 null)
"fb_brand_top1_name",
"fb_brand_top1_rate",
"fb_brand_top2_name",
"fb_brand_top2_rate",
"fb_brand_top3_name",
"fb_brand_top3_rate",
# 总销量同比/环比(完整真值表:null/0/±1000/真实比率)
"fb_shop_total_sales_yoy_rate",
"fb_shop_total_sales_mom_rate"
], ],
partition_dict={ partition_dict={
"site_name": site_name, "site_name": site_name,
......
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