Commit 7ef1fd6b by hejiangming

统一口径从dim_st_detail 拿数据 对于爬虫抓搜索词时 该搜索词没结果的情况用dim的数据对排名相关字段兜底 保证峰值月 常年可卖 趋势图这些字段的数据正常计算

parent 77d145f8
...@@ -123,6 +123,69 @@ class DwtAbaLast365(object): ...@@ -123,6 +123,69 @@ class DwtAbaLast365(object):
""" """
self.df_base = self.spark.sql(sql1).repartition(80, "search_term", "date_info").cache() self.df_base = self.spark.sql(sql1).repartition(80, "search_term", "date_info").cache()
# ============================================================
# 【补漏抓月骨架行】
# 为什么:爬虫某月漏抓 → 该词该月不在 dwt_aba_st_analytics → df_base 缺该(词,月) →
# 常年可卖(all_year_text_flag)误判否、peak_month/first_appear_month/total_appear_month/
# search_volume/st_num 等"按近12月"算的字段全被漏抓月带偏。
# 做法:从 dim_st_detail(ABA 清洗表,不经爬虫过滤,漏抓月仍有该词行)取近12月,
# 找出"已在 df_base 的词、但 df_base 缺的月",造骨架行(ABA 字段填真值、爬虫字段 null)
# union 进 df_base → 后面 handle_agg 一行不改,就按 ABA 的 presence + 真实 rank 算对。
# 骨架行打 is_skeleton=1 标记:只参与 handle_agg 的跨月聚合,不参与 handle_month_lastest
# 的单月"最新月"字段组(原因见 union 处注释)。
# 范围:inner join "已在 df_base 的词",只补它们的漏抓月;不引入从没被爬过的整词(那是方案①,不做)。
# 例:claw clips 2026-03 漏抓 → dim_st_detail 有 rank 850 → 补一行(rank/搜索量/st_num 有值,
# 销量/价格等爬虫字段 null) → 常年可卖数到12月判"是"、peak_month 纳入 850、2026-03 搜索量有值;
# 销量/毛利等那月仍空(爬虫无源,handle_agg 的 sum/max 自动忽略 null,不污染)。
# ============================================================
# 从 dim_st_detail 取近12月 ABA 字段(别名对齐 df_base 列名)
sql_gap = f"""
select
search_term,
st_bsr_cate_1_id_new as category_id,
st_bsr_cate_current_id_new as category_current_id,
st_search_num as search_volume,
st_appear_history_counts as st_num,
st_rank as rank,
st_is_search_text as is_search_text,
st_is_new_market_segment as is_new_market_segment,
st_is_ascending_text as is_ascending_text,
date_info
from dim_st_detail
where site_name = '{self.site_name}'
and date_type = '{self.date_type_original}'
and date_info in ({CommonUtil.list_to_insql(self.last_12_month)})
"""
df_dim = self.spark.sql(sql_gap)
# 已在 df_base 的词 + 其 id(骨架行沿用同一 id,保证聚合按 id 归到同一条)
df_exist_terms = self.df_base.select('search_term', 'id').dropDuplicates(['search_term'])
# df_base 已有的(词,月),用于排除"非漏抓月"
df_present = self.df_base.select('search_term', 'date_info').dropDuplicates()
# 漏抓月 = dim 有 且 词已在 df_base(inner) 且 该(词,月)不在 df_base(left_anti)
df_gap = df_dim.join(df_exist_terms, on='search_term', how='inner') \
.join(df_present, on=['search_term', 'date_info'], how='left_anti')
# 按 df_base 的列顺序与类型构造骨架:dim 提供的列取真值、其余列取 null,
# 每列都 cast 成 df_base 对应类型 → schema 完全一致,普通 unionByName 即可(不依赖 allowMissingColumns)
base_dtypes = dict(self.df_base.dtypes)
provided = set(df_gap.columns)
skeleton = df_gap.select([
(F.col(c) if c in provided else F.lit(None)).cast(base_dtypes[c]).alias(c)
for c in self.df_base.columns
])
# 打标 is_skeleton:骨架行只该参与跨月聚合(handle_agg 的峰值月/常年可卖/月列等),
# 不该参与单月口径的"最新月"字段组(handle_month_lastest)——否则当月正好是漏抓月时,
# 骨架行会被 filter(date_info=当月) 选成"最新月",落出 rank_lastest 有值、
# 价格/AO 等爬虫字段全 null 的半空行,违反"爬虫没抓到就不显示"的业务口径。
# 该列在 handle_agg 的 groupBy 聚合后自然消失,handle_save 又是显式 select,不会落表。
# union 补回 df_base。旧 df_base(sql1 的 cache) 在上面 union/join 里已被引用,
# 这里物化新 df_base 后显式 unpersist 旧的,避免新旧两份 df_base 同时占内存。
df_base_pre = self.df_base
self.df_base = df_base_pre.withColumn('is_skeleton', F.lit(0)) \
.unionByName(skeleton.withColumn('is_skeleton', F.lit(1))) \
.repartition(80, "search_term", "date_info").cache()
self.df_base.count() # 触发物化(读旧 cache 完成 union),之后旧 cache 可安全释放
df_base_pre.unpersist()
# 搜索词同比、环比 # 搜索词同比、环比
sql2 = f""" sql2 = f"""
select select
...@@ -191,17 +254,27 @@ class DwtAbaLast365(object): ...@@ -191,17 +254,27 @@ class DwtAbaLast365(object):
""" """
self.df_st_sv_rank = self.spark.sql(sql).na.fill({"last_rank": 0}).cache() self.df_st_sv_rank = self.spark.sql(sql).na.fill({"last_rank": 0}).cache()
# 历史新增词识别:读取月表所有历史数据,判断是否为历史新增 # ===== 原逻辑(保留备查):读 dwt_aba_st_analytics 求 min,因 dwt 月表老分区已删、最早只到~2023,首现月偏晚且会误判,已换源 =====
# # 历史新增词识别:读取月表所有历史数据,判断是否为历史新增
# sql = f"""
# select
# search_term,
# 0 as is_history_first_text,
# min(date_info) as history_first_appear_month
# from dwt_aba_st_analytics
# where site_name = '{self.site_name}'
# and date_type = '{self.date_type_original}'
# and date_info <= '{self.last_year_month}'
# group by search_term;
# """
# 历史新增词识别:改读全历史累加表 dim_st_detail_history 的 date_info_first(全站首次出现月,2020-05起),
# 与主表 is_first_ever_text 同源同口径;dim_st_detail_history 仅按 site_name 分区
sql = f""" sql = f"""
select select
search_term, search_term,
0 as is_history_first_text, date_info_first as history_first_appear_month
min(date_info) as history_first_appear_month from dim_st_detail_history
from dwt_aba_st_analytics
where site_name = '{self.site_name}' where site_name = '{self.site_name}'
and date_type = '{self.date_type_original}'
and date_info <= '{self.last_year_month}'
group by search_term;
""" """
self.df_history = self.spark.sql(sql).repartition(80, 'search_term').cache() self.df_history = self.spark.sql(sql).repartition(80, 'search_term').cache()
...@@ -261,7 +334,11 @@ class DwtAbaLast365(object): ...@@ -261,7 +334,11 @@ class DwtAbaLast365(object):
def handle_month_lastest(self): def handle_month_lastest(self):
# 保留最新的月数据 # 保留最新的月数据
# 注:st_attribute_label 已不在此处取(改到 handle_agg 按 12 个月合并去重,避免最新月没出现就丢标签) # 注:st_attribute_label 已不在此处取(改到 handle_agg 按 12 个月合并去重,避免最新月没出现就丢标签)
self.df_base_lastest = self.df_base.filter(f"date_info = '{self.date_info}'").select( # ===== 原逻辑(保留备查):self.df_base.filter(f"date_info = '{self.date_info}'") =====
# 【改】加 is_skeleton = 0:当月是漏抓月时,df_base 里该词的当月行是骨架行(ABA 有值、爬虫字段全 null),
# 不能被选成"最新月"——过滤后当月漏抓的词在 handle_save 的 left join 里 join 不上,
# 这组最新月字段整组 null,与改动前行为一致(爬虫没抓到就不显示)
self.df_base_lastest = self.df_base.filter(f"date_info = '{self.date_info}' and is_skeleton = 0").select(
"search_term", "rank", "color_proportion", "multi_color_proportion", "multi_size_proportion", "st_ao_avg", "search_term", "rank", "color_proportion", "multi_color_proportion", "multi_size_proportion", "st_ao_avg",
"st_ao_val_rate", "supply_demand", "total_asin_num", "new_asin_proportion", "bsr_orders", "st_ao_val_rate", "supply_demand", "total_asin_num", "new_asin_proportion", "bsr_orders",
"new_bsr_orders_proportion", "price_avg", "weight_avg", "volume_avg", "page3_brand_num", "brand_monopoly", "new_bsr_orders_proportion", "price_avg", "weight_avg", "volume_avg", "page3_brand_num", "brand_monopoly",
...@@ -414,14 +491,29 @@ class DwtAbaLast365(object): ...@@ -414,14 +491,29 @@ class DwtAbaLast365(object):
).fillna({'is_first_text': 1}) ).fillna({'is_first_text': 1})
# 历史新增词判断 # 历史新增词判断
# ===== 原逻辑(保留备查):靠"是否 join 上 df_history"打 0/1;history 源已换成全历史累加表 dim_st_detail_history =====
# self.df_base = self.df_base.join(
# self.df_history, on='search_term', how='left'
# ).fillna({
# 'is_history_first_text': 1
# }).withColumn(
# 'history_first_appear_month',
# F.when(F.col('history_first_appear_month').isNotNull(), F.col('history_first_appear_month'))
# .otherwise(F.col('first_appear_month'))
# )
# 【改】用全历史首现月 date_info_first 直接比较,更准
self.df_base = self.df_base.join( self.df_base = self.df_base.join(
self.df_history, on='search_term', how='left' self.df_history, on='search_term', how='left'
).fillna({ ).withColumn(
'is_history_first_text': 1 # 兜底:dwt 词⊆dim_st_detail_history(inner join 保证,已验证 missing=0),理论上必有值;
}).withColumn( # 万一拿不到,退回近12月首现 first_appear_month(与原逻辑同款兜底),保证不为 NULL
'history_first_appear_month', 'history_first_appear_month',
F.when(F.col('history_first_appear_month').isNotNull(), F.col('history_first_appear_month')) F.coalesce(F.col('history_first_appear_month'), F.col('first_appear_month'))
.otherwise(F.col('first_appear_month')) ).withColumn(
# is_history_first_text:首现月晚于12个月前(last_year_month)=近1年内才首现=1;否则=0(12+月前已存在)
# 口径与原逻辑一致(原按 date_info<=last_year_month 判0),只是"历史"从 dwt 有限窗口换成真全历史
'is_history_first_text',
F.when(F.col('history_first_appear_month') > F.lit(self.last_year_month), F.lit(1)).otherwise(F.lit(0))
) )
# 判断市场周期类型,优先保留最近月的数据,若为null则往前推 # 判断市场周期类型,优先保留最近月的数据,若为null则往前推
......
...@@ -48,6 +48,9 @@ class DwtAbaLastChangeRate(object): ...@@ -48,6 +48,9 @@ class DwtAbaLastChangeRate(object):
self.df_save = self.spark.sql(f"select 1+1;") self.df_save = self.spark.sql(f"select 1+1;")
# 需求1:近6月排名变化率/变化量(仅 month 类型生效) # 需求1:近6月排名变化率/变化量(仅 month 类型生效)
self.df_hist_rank = self.spark.sql(f"select 1+1;") # M-1~M-6 历史 rank self.df_hist_rank = self.spark.sql(f"select 1+1;") # M-1~M-6 历史 rank
# 爬虫漏抓补 rank:dim_st_detail 上月/去年同月 rank,给环比/同比 last_rank/last_year_rank 兜底(仅 month)
self.df_last_month_rank_dim = self.spark.sql(f"select 1+1;") # 上月 rank(dim)
self.df_last_year_rank_dim = self.spark.sql(f"select 1+1;") # 去年同月 rank(dim)
def handle_date_offset(self, handle_type: int): def handle_date_offset(self, handle_type: int):
# handle_type = 0 代表计算环比日期,等于 1 代表计算同比日期 # handle_type = 0 代表计算环比日期,等于 1 代表计算同比日期
...@@ -168,28 +171,79 @@ class DwtAbaLastChangeRate(object): ...@@ -168,28 +171,79 @@ class DwtAbaLastChangeRate(object):
month_list = [CommonUtil.get_month_offset(self.date_info, -i) for i in range(1, 7)] month_list = [CommonUtil.get_month_offset(self.date_info, -i) for i in range(1, 7)]
print(f"近6月历史月份列表: {month_list}") print(f"近6月历史月份列表: {month_list}")
# 一次读 dwt_aba_st_analytics 的 M-1~M-6 历史分区,仅取 search_term + rank # 一次读 M-1~M-6 历史分区的 rank,用于算 6 个变化量字段 + rank_rate_last_1_month
# 用于算 6 个变化量字段 + rank_rate_last_1_month # ===== 原逻辑(保留备查):读 dwt_aba_st_analytics 近6月 rank =====
# 问题:dwt = 爬虫∩ABA,某月漏抓 → 该月 rank 缺 → 近6月变化被误判成假下榜(+1e7)/假新进榜(-1e7)。
# sql_hist_rank = f"""
# select
# search_term,
# cast(rank as int) as rank,
# date_info
# from dwt_aba_st_analytics
# where site_name = '{self.site_name}'
# and date_type = '{self.date_type}'
# and date_info in ({CommonUtil.list_to_insql(month_list)})
# and rank > 0
# """
# 【改-一档】换源 dim_st_detail(ABA 清洗表,漏抓月仍有该词 st_rank),近6月排名变化不再出假尖刺
sql_hist_rank = f""" sql_hist_rank = f"""
select select
search_term, search_term,
cast(rank as int) as rank, cast(st_rank as int) as rank,
date_info date_info
from dwt_aba_st_analytics from dim_st_detail
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 in ({CommonUtil.list_to_insql(month_list)}) and date_info in ({CommonUtil.list_to_insql(month_list)})
and rank > 0 and st_rank > 0
""" """
self.df_hist_rank = self.spark.sql(sql_hist_rank).repartition(40, 'search_term').cache() self.df_hist_rank = self.spark.sql(sql_hist_rank).repartition(40, 'search_term').cache()
print("self.df_hist_rank:") print("self.df_hist_rank:")
self.df_hist_rank.show(10, truncate=True) self.df_hist_rank.show(10, truncate=True)
# 【二档】环比/同比 rank 补漏抓月:读 dim_st_detail 上月、去年同月的 st_rank,
# 在 handle_base / handle_year_ratio 里给 last_rank / last_year_rank 做 coalesce 兜底——
# 对比月被爬虫漏抓(dwt 无该词)时用 dim 真 rank,避免被算成假上升;
# 两端都无(真·新进榜/续断)时仍为 null → 沿用原 na.fill(-1000) 语义。
sql_last_month_rank = f"""
select search_term, cast(st_rank as int) as last_rank_dim
from dim_st_detail
where site_name = '{self.site_name}'
and date_type = '{self.date_type}'
and date_info = '{self.last_date_info}'
and st_rank > 0
"""
self.df_last_month_rank_dim = self.spark.sql(sql_last_month_rank).repartition(40, 'search_term').cache()
# 顺带多取 st_search_num,给同比搜索量变化率 last_year_search_volume 也做漏抓兜底(见 handle_year_ratio)
sql_last_year_rank = f"""
select search_term,
cast(st_rank as int) as last_year_rank_dim,
st_search_num as last_year_search_volume_dim
from dim_st_detail
where site_name = '{self.site_name}'
and date_type = '{self.date_type}'
and date_info = '{self.last_year_date_info}'
and st_rank > 0
"""
self.df_last_year_rank_dim = self.spark.sql(sql_last_year_rank).repartition(40, 'search_term').cache()
def handle_base(self): def handle_base(self):
self.df_st_base_data = self.df_aba_analytics.join( self.df_st_base_data = self.df_aba_analytics.join(
self.df_aba_analytics_old, on='id', how='left' self.df_aba_analytics_old, on='id', how='left'
) )
# 【二档-环比】last_rank 补漏抓月(仅 month):上月被爬虫漏抓时 dwt 的 last_rank 为 null,
# 用 dim 上月 st_rank(df_last_month_rank_dim,按 search_term)coalesce 兜底 → "漏抓续榜"不再算成假上升;
# 两端都无(真新进榜)时 last_rank 仍 null → 后续 na.fill(-1000) 语义不变。
# (df_st_base_data 已带 search_term,来自当月 df_aba_analytics;按 search_term left join,1:1 不放大)
if self.date_type == DateTypes.month.name:
self.df_st_base_data = self.df_st_base_data.join(
self.df_last_month_rank_dim, on='search_term', how='left'
).withColumn(
'last_rank', F.coalesce(F.col('last_rank'), F.col('last_rank_dim'))
).drop('last_rank_dim')
self.df_last_month_rank_dim.unpersist()
self.df_st_base_data = self.df_st_base_data.withColumn( self.df_st_base_data = self.df_st_base_data.withColumn(
'rank_rate_of_change', 'rank_rate_of_change',
F.round((F.col('rank') - F.col('last_rank')) / F.col('last_rank'), 4) F.round((F.col('rank') - F.col('last_rank')) / F.col('last_rank'), 4)
...@@ -218,6 +272,22 @@ class DwtAbaLastChangeRate(object): ...@@ -218,6 +272,22 @@ class DwtAbaLastChangeRate(object):
df_year_ratio = self.df_st_base_data.join( df_year_ratio = self.df_st_base_data.join(
self.df_st_last_year_data, on='search_term', how='left' self.df_st_last_year_data, on='search_term', how='left'
) )
# 【二档-同比】last_year_rank 补漏抓月(仅 month):去年同月被漏抓时 dwt 无该词 → last_year_rank null,
# 用 dim 去年同月 st_rank(df_last_year_rank_dim,按 search_term)coalesce 兜底 → 同比 rank 不再假上升;
# 真·去年没有(dim 也无)时仍 null → 沿用原 na.fill(-1000) 语义。
if self.date_type == DateTypes.month.name:
df_year_ratio = df_year_ratio.join(
self.df_last_year_rank_dim, on='search_term', how='left'
).withColumn(
'last_year_rank', F.coalesce(F.col('last_year_rank'), F.col('last_year_rank_dim'))
).withColumn(
# 同比搜索量变化率同样补漏抓月:去年同月漏抓时 dwt 的 last_year_search_volume 为 null,
# 用 dim 的 st_search_num(last_year_search_volume_dim)coalesce 兜底,避免 search_volume_change_rate 假上升(na.fill 1000);
# 真·去年没有(dim 也无)时仍 null → 沿用原 na.fill 语义。
'last_year_search_volume',
F.coalesce(F.col('last_year_search_volume'), F.col('last_year_search_volume_dim'))
).drop('last_year_rank_dim', 'last_year_search_volume_dim')
self.df_last_year_rank_dim.unpersist()
df_year_ratio = df_year_ratio.withColumn( df_year_ratio = df_year_ratio.withColumn(
"rank_change_rate", "rank_change_rate",
F.round(F.expr("(rank - last_year_rank) / last_year_rank"), 4) F.round(F.expr("(rank - last_year_rank) / last_year_rank"), 4)
......
...@@ -535,12 +535,28 @@ class DwtAbaStAnalytics(Templates): ...@@ -535,12 +535,28 @@ class DwtAbaStAnalytics(Templates):
# 读自身历史 11 个月分区的 rank,用于计算峰值月 peak_month / 常年可卖 all_year_text_flag # 读自身历史 11 个月分区的 rank,用于计算峰值月 peak_month / 常年可卖 all_year_text_flag
# ============================================================ # ============================================================
last_11_month = [CommonUtil.get_month_offset(self.date_info, -i) for i in range(1, 12)] last_11_month = [CommonUtil.get_month_offset(self.date_info, -i) for i in range(1, 12)]
# ===== 原逻辑(保留备查):读 dwt_aba_st_analytics 近11月 rank =====
# 问题:dwt_aba_st_analytics = 爬虫 ∩ ABA(inner join),某月爬虫漏抓 → 该词该月无行 →
# peak_month/常年可卖(all_year_text_flag)被漏抓月带偏(峰值月算错、常年可卖误判为否)。
# sql = f"""
# select
# search_term,
# rank,
# date_info
# from dwt_aba_st_analytics
# where site_name = '{self.site_name}'
# and date_type = 'month'
# and date_info in ({CommonUtil.list_to_insql(last_11_month)})
# """
# 【改】换源 dim_st_detail(ABA 清洗表,dwt 的 rank 本就取自它的 st_rank):
# dim_st_detail 不经爬虫过滤,漏抓月 ABA 侧仍有该词行 →
# 峰值月/常年可卖按 ABA 的 presence + 真实 rank 计算,不再受爬虫漏抓影响。
sql = f""" sql = f"""
select select
search_term, search_term,
rank, st_rank as rank,
date_info date_info
from dwt_aba_st_analytics from dim_st_detail
where site_name = '{self.site_name}' where site_name = '{self.site_name}'
and date_type = 'month' and date_type = 'month'
and date_info in ({CommonUtil.list_to_insql(last_11_month)}) and date_info in ({CommonUtil.list_to_insql(last_11_month)})
......
...@@ -60,25 +60,51 @@ class DwtSTBaseReport(object): ...@@ -60,25 +60,51 @@ class DwtSTBaseReport(object):
df_rank_sv.show(10, False) df_rank_sv.show(10, False)
# 搜索词主表 # 搜索词主表
# ===== 原逻辑(保留备查):从 dwt_aba_st_analytics 取词集合(st_key + search_term) =====
# 问题:dwt_aba_st_analytics = 爬虫 ∩ ABA,某月爬虫漏抓 → 该词该月不在 → 下面与排名 inner join 后,
# 趋势图那一月断格(排名很好的词突然缺一点)。
# sql2 = f"""
# select
# id as st_key,
# search_term
# from dwt_aba_st_analytics
# where site_name = '{self.site_name}'
# and date_type = '{self.date_type}'
# and date_info = '{self.date_info}';
# """
# 【改】st_key 改从 ods_st_key 取(全站"词→key"映射,与 dwt_st_sv_last365 / backfill 同款);
# 词集合不再被 dwt 卡,改由下面 dim_st_detail 的排名表(sql3)驱动 → inner join 后 = 当月 ABA 有 key 的词。
sql2 = f""" sql2 = f"""
select select
id as st_key, st_key,
search_term search_term
from dwt_aba_st_analytics from ods_st_key
where site_name = '{self.site_name}' where site_name = '{self.site_name}';
and date_type = '{self.date_type}'
and date_info = '{self.date_info}';
""" """
df_st_base = self.spark.sql(sql2).repartition(40, 'search_term').cache() df_st_base = self.spark.sql(sql2).repartition(40, 'search_term').cache()
print("搜索词主表:") print("搜索词主表:")
df_st_base.show(10, False) df_st_base.show(10, False)
# 读ods_brand_analytics表,获取报告中的搜索词+排名 # 读ABA搜索词+排名(口径统一到 dim_st_detail,词集合由它驱动)
# ===== 原逻辑(保留备查):从 ods_brand_analytics 取 search_term + rank =====
# 原本词集合被上面 dwt(sql2) 卡住,这里的 ods rank 只作补值;为口径统一、且 ods 月分区可能一词多行,
# 改从 dim_st_detail 取(排名同源:dim.st_rank ← ods.rank,且已按词去重、一词一月一行)。
# sql3 = f"""
# select
# search_term,
# rank as st_rank
# from ods_brand_analytics
# where site_name = '{self.site_name}'
# and date_type = '{self.date_type}'
# and date_info = '{self.date_info}';
# """
# 【改】改读 dim_st_detail(ABA 清洗表,st_rank 即排名):由它驱动词集合 →
# 当月 ABA 有的词趋势图都有排名点,漏抓月不再断格;与月表/年表口径统一到 dim_st_detail。
sql3 = f""" sql3 = f"""
select select
search_term, search_term,
rank as st_rank st_rank as st_rank
from ods_brand_analytics from dim_st_detail
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}';
......
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