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
964d93bf
Commit
964d93bf
authored
Sep 11, 2026
by
hejiangming
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
fix
parent
0f7e622c
Hide whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
28 additions
and
22 deletions
+28
-22
dwt_bs_top100.py
Pyspark_job/sqoop_export/dwt_bs_top100.py
+14
-11
dwt_bs_top100_change_rate.py
Pyspark_job/sqoop_export/dwt_bs_top100_change_rate.py
+14
-11
No files found.
Pyspark_job/sqoop_export/dwt_bs_top100.py
View file @
964d93bf
...
@@ -166,19 +166,13 @@ if __name__ == '__main__':
...
@@ -166,19 +166,13 @@ if __name__ == '__main__':
# 导出完成状态更新:主表 + 同比环比表两个导出脚本共用同一 belonging_to_process(品类调研流程_{date_type})。
# 导出完成状态更新:主表 + 同比环比表两个导出脚本共用同一 belonging_to_process(品类调研流程_{date_type})。
# modify_export_workflow_status 会把本脚本置完成(status=3)、统计同组未完成数——只有最后一个跑完的
# modify_export_workflow_status 会把本脚本置完成(status=3)、统计同组未完成数——只有最后一个跑完的
# (未完成数=0)才执行下面的 update_workflow_sql,即"两个都导完
,页面才展示
"(机制同 dwt_aba_last_change_rate.py)。
# (未完成数=0)才执行下面的 update_workflow_sql,即"两个都导完
才写
"(机制同 dwt_aba_last_change_rate.py)。
# 两脚本传同一条 update SQL,谁最后完成都能正确触发。本地/test 运行时该函数自动跳过。
# 两脚本传同一条 update SQL,谁最后完成都能正确触发。本地/test 运行时该函数自动跳过。
# 本功能不像 aba 会在流程启动时预写一条 workflow_everyday,故用 replace INTO 直接插入/覆盖该行(不存在则建、存在则替换)。
# 页面只查 workflow_progress(有当月行即放出当月数据),故它必须严格门控——放进 modify_export_workflow_status,
# 只有主表+同比环比两个导出都完成时,modify_export_workflow_status 才会执行这条(即页面此时才出现"完成"行)。
# 由"两个导出都完成"这一刻写。这里只放这一条 REPLACE(单语句):modify 内部走的 engine_exec_sql 是单次
# 同一条 SQL 里再拼一条 REPLACE,往 selection.workflow_progress 写"品类调研 计算完成"进度记录:
# execute,MySQL/pymysql 不支持一次多条语句(会 1064),所以不能再拼第二条。table_name 不带月份;
# 和 workflow_everyday 同一时机(两个导出都完成才触发)、同一次执行、只写 1 条;
# status/status_val 用 workflow_progress 自己口径(计算完成/3);created_at/updated_at 用 MySQL NOW() 取执行时刻。
# table_name/remark 与 workflow_everyday 一致(table_name 不带月份);status/status_val 用 workflow_progress 自己口径(计算完成/3)。
# engine_exec_sql 支持一次多条语句;created_at/updated_at 用 MySQL NOW() 取执行时刻。
update_workflow_sql
=
f
"""
update_workflow_sql
=
f
"""
replace INTO selection.workflow_everyday
(site_name, report_date, status, status_val, table_name, date_type, page, is_end, remark, export_db_type)
VALUES('{site_name}', '{date_info}', '导出PG数据库完成', 14, '{site_name}_category_top_analysis', '{date_type}', '品类调研', '是', '品类调研月表', '{db_type}');
REPLACE INTO selection.workflow_progress
REPLACE INTO selection.workflow_progress
(site_name, page, table_name, date_type, date_info, status, status_val, is_end,
(site_name, page, table_name, date_type, date_info, status, status_val, is_end,
over_date, total, crawling_responsible, calculate_responsible, exhibition_responsible,
over_date, total, crawling_responsible, calculate_responsible, exhibition_responsible,
...
@@ -189,4 +183,13 @@ if __name__ == '__main__':
...
@@ -189,4 +183,13 @@ if __name__ == '__main__':
NOW(), NOW(), 1, 1, NULL, 1, 1, 1);
NOW(), NOW(), 1, 1, NULL, 1, 1, 1);
"""
"""
CommonUtil
.
modify_export_workflow_status
(
update_workflow_sql
,
site_name
,
date_type
,
date_info
)
CommonUtil
.
modify_export_workflow_status
(
update_workflow_sql
,
site_name
,
date_type
,
date_info
)
# workflow_everyday 只是流程记录、页面不读,不需要严格门控:本脚本导出成功即写(REPLACE 幂等,两个脚本各写一次也只覆盖同一行)。
# 单独用 exec_sql 发(内部按 ; 拆分逐条执行),不与上面的 workflow_progress 拼成多语句,避免 pymysql 1064。
if
test_flag
!=
'test'
:
DBUtil
.
exec_sql
(
'mysql'
,
'us'
,
f
"""
REPLACE INTO selection.workflow_everyday
(site_name, report_date, status, status_val, table_name, date_type, page, is_end, remark, export_db_type)
VALUES('{site_name}', '{date_info}', '导出PG数据库完成', 14, '{site_name}_category_top_analysis', '{date_type}', '品类调研', '是', '品类调研月表', '{db_type}');
"""
)
print
(
"success"
)
print
(
"success"
)
Pyspark_job/sqoop_export/dwt_bs_top100_change_rate.py
View file @
964d93bf
...
@@ -130,19 +130,13 @@ if __name__ == '__main__':
...
@@ -130,19 +130,13 @@ if __name__ == '__main__':
# 导出完成状态更新:主表 + 同比环比表两个导出脚本共用同一 belonging_to_process(品类调研流程_{date_type})。
# 导出完成状态更新:主表 + 同比环比表两个导出脚本共用同一 belonging_to_process(品类调研流程_{date_type})。
# modify_export_workflow_status 会把本脚本置完成(status=3)、统计同组未完成数——只有最后一个跑完的
# modify_export_workflow_status 会把本脚本置完成(status=3)、统计同组未完成数——只有最后一个跑完的
# (未完成数=0)才执行下面的 update_workflow_sql,即"两个都导完
,页面才展示
"(机制同 dwt_aba_last_change_rate.py)。
# (未完成数=0)才执行下面的 update_workflow_sql,即"两个都导完
才写
"(机制同 dwt_aba_last_change_rate.py)。
# 两脚本传同一条 update SQL,谁最后完成都能正确触发。本地/test 运行时该函数自动跳过。
# 两脚本传同一条 update SQL,谁最后完成都能正确触发。本地/test 运行时该函数自动跳过。
# 本功能不像 aba 会在流程启动时预写一条 workflow_everyday,故用 replace INTO 直接插入/覆盖该行(不存在则建、存在则替换)。
# 页面只查 workflow_progress(有当月行即放出当月数据),故它必须严格门控——放进 modify_export_workflow_status,
# 只有主表+同比环比两个导出都完成时,modify_export_workflow_status 才会执行这条(即页面此时才出现"完成"行)。
# 由"两个导出都完成"这一刻写。这里只放这一条 REPLACE(单语句):modify 内部走的 engine_exec_sql 是单次
# 同一条 SQL 里再拼一条 REPLACE,往 selection.workflow_progress 写"品类调研 计算完成"进度记录:
# execute,MySQL/pymysql 不支持一次多条语句(会 1064),所以不能再拼第二条。table_name 不带月份;
# 和 workflow_everyday 同一时机(两个导出都完成才触发)、同一次执行、只写 1 条;
# status/status_val 用 workflow_progress 自己口径(计算完成/3);created_at/updated_at 用 MySQL NOW() 取执行时刻。
# table_name/remark 与 workflow_everyday 一致(table_name 不带月份);status/status_val 用 workflow_progress 自己口径(计算完成/3)。
# 两个导出脚本这条 SQL 完全一致,谁最后完成谁触发,结果只 1 条。created_at/updated_at 用 MySQL NOW()。
update_workflow_sql
=
f
"""
update_workflow_sql
=
f
"""
replace INTO selection.workflow_everyday
(site_name, report_date, status, status_val, table_name, date_type, page, is_end, remark, export_db_type)
VALUES('{site_name}', '{date_info}', '导出PG数据库完成', 14, '{site_name}_category_top_analysis', '{date_type}', '品类调研', '是', '品类调研月表', '{db_type}');
REPLACE INTO selection.workflow_progress
REPLACE INTO selection.workflow_progress
(site_name, page, table_name, date_type, date_info, status, status_val, is_end,
(site_name, page, table_name, date_type, date_info, status, status_val, is_end,
over_date, total, crawling_responsible, calculate_responsible, exhibition_responsible,
over_date, total, crawling_responsible, calculate_responsible, exhibition_responsible,
...
@@ -153,4 +147,13 @@ if __name__ == '__main__':
...
@@ -153,4 +147,13 @@ if __name__ == '__main__':
NOW(), NOW(), 1, 1, NULL, 1, 1, 1);
NOW(), NOW(), 1, 1, NULL, 1, 1, 1);
"""
"""
CommonUtil
.
modify_export_workflow_status
(
update_workflow_sql
,
site_name
,
date_type
,
date_info
)
CommonUtil
.
modify_export_workflow_status
(
update_workflow_sql
,
site_name
,
date_type
,
date_info
)
# workflow_everyday 只是流程记录、页面不读,不需要严格门控:本脚本导出成功即写(REPLACE 幂等,两个脚本各写一次也只覆盖同一行)。
# 单独用 exec_sql 发(内部按 ; 拆分逐条执行),不与上面的 workflow_progress 拼成多语句,避免 pymysql 1064。
if
test_flag
!=
'test'
:
DBUtil
.
exec_sql
(
'mysql'
,
'us'
,
f
"""
REPLACE INTO selection.workflow_everyday
(site_name, report_date, status, status_val, table_name, date_type, page, is_end, remark, export_db_type)
VALUES('{site_name}', '{date_info}', '导出PG数据库完成', 14, '{site_name}_category_top_analysis', '{date_type}', '品类调研', '是', '品类调研月表', '{db_type}');
"""
)
print
(
"success"
)
print
(
"success"
)
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