Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
A
Amazon-Selection-Data
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
abel_cjy
Amazon-Selection-Data
Commits
7420cbb8
Commit
7420cbb8
authored
Aug 18, 2026
by
hejiangming
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
周搜索词新增创历史新高 + 历史首次出现
parent
72782693
Hide whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
133 additions
and
2 deletions
+133
-2
dim_st_detail_week.py
Pyspark_job/dim/dim_st_detail_week.py
+127
-1
dim_st_detail_week.py
Pyspark_job/sqoop_export/dim_st_detail_week.py
+6
-1
No files found.
Pyspark_job/dim/dim_st_detail_week.py
View file @
7420cbb8
...
@@ -16,6 +16,10 @@ class DwtStDetailWeek(object):
...
@@ -16,6 +16,10 @@ class DwtStDetailWeek(object):
def
__init__
(
self
,
site_name
,
date_type
,
date_info
):
def
__init__
(
self
,
site_name
,
date_type
,
date_info
):
super
()
.
__init__
()
super
()
.
__init__
()
# 传入周号可能是 2026-1 也可能是 2026-01,统一补零成 YYYY-WW 再往下走
# 下面所有时间判断都是字符串比较,不补零会比错:'2026-32' <= '2026-1' 的结果是 false
year_part
,
week_part
=
date_info
.
split
(
'-'
)
date_info
=
f
"{int(year_part)}-{int(week_part):02d}"
self
.
site_name
=
site_name
self
.
site_name
=
site_name
self
.
date_type
=
date_type
self
.
date_type
=
date_type
self
.
date_info
=
date_info
self
.
date_info
=
date_info
...
@@ -59,6 +63,14 @@ class DwtStDetailWeek(object):
...
@@ -59,6 +63,14 @@ class DwtStDetailWeek(object):
'asin3'
,
'product_title3'
,
'click_share3'
,
'conversion_share3'
,
'brand3'
,
'category3'
,
'quantity_being_sold'
]
'asin3'
,
'product_title3'
,
'click_share3'
,
'conversion_share3'
,
'brand3'
,
'category3'
,
'quantity_being_sold'
]
self
.
sp_symbols
=
[]
self
.
sp_symbols
=
[]
# 全历史首次出现 + 创历史新高:仅 us 且 >= 2026 年第 1 周计算真值,其余占位 -1
date_year
=
int
(
self
.
date_info
.
split
(
'-'
)[
0
])
self
.
is_new_field_calc
=
(
self
.
site_name
==
'us'
and
date_year
>=
2026
)
# is_first_ever_text 用:2023-01 之前出现过的 search_term(2023-01 之后那段由 df_dim_st_rank 的聚合结果覆盖)
self
.
df_dim_st_old
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
# is_first_ever_text + new_high 共用:2023-01 周起至当周的 (词, 周, 排名)
self
.
df_dim_st_rank
=
self
.
spark
.
sql
(
f
"select 1+1;"
)
def
st_word_count
(
self
,
sp_symbols
):
def
st_word_count
(
self
,
sp_symbols
):
def
udf_st_word_count
(
name
):
def
udf_st_word_count
(
name
):
# 特殊字符基准列表---迁移到数据库维护 -已处理
# 特殊字符基准列表---迁移到数据库维护 -已处理
...
@@ -223,6 +235,42 @@ class DwtStDetailWeek(object):
...
@@ -223,6 +235,42 @@ class DwtStDetailWeek(object):
# self.df_st_detail_3_week_ago.show(10, True)
# self.df_st_detail_3_week_ago.show(10, True)
self
.
df_st_detail_week
.
unpersist
()
self
.
df_st_detail_week
.
unpersist
()
# is_first_ever_text / new_high 的历史对比数据,都取自 ods_brand_analytics(与趋势图 dwt_st_base_report_week 同源)
# 以 2023-01 为界读两段:2023-01 之前只要词集合,2023-01 起要带 rank 的明细
# 分界点选 2023-01 是因为 new_high 的对比范围就是从这周起,两段拼起来正好等于 is_first_ever_text 要的全历史
if
self
.
is_new_field_calc
:
# 2023-01 之前(us 最早 2020-44)出现过的搜索词,只服务 is_first_ever_text
# 只要词集合、不取 rank,distinct 本身就去重,不需要按 updated_time 挑最新那条
# search_term 要按上面本周数据同样的规则清掉 \x00,否则带脏字节的词两边 join 不上、被误判成首次出现
# 清洗必须在 distinct 之前:'abc' 和 'abc\x00' 清洗后同名,先 distinct 会留下重复行,left join 时把左表撑出多行
self
.
df_dim_st_old
=
self
.
spark
.
sql
(
f
"""
select distinct search_term
from ods_brand_analytics
where site_name = '{self.site_name}'
and date_type = 'week'
and date_info < '2023-01'
and search_term is not null
and rank > 0
"""
)
.
withColumn
(
'search_term'
,
F
.
regexp_replace
(
F
.
col
(
'search_term'
),
'
\x00
'
,
''
)
)
.
distinct
()
# 2023-01 周起至当周的 (词, 周, 排名),new_high 和 is_first_ever_text 共用这一份
# 带上 updated_time:同一 (search_term, date_info) 万一有重复行,靠它挑最新那条,口径和上面本周数据的去重一致
# search_term 同样清 \x00 对齐本周数据;这份不用去重,后面 groupBy('search_term') 会把清洗后同名的行并到一组
self
.
df_dim_st_rank
=
self
.
spark
.
sql
(
f
"""
select search_term, rank, date_info, updated_time
from ods_brand_analytics
where site_name = '{self.site_name}'
and date_type = 'week'
and date_info >= '2023-01'
and date_info <= '{self.date_info}'
and search_term is not null
and rank > 0
"""
)
.
withColumn
(
'search_term'
,
F
.
regexp_replace
(
F
.
col
(
'search_term'
),
'
\x00
'
,
''
)
)
def
handle_st_flag
(
self
):
def
handle_st_flag
(
self
):
# 热搜词:最近4周中,出现的次数大于80%,即近4周都出现
# 热搜词:最近4周中,出现的次数大于80%,即近4周都出现
df_hot_search_term
=
self
.
df_st_detail_last_4_week
.
groupBy
(
'search_term'
)
.
agg
(
df_hot_search_term
=
self
.
df_st_detail_last_4_week
.
groupBy
(
'search_term'
)
.
agg
(
...
@@ -355,6 +403,74 @@ class DwtStDetailWeek(object):
...
@@ -355,6 +403,74 @@ class DwtStDetailWeek(object):
F
.
round
(
F
.
col
(
'conversion_share1'
)
+
F
.
col
(
'conversion_share2'
)
+
F
.
col
(
'conversion_share3'
),
4
)
F
.
round
(
F
.
col
(
'conversion_share1'
)
+
F
.
col
(
'conversion_share2'
)
+
F
.
col
(
'conversion_share3'
),
4
)
)
.
cache
()
)
.
cache
()
NEW_HIGH_COLS
=
[
'new_high_this_week'
,
'new_high_1_week_ago'
,
'new_high_2_week_ago'
,
'new_high_3_week_ago'
]
# 全历史首次出现 + 创历史新高,5 个字段共用一次历史聚合
# 合起来算是因为两者的历史范围重叠:is_first_ever_text 要"当周之前所有周",new_high 要"2023-01 起"。
# 下面 groupBy 出的 min_before_w0(当周之前的历史最佳 rank)非空,就等价于"该词 2023-01 起、当周之前上过榜",
# 所以 is_first_ever_text 只要再补一个"2023-01 之前出现过"的词集合,不必把 2023 年以后的数据再扫一遍
def
handle_new_fields
(
self
):
target_weeks
=
[
self
.
date_info
,
self
.
date_info_last_week
,
self
.
date_info_2_week_ago
,
self
.
date_info_3_week_ago
]
# 一次 groupBy 出 8 列:每个目标周的 rank + 该周之前的历史最佳 rank
# 用条件聚合而不是开窗,是因为聚合带 map 端预聚合、且只 shuffle 一次,全历史数据不会被反复扫
# rank_w:max(struct(updated_time, rank)) 取 updated_time 最新那条的 rank——struct 按字段顺序比大小,
# 等价于本周数据那套 row_number(order by updated_time desc) 去重。当前 ods 无重复行,这里是防未来出现重复
# min_before:历史最佳排名,重复行取 min 本就是要的语义,不用挑最新
agg_cols
=
[]
for
i
,
week
in
enumerate
(
target_weeks
):
if
week
is
None
:
continue
agg_cols
.
append
(
F
.
max
(
F
.
when
(
F
.
col
(
'date_info'
)
==
week
,
F
.
struct
(
'updated_time'
,
'rank'
)))
.
getField
(
'rank'
)
.
alias
(
f
'rank_w{i}'
)
)
agg_cols
.
append
(
F
.
min
(
F
.
when
(
F
.
col
(
'date_info'
)
<
week
,
F
.
col
(
'rank'
)))
.
alias
(
f
'min_before_w{i}'
))
df_agg
=
self
.
df_dim_st_rank
.
groupBy
(
'search_term'
)
.
agg
(
*
agg_cols
)
# 逐周判创历史新高标记,rank 越小越好
for
i
,
week
in
enumerate
(
target_weeks
):
if
week
is
None
:
# 日历表往前查不到该周(数据起点边界),无从判断,按没创新高处理
df_agg
=
df_agg
.
withColumn
(
self
.
NEW_HIGH_COLS
[
i
],
F
.
lit
(
0
))
continue
df_agg
=
df_agg
.
withColumn
(
self
.
NEW_HIGH_COLS
[
i
],
F
.
when
(
F
.
col
(
f
'rank_w{i}'
)
.
isNull
(),
F
.
lit
(
0
))
# 那周没上榜
.
when
(
F
.
col
(
f
'min_before_w{i}'
)
.
isNull
(),
F
.
lit
(
1
))
# 那周之前没有历史 = 首次上榜
.
when
(
F
.
col
(
f
'rank_w{i}'
)
<=
F
.
col
(
f
'min_before_w{i}'
),
F
.
lit
(
1
))
# 追平历史最佳也算创新高
.
otherwise
(
F
.
lit
(
0
))
)
# target_weeks[0] 恒为当周、不可能是 None,所以 min_before_w0 一定存在
df_agg
=
df_agg
.
select
(
'search_term'
,
F
.
col
(
'min_before_w0'
)
.
isNotNull
()
.
alias
(
'seen_since_2023'
),
*
self
.
NEW_HIGH_COLS
)
df_old
=
self
.
df_dim_st_old
.
withColumn
(
'seen_before_2023'
,
F
.
lit
(
True
))
self
.
df_st_detail
=
self
.
df_st_detail
\
.
join
(
df_agg
,
on
=
'search_term'
,
how
=
'left'
)
\
.
join
(
df_old
,
on
=
'search_term'
,
how
=
'left'
)
# 两段历史任一段上过榜 = 不是首次;两段都 join 不上(coalesce 成 false)= 全历史首次
self
.
df_st_detail
=
self
.
df_st_detail
.
withColumn
(
'is_first_ever_text'
,
F
.
when
(
F
.
coalesce
(
F
.
col
(
'seen_since_2023'
),
F
.
lit
(
False
))
|
F
.
coalesce
(
F
.
col
(
'seen_before_2023'
),
F
.
lit
(
False
)),
F
.
lit
(
0
)
)
.
otherwise
(
F
.
lit
(
1
))
)
.
drop
(
'seen_since_2023'
,
'seen_before_2023'
)
\
.
fillna
({
col
:
0
for
col
in
self
.
NEW_HIGH_COLS
})
# 非 us 或早于 2026-01:5 个字段全填 -1 占位
def
handle_new_field_padding
(
self
):
self
.
df_st_detail
=
self
.
df_st_detail
\
.
withColumn
(
'is_first_ever_text'
,
F
.
lit
(
-
1
))
for
col
in
self
.
NEW_HIGH_COLS
:
self
.
df_st_detail
=
self
.
df_st_detail
.
withColumn
(
col
,
F
.
lit
(
-
1
))
def
save_data
(
self
):
def
save_data
(
self
):
self
.
df_save
=
self
.
df_st_detail
.
filter
(
self
.
df_save
=
self
.
df_st_detail
.
filter
(
'length(asin1) <= 10 AND length(asin2) <= 10 AND length(asin3) <= 10'
'length(asin1) <= 10 AND length(asin2) <= 10 AND length(asin3) <= 10'
...
@@ -420,7 +536,12 @@ class DwtStDetailWeek(object):
...
@@ -420,7 +536,12 @@ class DwtStDetailWeek(object):
'rank_change_3_week_ago'
,
'rank_change_3_week_ago'
,
'rank_rate_1_week_ago'
,
'rank_rate_1_week_ago'
,
'rank_rate_2_week_ago'
,
'rank_rate_2_week_ago'
,
'rank_rate_3_week_ago'
'rank_rate_3_week_ago'
,
'is_first_ever_text'
,
'new_high_this_week'
,
'new_high_1_week_ago'
,
'new_high_2_week_ago'
,
'new_high_3_week_ago'
)
.
withColumn
(
)
.
withColumn
(
'site_name'
,
F
.
lit
(
self
.
site_name
)
'site_name'
,
F
.
lit
(
self
.
site_name
)
)
.
withColumn
(
)
.
withColumn
(
...
@@ -453,6 +574,11 @@ class DwtStDetailWeek(object):
...
@@ -453,6 +574,11 @@ class DwtStDetailWeek(object):
self
.
handle_rank_rate
()
self
.
handle_rank_rate
()
# 计算其他
# 计算其他
self
.
handle_other
()
self
.
handle_other
()
# 全历史首次出现 + 创历史新高
if
self
.
is_new_field_calc
:
self
.
handle_new_fields
()
else
:
self
.
handle_new_field_padding
()
# 数据落盘
# 数据落盘
self
.
save_data
()
self
.
save_data
()
...
...
Pyspark_job/sqoop_export/dim_st_detail_week.py
View file @
7420cbb8
...
@@ -102,7 +102,12 @@ if __name__ == '__main__':
...
@@ -102,7 +102,12 @@ if __name__ == '__main__':
'rank_change_3_week_ago'
,
'rank_change_3_week_ago'
,
'rank_rate_1_week_ago'
,
'rank_rate_1_week_ago'
,
'rank_rate_2_week_ago'
,
'rank_rate_2_week_ago'
,
'rank_rate_3_week_ago'
'rank_rate_3_week_ago'
,
'is_first_ever_text'
,
'new_high_this_week'
,
'new_high_1_week_ago'
,
'new_high_2_week_ago'
,
'new_high_3_week_ago'
],
],
partition_dict
=
{
partition_dict
=
{
"site_name"
:
site_name
,
"site_name"
:
site_name
,
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment