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
84e3915a
Commit
84e3915a
authored
Aug 13, 2026
by
chenyuanjie
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
年度流量选品-上架/追踪时间类型修复
parent
93af9a0c
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
31 additions
and
20 deletions
+31
-20
dwt_flow_asin_month.py
Pyspark_job/doris_handle/dwt_flow_asin_month.py
+31
-20
No files found.
Pyspark_job/doris_handle/dwt_flow_asin_month.py
View file @
84e3915a
...
@@ -1200,13 +1200,19 @@ PROPERTIES (
...
@@ -1200,13 +1200,19 @@ PROPERTIES (
"""
"""
def
_year_select_from_join_sql
(
site_name
):
def
_year_select_from_join_sql
(
site_name
,
latest_date_info
):
"""年表 SELECT 主体(不含 INSERT 前缀,不含 WHERE),供整表覆盖 / 分批插入两种场景共用:
"""年表 SELECT 主体(不含 INSERT 前缀,不含 WHERE),供整表覆盖 / 分批插入两种场景共用:
- 主体: dwt.{site}_flow_asin_365day (每 asin 最新月快照,字段名沿用 dwt 月表)
- 主体: dwt.{site}_flow_asin_365day (每 asin 最新月快照,字段名沿用 dwt 月表)
- 年度聚合: dwt.{site}_flow_asin_365day_extra (Hive 算好经 sync_extra_table 同步来的年度指标,
- 年度聚合: dwt.{site}_flow_asin_365day_extra (Hive 算好经 sync_extra_table 同步来的年度指标,
含周期性/季节性/峰值)
含周期性/季节性/峰值)
- 外层 LEFT JOIN: profit_rate / keepa / brand_alert / self_asin / category_hide / user_mask / auction 等
- 外层 LEFT JOIN: profit_rate / keepa / brand_alert / self_asin / category_hide / user_mask / auction 等
- launch_time_type / tracking_since_type 的天数基准统一用 latest_date_info(本次跑批最新月份)
的次月初,而不是各批次自身的 f.asin_crawl_date:dwt.{site}_flow_asin_365day 是按 asin 去重只保留
最新一次快照的表,长期未被重新抓取的 asin 的 asin_crawl_date 会停留在很久之前,用它做基准会导致
这两个分档字段冻结在过去、不随时间推进
:param latest_date_info: main() 的 date_info 参数(yyyy-MM),与分批用的历史 date_info_batch 无关
"""
"""
type_anchor_date
=
f
"DATE_ADD(CAST(CONCAT('{latest_date_info}', '-01') AS DATE), INTERVAL 1 MONTH)"
return
f
"""
return
f
"""
SELECT
SELECT
f.asin,
f.asin,
...
@@ -1253,12 +1259,12 @@ SELECT
...
@@ -1253,12 +1259,12 @@ SELECT
COALESCE(f.launch_time, kp.keepa_launch_time) AS launch_time,
COALESCE(f.launch_time, kp.keepa_launch_time) AS launch_time,
CASE
CASE
WHEN COALESCE(f.launch_time, kp.keepa_launch_time) IS NULL THEN 0
WHEN COALESCE(f.launch_time, kp.keepa_launch_time) IS NULL THEN 0
WHEN DATEDIFF(
f.asin_crawl_date
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 30 THEN 1
WHEN DATEDIFF(
{type_anchor_date}
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 30 THEN 1
WHEN DATEDIFF(
f.asin_crawl_date
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 90 THEN 2
WHEN DATEDIFF(
{type_anchor_date}
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 90 THEN 2
WHEN DATEDIFF(
f.asin_crawl_date
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 180 THEN 3
WHEN DATEDIFF(
{type_anchor_date}
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 180 THEN 3
WHEN DATEDIFF(
f.asin_crawl_date
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 360 THEN 4
WHEN DATEDIFF(
{type_anchor_date}
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 360 THEN 4
WHEN DATEDIFF(
f.asin_crawl_date
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 720 THEN 5
WHEN DATEDIFF(
{type_anchor_date}
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 720 THEN 5
WHEN DATEDIFF(
f.asin_crawl_date
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 1080 THEN 6
WHEN DATEDIFF(
{type_anchor_date}
, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 1080 THEN 6
ELSE 7
ELSE 7
END AS launch_time_type,
END AS launch_time_type,
f.img_url,
f.img_url,
...
@@ -1294,12 +1300,12 @@ SELECT
...
@@ -1294,12 +1300,12 @@ SELECT
FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60) AS tracking_since,
FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60) AS tracking_since,
CASE
CASE
WHEN kp.tracking_since IS NULL OR kp.tracking_since <= 0 THEN 0
WHEN kp.tracking_since IS NULL OR kp.tracking_since <= 0 THEN 0
WHEN DATEDIFF(
f.asin_crawl_date
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 30 THEN 1
WHEN DATEDIFF(
{type_anchor_date}
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 30 THEN 1
WHEN DATEDIFF(
f.asin_crawl_date
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 90 THEN 2
WHEN DATEDIFF(
{type_anchor_date}
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 90 THEN 2
WHEN DATEDIFF(
f.asin_crawl_date
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 180 THEN 3
WHEN DATEDIFF(
{type_anchor_date}
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 180 THEN 3
WHEN DATEDIFF(
f.asin_crawl_date
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 360 THEN 4
WHEN DATEDIFF(
{type_anchor_date}
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 360 THEN 4
WHEN DATEDIFF(
f.asin_crawl_date
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 720 THEN 5
WHEN DATEDIFF(
{type_anchor_date}
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 720 THEN 5
WHEN DATEDIFF(
f.asin_crawl_date
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 1080 THEN 6
WHEN DATEDIFF(
{type_anchor_date}
, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 1080 THEN 6
ELSE 7
ELSE 7
END AS tracking_since_type,
END AS tracking_since_type,
kp.package_length,
kp.package_length,
...
@@ -1413,19 +1419,22 @@ LEFT JOIN (
...
@@ -1413,19 +1419,22 @@ LEFT JOIN (
"""
"""
def
build_year_insert_overwrite_sql
(
site_name
,
table_name
):
def
build_year_insert_overwrite_sql
(
site_name
,
table_name
,
latest_date_info
):
"""整表覆盖版(非分批):保留供手工排查/小站点场景使用,正式流程走下面的分批版本
"""整表覆盖版(非分批):保留供手工排查/小站点场景使用,正式流程走下面的分批版本
write_year_table_batched,不再直接调用这个函数"""
write_year_table_batched,不再直接调用这个函数"""
return
f
"INSERT OVERWRITE TABLE `selection`.`{table_name}`
\n
{_year_select_from_join_sql(site_name)}"
return
f
"INSERT OVERWRITE TABLE `selection`.`{table_name}`
\n
{_year_select_from_join_sql(site_name
, latest_date_info
)}"
def
build_year_batch_insert_sql
(
site_name
,
table_name
,
date_info_batch
):
def
build_year_batch_insert_sql
(
site_name
,
table_name
,
date_info_batch
,
latest_date_info
):
"""构造单个 date_info 批次的 INSERT INTO SQL
"""构造单个 date_info 批次的 INSERT INTO SQL
:param date_info_batch: dwt.{site}_flow_asin_365day 里的某个 date_info 取值(DATE,如 '2026-05-01')
:param date_info_batch: dwt.{site}_flow_asin_365day 里的某个 date_info 取值(DATE,如 '2026-05-01')
:param latest_date_info: main() 的 date_info 参数(yyyy-MM,本次跑批最新月份),传给
_year_select_from_join_sql 做 launch_time_type/tracking_since_type 的统一天数基准,
与 date_info_batch(历史批次键)是两个不同的概念
"""
"""
return
(
return
(
f
"INSERT INTO `selection`.`{table_name}`
\n
"
f
"INSERT INTO `selection`.`{table_name}`
\n
"
f
"{_year_select_from_join_sql(site_name)}"
f
"{_year_select_from_join_sql(site_name
, latest_date_info
)}"
f
"WHERE f.date_info = '{date_info_batch}'
\n
"
f
"WHERE f.date_info = '{date_info_batch}'
\n
"
)
)
...
@@ -1438,13 +1447,15 @@ def get_year_batch_date_infos(site_name):
...
@@ -1438,13 +1447,15 @@ def get_year_batch_date_infos(site_name):
return
[
str
(
row
[
0
])
for
row
in
rows
]
return
[
str
(
row
[
0
])
for
row
in
rows
]
def
write_year_table_batched
(
site_name
,
live_table
):
def
write_year_table_batched
(
site_name
,
live_table
,
latest_date_info
):
"""按 date_info 分批写入 selection 年表,避免一次性 13 路 JOIN 大 SQL 导致的内存爆炸/超时:
"""按 date_info 分批写入 selection 年表,避免一次性 13 路 JOIN 大 SQL 导致的内存爆炸/超时:
1) DROP + 重建 copy 表(保证 schema 跟当前 DDL 一致)
1) DROP + 重建 copy 表(保证 schema 跟当前 DDL 一致)
2) 按 dwt.{site}_flow_asin_365day 现有的 date_info 取值逐批 INSERT INTO copy 表
2) 按 dwt.{site}_flow_asin_365day 现有的 date_info 取值逐批 INSERT INTO copy 表
3) 数据量校验:copy 表最终行数 应等于 dwt.{site}_flow_asin_365day 原表行数——驱动表
3) 数据量校验:copy 表最终行数 应等于 dwt.{site}_flow_asin_365day 原表行数——驱动表
4) 校验通过后 ALTER TABLE ... REPLACE WITH TABLE ... PROPERTIES('swap'='true') 原子
4) 校验通过后 ALTER TABLE ... REPLACE WITH TABLE ... PROPERTIES('swap'='true') 原子
切换成正式表,copy 表则拿到正式表交换前的旧数据(可回滚,不会立刻销毁)
切换成正式表,copy 表则拿到正式表交换前的旧数据(可回滚,不会立刻销毁)
:param latest_date_info: main() 的 date_info 参数(yyyy-MM,本次跑批最新月份),所有批次统一用它
算 launch_time_type/tracking_since_type 的天数基准(次月初),而不是各批次自身的历史 date_info_batch
"""
"""
copy_table
=
f
"{live_table}_copy"
copy_table
=
f
"{live_table}_copy"
stg_table
=
f
"{site_name}_flow_asin_365day"
stg_table
=
f
"{site_name}_flow_asin_365day"
...
@@ -1457,7 +1468,7 @@ def write_year_table_batched(site_name, live_table):
...
@@ -1457,7 +1468,7 @@ def write_year_table_batched(site_name, live_table):
print
(
f
"[Step 7] 待分批写入 selection.{copy_table} 的 date_info 批次(共 {len(date_infos)} 批):{date_infos}"
)
print
(
f
"[Step 7] 待分批写入 selection.{copy_table} 的 date_info 批次(共 {len(date_infos)} 批):{date_infos}"
)
for
date_info_batch
in
date_infos
:
for
date_info_batch
in
date_infos
:
print
(
f
"[Step 7] INSERT INTO selection.{copy_table} WHERE f.date_info = '{date_info_batch}'"
)
print
(
f
"[Step 7] INSERT INTO selection.{copy_table} WHERE f.date_info = '{date_info_batch}'"
)
_exec_doris_sql
([
build_year_batch_insert_sql
(
site_name
,
copy_table
,
date_info_batch
)])
_exec_doris_sql
([
build_year_batch_insert_sql
(
site_name
,
copy_table
,
date_info_batch
,
latest_date_info
)])
dwt_count
=
_query_doris
(
f
"SELECT COUNT(1) FROM `dwt`.`{stg_table}`"
)[
0
][
0
]
dwt_count
=
_query_doris
(
f
"SELECT COUNT(1) FROM `dwt`.`{stg_table}`"
)[
0
][
0
]
copy_count
=
_query_doris
(
f
"SELECT COUNT(1) FROM `selection`.`{copy_table}`"
)[
0
][
0
]
copy_count
=
_query_doris
(
f
"SELECT COUNT(1) FROM `selection`.`{copy_table}`"
)[
0
][
0
]
...
@@ -1566,7 +1577,7 @@ def main(site_name, date_info, result_type='formal'):
...
@@ -1566,7 +1577,7 @@ def main(site_name, date_info, result_type='formal'):
# ===== [Step 7] 按 date_info 分批写入 selection 年表的 copy 表,再原子交换成正式表 =====
# ===== [Step 7] 按 date_info 分批写入 selection 年表的 copy 表,再原子交换成正式表 =====
print
(
f
"[Step 7] 分批写入 selection.{year_selection_table} 的 copy 表并原子交换"
)
print
(
f
"[Step 7] 分批写入 selection.{year_selection_table} 的 copy 表并原子交换"
)
write_year_table_batched
(
site_name
,
year_selection_table
)
write_year_table_batched
(
site_name
,
year_selection_table
,
date_info
)
# ===== [Step 8] 流程记录表更新(月+年各一条,仅 formal 模式)=====
# ===== [Step 8] 流程记录表更新(月+年各一条,仅 formal 模式)=====
modify_mission_record_status
(
site_name
,
date_info
,
result_type
)
modify_mission_record_status
(
site_name
,
date_info
,
result_type
)
...
...
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