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
b3463392
Commit
b3463392
authored
Jul 23, 2026
by
chenyuanjie
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
年度流量选品数据分批写入
parent
eeb57a26
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
90 additions
and
19 deletions
+90
-19
dwt_flow_asin_month.py
Pyspark_job/doris_handle/dwt_flow_asin_month.py
+90
-19
No files found.
Pyspark_job/doris_handle/dwt_flow_asin_month.py
View file @
b3463392
"""
"""
@Author : CT
@Author : CT
@Description : 月流量选品 Doris 落地脚本(月流程),同时关联产出年度流量选品数据
@Description : 月流量选品 Doris 落地脚本,同时关联产出年度流量选品数据
(原 doris_handle/dwt_flow_asin_year.py 的 Doris 编排部分合并进来,
因为月流程本身就需要关联部分年度聚合指标,拆开两个脚本维护反而麻烦)
流程:
流程:
[Step 1] Doris 建表 selection.{site}_flow_asin_month_{yyyy_mm}[_test]
[Step 1] Doris 建表 selection.{site}_flow_asin_month_{yyyy_mm}[_test]
[Step 2] 读 Hive dwt_flow_asin 月数据 + 字段规范化(与 dwt.us_flow_asin_30day
[Step 2] 读 Hive dwt_flow_asin 月数据 + 字段规范化(与 dwt.us_flow_asin_30day
...
@@ -16,7 +14,7 @@
...
@@ -16,7 +14,7 @@
site_name+date_info 分区同步过来,TRUNCATE 后全量重算
site_name+date_info 分区同步过来,TRUNCATE 后全量重算
[Step 6] Doris INSERT OVERWRITE 到 selection 月物化表
[Step 6] Doris INSERT OVERWRITE 到 selection 月物化表
selection.{site}_flow_asin_month_{yyyy_mm}[_test]
selection.{site}_flow_asin_month_{yyyy_mm}[_test]
[Step 7] Doris INSERT
OVERWRITE
到 selection 年表
[Step 7] Doris INSERT 到 selection 年表
selection.{site}_flow_asin_365day[_test]
selection.{site}_flow_asin_365day[_test]
[Step 8] 更新 MySQL workflow_everyday 流程记录表(月/年各写一条,仅 formal 模式)
[Step 8] 更新 MySQL workflow_everyday 流程记录表(月/年各写一条,仅 formal 模式)
依赖:dwt/dwt_flow_asin_year.py 必须已经跑完同一个 site_name+date_info,
依赖:dwt/dwt_flow_asin_year.py 必须已经跑完同一个 site_name+date_info,
...
@@ -776,7 +774,7 @@ CREATE TABLE IF NOT EXISTS `dwt`.`{table_name}`
...
@@ -776,7 +774,7 @@ CREATE TABLE IF NOT EXISTS `dwt`.`{table_name}`
) ENGINE=OLAP
) ENGINE=OLAP
UNIQUE KEY(`asin`)
UNIQUE KEY(`asin`)
COMMENT '流量选品年度聚合指标'
COMMENT '流量选品年度聚合指标'
DISTRIBUTED BY HASH(`asin`) BUCKETS
32
DISTRIBUTED BY HASH(`asin`) BUCKETS
8
PROPERTIES (
PROPERTIES (
"replication_num" = "3",
"replication_num" = "3",
"enable_unique_key_merge_on_write" = "true"
"enable_unique_key_merge_on_write" = "true"
...
@@ -807,8 +805,6 @@ def sync_extra_table(spark, site_name, date_info):
...
@@ -807,8 +805,6 @@ def sync_extra_table(spark, site_name, date_info):
df_raw
=
spark
.
sql
(
sqlQuery
=
sql
)
df_raw
=
spark
.
sql
(
sqlQuery
=
sql
)
# total_appear_month / peak_month_arr 在 Hive 那边存的是逗号拼接的 STRING
# total_appear_month / peak_month_arr 在 Hive 那边存的是逗号拼接的 STRING
# (dwt/dwt_flow_asin_year.py 里的说明),这里用 Spark 算子转回 ARRAY<INT> 再 to_json
# 成 "[1,2,3]"(与 img_type 处理方式一致,Doris StreamLoad ARRAY<INT> 要求)
df_extra
=
df_raw
.
withColumn
(
df_extra
=
df_raw
.
withColumn
(
'total_appear_month'
,
F
.
split
(
F
.
col
(
'total_appear_month'
),
','
)
'total_appear_month'
,
F
.
split
(
F
.
col
(
'total_appear_month'
),
','
)
)
.
withColumn
(
)
.
withColumn
(
...
@@ -827,8 +823,6 @@ def sync_extra_table(spark, site_name, date_info):
...
@@ -827,8 +823,6 @@ def sync_extra_table(spark, site_name, date_info):
print
(
f
"[同步365day_extra] 读取 Hive dwt_flow_asin_year[site_name={site_name}, date_info={date_info}]:"
print
(
f
"[同步365day_extra] 读取 Hive dwt_flow_asin_year[site_name={site_name}, date_info={date_info}]:"
f
"{row_count} 条"
)
f
"{row_count} 条"
)
if
row_count
==
0
:
if
row_count
==
0
:
# 读到0条大概率是 dwt/dwt_flow_asin_year.py 还没跑这个 site_name+date_info,
# 先报错,不要往下 TRUNCATE,避免把上个月还有效的数据清空
raise
ValueError
(
raise
ValueError
(
f
"Hive dwt_flow_asin_year[site_name={site_name}, date_info={date_info}] 读到0条,"
f
"Hive dwt_flow_asin_year[site_name={site_name}, date_info={date_info}] 读到0条,"
f
"请确认 dwt/dwt_flow_asin_year.py 是否已经跑完这个 site_name+date_info"
f
"请确认 dwt/dwt_flow_asin_year.py 是否已经跑完这个 site_name+date_info"
...
@@ -1276,15 +1270,14 @@ PROPERTIES (
...
@@ -1276,15 +1270,14 @@ PROPERTIES (
"""
"""
def
build_year_insert_overwrite_sql
(
site_name
,
tabl
e_name
):
def
_year_select_from_join_sql
(
sit
e_name
):
"""
构造 INSERT OVERWRITE SQL
"""
年表 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 等
"""
"""
return
f
"""
return
f
"""
INSERT OVERWRITE TABLE `selection`.`{table_name}`
SELECT
SELECT
f.asin,
f.asin,
f.parent_asin,
f.parent_asin,
...
@@ -1487,9 +1480,13 @@ LEFT JOIN `selection`.`user_mask_asin` uma ON f.asin = uma.asin
...
@@ -1487,9 +1480,13 @@ LEFT JOIN `selection`.`user_mask_asin` uma ON f.asin = uma.asin
LEFT JOIN `selection`.`user_mask_category` umc ON f.category_id = umc.category_id
LEFT JOIN `selection`.`user_mask_category` umc ON f.category_id = umc.category_id
-- ===== 拍卖/SKU =====
-- ===== 拍卖/SKU =====
LEFT JOIN `dwd`.`dwd_asin_auction` aa ON f.asin = aa.asin
LEFT JOIN `dwd`.`dwd_asin_auction` aa ON f.asin = aa.asin
-- ===== 品牌背书原因(按 asin 自己最新出现的月份匹配,f.date_info 每个 asin 可能不是同一个月)=====
-- ===== 品牌推荐原因——
LEFT JOIN `dwd`.`dwd_st_brand_badge` bb
LEFT JOIN (
ON f.brand = bb.brand AND bb.site_name = '{site_name}' AND bb.date_info = DATE_FORMAT(f.date_info, '
%
Y-
%
m')
SELECT brand, brand_badge_reason
FROM `dwd`.`dwd_st_brand_badge`
WHERE site_name = '{site_name}'
AND date_info = (SELECT MAX(date_info) FROM `dwd`.`dwd_st_brand_badge` WHERE site_name = '{site_name}')
) bb ON f.brand = bb.brand
-- ===== AI分析数据 =====
-- ===== AI分析数据 =====
LEFT JOIN (
LEFT JOIN (
SELECT
SELECT
...
@@ -1519,6 +1516,81 @@ LEFT JOIN (
...
@@ -1519,6 +1516,81 @@ LEFT JOIN (
"""
"""
def
build_year_insert_overwrite_sql
(
site_name
,
table_name
):
"""整表覆盖版(非分批):保留供手工排查/小站点场景使用,正式流程走下面的分批版本
write_year_table_batched,不再直接调用这个函数"""
return
f
"INSERT OVERWRITE TABLE `selection`.`{table_name}`
\n
{_year_select_from_join_sql(site_name)}"
def
build_year_batch_insert_sql
(
site_name
,
table_name
,
date_info_batch
):
"""构造单个 date_info 批次的 INSERT INTO SQL
:param date_info_batch: dwt.{site}_flow_asin_365day 里的某个 date_info 取值(DATE,如 '2026-05-01')
"""
return
(
f
"INSERT INTO `selection`.`{table_name}`
\n
"
f
"{_year_select_from_join_sql(site_name)}"
f
"WHERE f.date_info = '{date_info_batch}'
\n
"
)
def
get_year_batch_date_infos
(
site_name
):
"""获取 dwt.{site_name}_flow_asin_365day 当前实际存在的 date_info 取值(DATE,如
'2026-05-01'),升序排列,作为年表分批 INSERT 的批次键(近12个月窗口,最多约12批)"""
stg_table
=
f
"{site_name}_flow_asin_365day"
rows
=
_query_doris
(
f
"SELECT DISTINCT date_info FROM `dwt`.`{stg_table}` ORDER BY date_info"
)
return
[
str
(
row
[
0
])
for
row
in
rows
]
def
write_year_table_batched
(
site_name
,
live_table
):
"""按 date_info 分批写入 selection 年表,避免一次性 13 路 JOIN 大 SQL 导致的内存爆炸/超时:
1) DROP + 重建 copy 表(保证 schema 跟当前 DDL 一致)
2) 按 dwt.{site}_flow_asin_365day 现有的 date_info 取值逐批 INSERT INTO copy 表
3) 数据量校验:copy 表最终行数 应等于 dwt.{site}_flow_asin_365day 原表行数——驱动表
4) 校验通过后 ALTER TABLE ... REPLACE WITH TABLE ... PROPERTIES('swap'='true') 原子
切换成正式表,copy 表则拿到正式表交换前的旧数据(可回滚,不会立刻销毁)
"""
copy_table
=
f
"{live_table}_copy"
stg_table
=
f
"{site_name}_flow_asin_365day"
print
(
f
"[Step 7] 重建 copy 表 selection.{copy_table}"
)
_exec_doris_sql
([
f
"DROP TABLE IF EXISTS `selection`.`{copy_table}`"
])
_exec_doris_sql
([
build_year_create_table_sql
(
copy_table
)])
date_infos
=
get_year_batch_date_infos
(
site_name
)
print
(
f
"[Step 7] 待分批写入 selection.{copy_table} 的 date_info 批次(共 {len(date_infos)} 批):{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}'"
)
_exec_doris_sql
([
build_year_batch_insert_sql
(
site_name
,
copy_table
,
date_info_batch
)])
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
]
print
(
f
"[Step 7] 数据量校验:dwt.{stg_table}={dwt_count},selection.{copy_table}={copy_count}"
)
if
copy_count
==
0
:
# 兜底:即使 dwt_count 恰好也是0导致下面的相等校验通过,也绝不能把空表swap成正式表
raise
ValueError
(
f
"[Step 7] copy 表 selection.{copy_table} 写入后行数为0,已中止交换,"
f
"避免把正式表 selection.{live_table} 清空,请检查 dwt.{stg_table} 是否为空或分批写入是否异常"
)
if
copy_count
!=
dwt_count
:
raise
ValueError
(
f
"[Step 7] copy 表行数({copy_count})与 dwt.{stg_table} 原表行数({dwt_count})不一致,"
f
"疑似分批写入有遗漏/重复,已中止交换,请人工排查(copy 表已保留,不会自动清理)"
)
print
(
f
"[Step 7] 确保正式表 selection.{live_table} 存在"
)
_exec_doris_sql
([
build_year_create_table_sql
(
live_table
)])
print
(
f
"[Step 7] REPLACE WITH TABLE:selection.{live_table} <-> selection.{copy_table}"
)
# 注意:REPLACE WITH TABLE 后面的表名不能带库名前缀,只能是裸表名(同库),
# 带上 `selection`. 前缀会导致 Doris 解析报错 "mismatched input '.'"
_exec_doris_sql
([
f
"ALTER TABLE `selection`.`{live_table}` REPLACE WITH TABLE `{copy_table}` "
f
"PROPERTIES('swap' = 'false')"
])
print
(
f
"[Step 7] selection.{live_table} 已切换为最新数据,"
f
"selection.{copy_table} 保留的是交换前的旧数据(可用于回滚)"
)
# ============================================================
# ============================================================
# [Step 8] 流程记录表更新(月+年各写一条,仅 formal 模式)
# [Step 8] 流程记录表更新(月+年各写一条,仅 formal 模式)
# ============================================================
# ============================================================
...
@@ -1596,10 +1668,9 @@ def main(site_name, date_info, result_type='formal'):
...
@@ -1596,10 +1668,9 @@ def main(site_name, date_info, result_type='formal'):
print
(
f
"[Step 6] Doris INSERT OVERWRITE selection.{month_selection_table}"
)
print
(
f
"[Step 6] Doris INSERT OVERWRITE selection.{month_selection_table}"
)
_exec_doris_sql
([
build_month_insert_overwrite_sql
(
site_name
,
month_selection_table
,
date_info
)])
_exec_doris_sql
([
build_month_insert_overwrite_sql
(
site_name
,
month_selection_table
,
date_info
)])
# ===== [Step 7] Doris INSERT OVERWRITE 到 selection 年表 =====
# ===== [Step 7] 按 date_info 分批写入 selection 年表的 copy 表,再原子交换成正式表 =====
print
(
f
"[Step 7] Doris INSERT OVERWRITE selection.{year_selection_table}"
)
print
(
f
"[Step 7] 分批写入 selection.{year_selection_table} 的 copy 表并原子交换"
)
_exec_doris_sql
([
build_year_create_table_sql
(
year_selection_table
)])
write_year_table_batched
(
site_name
,
year_selection_table
)
_exec_doris_sql
([
build_year_insert_overwrite_sql
(
site_name
,
year_selection_table
)])
# ===== [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