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
e9f81eb6
Commit
e9f81eb6
authored
Sep 21, 2026
by
chenyuanjie
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
流量选品-更新Doris兼容重复导出场景
parent
6a615706
Show whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
13 additions
and
4 deletions
+13
-4
dwt_flow_asin_month.py
Pyspark_job/doris_handle/dwt_flow_asin_month.py
+13
-4
No files found.
Pyspark_job/doris_handle/dwt_flow_asin_month.py
View file @
e9f81eb6
...
@@ -473,8 +473,14 @@ def read_and_normalize_month_data(spark, site_name, date_info):
...
@@ -473,8 +473,14 @@ def read_and_normalize_month_data(spark, site_name, date_info):
# [Step 3] 写入 Doris dwt 主表
# [Step 3] 写入 Doris dwt 主表
# ============================================================
# ============================================================
def
write_dwt_month_table
(
df_save
,
doris_table
):
def
write_dwt_month_table
(
df_save
,
doris_table
,
date_info
):
"""写入 Doris dwt 主表 dwt.{site}_flow_asin_month,写完 unpersist 传入的 df_save"""
"""写入 Doris dwt 主表 dwt.{site}_flow_asin_month:
先按 date_info 精确清理该月分区旧数据,再整批写入,保证与 Hive dwt_flow_asin 当月数据完全一致
(写入方式是 append/upsert,不会自动删除多余行;若不先清理,重跑时 Hive 侧本月 asin 集合缩小会导致 Doris 残留多余 asin)
写完 unpersist 传入的 df_save"""
print
(
f
"[Step 3] 清理 {DORIS_DB}.{doris_table} 旧数据 WHERE date_info='{date_info}'"
)
_exec_doris_sql
([
f
"DELETE FROM `{DORIS_DB}`.`{doris_table}` WHERE date_info = '{date_info}'"
])
table_columns
=
(
table_columns
=
(
"date_info, asin, ao_val, zr_counts, sp_counts, sb_counts, vi_counts, bs_counts, ac_counts, tr_counts, er_counts, "
"date_info, asin, ao_val, zr_counts, sp_counts, sb_counts, vi_counts, bs_counts, ac_counts, tr_counts, er_counts, "
"bsr_orders, bsr_orders_sale, title, title_len, price, rating, total_comments, buy_box_seller_type, page_inventory, "
"bsr_orders, bsr_orders_sale, title, title_len, price, rating, total_comments, buy_box_seller_type, page_inventory, "
...
@@ -682,7 +688,7 @@ STG_INSERT_COLUMNS = """
...
@@ -682,7 +688,7 @@ STG_INSERT_COLUMNS = """
def
maintain_stg_table
(
site_name
,
months
,
date_info
):
def
maintain_stg_table
(
site_name
,
months
,
date_info
):
"""维护基础详情快照中间表 dwt.{site}_flow_asin_365day:
"""维护基础详情快照中间表 dwt.{site}_flow_asin_365day:
1)
强制重写入参 date_info 当月数据
1)
清理 + 强制重写入参 date_info 当月数据(先 DELETE 再 INSERT,避免残留上次写入的多余 asin)
2) 探测缺失月份,逐月 INSERT INTO 补齐
2) 探测缺失月份,逐月 INSERT INTO 补齐
3) 清理 < 近12月窗口起点的过期数据
3) 清理 < 近12月窗口起点的过期数据
:param months: 近12个月列表,需按 yyyy-MM 升序排列
:param months: 近12个月列表,需按 yyyy-MM 升序排列
...
@@ -706,6 +712,9 @@ def maintain_stg_table(site_name, months, date_info):
...
@@ -706,6 +712,9 @@ def maintain_stg_table(site_name, months, date_info):
WHERE date_info = '{month_yyyymm}'
WHERE date_info = '{month_yyyymm}'
"""
"""
current_month_date
=
f
"{date_info}-01"
print
(
f
"[维护365day] 清理 dwt.{stg_table} 旧数据 WHERE date_info='{current_month_date}',避免残留上次写入的多余 asin"
)
_exec_doris_sql
([
f
"DELETE FROM `dwt`.`{stg_table}` WHERE date_info = '{current_month_date}'"
])
print
(
f
"[维护365day] 强制重写入参当月 {date_info} 数据到 dwt.{stg_table}"
)
print
(
f
"[维护365day] 强制重写入参当月 {date_info} 数据到 dwt.{stg_table}"
)
_exec_doris_sql
([
_build_insert_sql
(
date_info
)])
_exec_doris_sql
([
_build_insert_sql
(
date_info
)])
...
@@ -1578,7 +1587,7 @@ def main(site_name, date_info, result_type='formal'):
...
@@ -1578,7 +1587,7 @@ def main(site_name, date_info, result_type='formal'):
df_save
=
read_and_normalize_month_data
(
spark
,
site_name
,
date_info
)
df_save
=
read_and_normalize_month_data
(
spark
,
site_name
,
date_info
)
# ===== [Step 3] 写入 Doris dwt 主表 =====
# ===== [Step 3] 写入 Doris dwt 主表 =====
write_dwt_month_table
(
df_save
,
doris_table
)
write_dwt_month_table
(
df_save
,
doris_table
,
date_info
)
# ===== [Step 4] 写入 Doris dwt 年表 dwt.{site}_flow_asin_365day(自动聚合+清理过期数据)=====
# ===== [Step 4] 写入 Doris dwt 年表 dwt.{site}_flow_asin_365day(自动聚合+清理过期数据)=====
print
(
f
"[Step 4] 写入 Doris dwt 年表 dwt.{site_name}_flow_asin_365day"
)
print
(
f
"[Step 4] 写入 Doris dwt 年表 dwt.{site_name}_flow_asin_365day"
)
...
...
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