Commit 6bc20288 by chenyuanjie

流量选品关联信息库数据+年度流量选品导出

parent 747d60b6
""" """
author: CT @Author : CT
description: 同步 Hive dwt_flow_asin 月维度数据到 Doris dwt.{site}_flow_asin_month @Description : 月流量选品 Doris 落地脚本(月流程),同时关联产出年度流量选品数据
流程: (原 doris_handle/dwt_flow_asin_year.py 的 Doris 编排部分合并进来,
Step 0) 建 selection 月物化表 selection.{site}_flow_asin_month_{yyyy_mm}[_test] 因为月流程本身就需要关联部分年度聚合指标,拆开两个脚本维护反而麻烦)
Step 1~3) Spark 读 Hive → 规范化 → 写 Doris dwt.{site}_flow_asin_month 流程:
Step 4) Doris INSERT OVERWRITE 物化到 selection 月物化表 [Step 1] Doris 建表 selection.{site}_flow_asin_month_{yyyy_mm}[_test]
Step 5) 更新 MySQL workflow_everyday 流程记录表(仅 formal 模式) [Step 2] 读 Hive dwt_flow_asin 月数据 + 字段规范化(与 dwt.us_flow_asin_30day
支持 us / uk / de 三站点 DDL 对齐),合并为一步
支持 formal / test 模式: [Step 3] 写入 Doris dwt 主表 dwt.{site}_flow_asin_month
- formal:selection 表名无后缀,更新流程记录表 [Step 4] 写入 Doris dwt 年表 dwt.{site}_flow_asin_365day:强制重写当月数据 +
- test :selection 表名加 _test 后缀,不更新流程记录表(dwt 主表不变) 补齐缺失月份,靠 UNIQUE KEY+MOW+sequence_col(date_info) 自动去重只
保留每个 asin 最新一条快照,并删除近12月窗口之外的过期数据
[Step 5] 同步年度聚合指标 dwt.{site}_flow_asin_365day_extra:从 Hive
dwt_flow_asin_year(由 dwt/dwt_flow_asin_year.py 算好写入)对应
site_name+date_info 分区同步过来,TRUNCATE 后全量重算
[Step 6] Doris INSERT OVERWRITE 到 selection 月物化表
selection.{site}_flow_asin_month_{yyyy_mm}[_test]
[Step 7] Doris INSERT OVERWRITE 到 selection 年表
selection.{site}_flow_asin_365day[_test]
[Step 8] 更新 MySQL workflow_everyday 流程记录表(月/年各写一条,仅 formal 模式)
依赖:dwt/dwt_flow_asin_year.py 必须已经跑完同一个 site_name+date_info,
Step 5 才能同步到有效的年度聚合数据
支持 us / uk / de 三站点
支持 formal / test 模式:
- formal:selection 表名无后缀,更新流程记录表
- test :selection 表名加 _test 后缀,不更新流程记录表(dwt 主表不变)
执行示例: 执行示例:
spark-submit dwt_flow_asin_month.py us 2026-05 # 默认 formal spark-submit dwt_flow_asin_month.py us 2026-05 # 默认 formal
spark-submit dwt_flow_asin_month.py us 2026-05 formal spark-submit dwt_flow_asin_month.py us 2026-05 formal
...@@ -25,6 +40,7 @@ from pyspark.sql import functions as F ...@@ -25,6 +40,7 @@ from pyspark.sql import functions as F
from utils.spark_util import SparkUtil from utils.spark_util import SparkUtil
from utils.DorisHelper import DorisHelper from utils.DorisHelper import DorisHelper
from utils.common_util import CommonUtil
from utils.db_util import DBUtil, DbTypes from utils.db_util import DBUtil, DbTypes
...@@ -32,12 +48,11 @@ DORIS_DB = "dwt" ...@@ -32,12 +48,11 @@ DORIS_DB = "dwt"
SUPPORTED_SITES = ("us", "uk", "de") SUPPORTED_SITES = ("us", "uk", "de")
def _exec_doris_sql(sql_list, use_type='selection'): def _doris_connect(use_type='selection'):
"""通过 pymysql 走 Doris jdbc_port 执行 DDL / INSERT 等 """统一 pymysql 连接,database='selection' 让 UDF(如 stem_en)与 selection 库可见"""
指定 database=selection 让 UDF(如 stem_en)可见"""
import pymysql import pymysql
conn_info = DorisHelper.get_connection_info(use_type) conn_info = DorisHelper.get_connection_info(use_type)
conn = pymysql.connect( return pymysql.connect(
host=conn_info['ip'], host=conn_info['ip'],
port=conn_info['jdbc_port'], port=conn_info['jdbc_port'],
user=conn_info['user'], user=conn_info['user'],
...@@ -46,6 +61,11 @@ def _exec_doris_sql(sql_list, use_type='selection'): ...@@ -46,6 +61,11 @@ def _exec_doris_sql(sql_list, use_type='selection'):
charset='utf8mb4', charset='utf8mb4',
autocommit=True, autocommit=True,
) )
def _exec_doris_sql(sql_list, use_type='selection'):
"""通过 pymysql 走 Doris jdbc_port 执行 DDL / INSERT 等"""
conn = _doris_connect(use_type)
try: try:
cur = conn.cursor() cur = conn.cursor()
for sql in sql_list: for sql in sql_list:
...@@ -56,7 +76,24 @@ def _exec_doris_sql(sql_list, use_type='selection'): ...@@ -56,7 +76,24 @@ def _exec_doris_sql(sql_list, use_type='selection'):
conn.close() conn.close()
def build_create_table_sql(table_name, date_info): def _query_doris(sql, use_type='selection'):
"""执行查询并返回所有行"""
conn = _doris_connect(use_type)
try:
cur = conn.cursor()
cur.execute(sql)
rows = cur.fetchall()
cur.close()
return rows
finally:
conn.close()
# ============================================================
# [Step 1] selection 月物化表建表
# ============================================================
def build_month_create_table_sql(table_name, date_info):
"""构建 selection.{table_name} 建表语句(与 DDL 一致); """构建 selection.{table_name} 建表语句(与 DDL 一致);
table_name 由外层拼接:{site}_flow_asin_month_{yyyy_mm}[_test]""" table_name 由外层拼接:{site}_flow_asin_month_{yyyy_mm}[_test]"""
return f""" return f"""
...@@ -198,6 +235,22 @@ CREATE TABLE IF NOT EXISTS `selection`.`{table_name}` ...@@ -198,6 +235,22 @@ CREATE TABLE IF NOT EXISTS `selection`.`{table_name}`
`title_matching_degree` DECIMAL(20,4) NULL, `title_matching_degree` DECIMAL(20,4) NULL,
`brand_badge_reason` STRING NULL, `brand_badge_reason` STRING NULL,
`title_stem_15` STRING NULL COMMENT '标题词干前15个词(按空格分词截取)', `title_stem_15` STRING NULL COMMENT '标题词干前15个词(按空格分词截取)',
`ai_package_quantity` STRING NULL COMMENT 'AI分析-包装数量描述',
`ai_package_quantity_arr` ARRAY<INT> NULL COMMENT 'AI分析-包装数量数组',
`ai_material` STRING NULL COMMENT 'AI分析-材质',
`ai_color` STRING NULL COMMENT 'AI分析-颜色',
`ai_appearance` STRING NULL COMMENT 'AI分析-外观',
`ai_size` STRING NULL COMMENT 'AI分析-尺寸',
`ai_shape` STRING NULL COMMENT 'AI分析-形状',
`ai_function` STRING NULL COMMENT 'AI分析-功能',
`ai_scene_title` STRING NULL COMMENT 'AI分析-使用场景(标题)',
`ai_scene_comment` STRING NULL COMMENT 'AI分析-使用场景(评论)',
`ai_uses` STRING NULL COMMENT 'AI分析-用途',
`ai_theme` STRING NULL COMMENT 'AI分析-主题',
`ai_crowd` STRING NULL COMMENT 'AI分析-目标人群',
`ai_short_desc` STRING NULL COMMENT 'AI分析-简短描述',
`seller_province` STRING NULL COMMENT '卖家所在省份(dim_seller_address关联,仅国内卖家有值)',
`seller_city` STRING NULL COMMENT '卖家所在城市(dim_seller_address关联,仅国内卖家有值)',
INDEX idx_title (`title`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题倒排索引', INDEX idx_title (`title`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题倒排索引',
INDEX idx_title_stem (`title_stem`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题词干倒排索引', INDEX idx_title_stem (`title_stem`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题词干倒排索引',
INDEX idx_title_stem_15 (`title_stem_15`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题词干前15个词倒排索引', INDEX idx_title_stem_15 (`title_stem_15`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题词干前15个词倒排索引',
...@@ -214,239 +267,13 @@ PROPERTIES ( ...@@ -214,239 +267,13 @@ PROPERTIES (
""" """
def build_insert_overwrite_sql(site_name, table_name, date_info): # ============================================================
"""构建 INSERT OVERWRITE 到 selection.{table_name} 的 SQL; # [Step 2] 读 Hive dwt_flow_asin 月数据 + 字段规范化
table_name 由外层拼接:{site}_flow_asin_month_{yyyy_mm}[_test]""" # ============================================================
return f"""
INSERT OVERWRITE TABLE `selection`.`{table_name}`
SELECT
f.asin,
f.parent_asin,
f.collapse_asin,
f.asin_crawl_date,
f.title,
stem_en(f.title) AS title_stem,
f.title_len,
f.brand,
f.asin_describe,
f.describe_len,
f.product_features,
f.together_asin,
f.price,
f.fbm_price,
f.rating,
f.total_comments,
f.one_star, f.two_star, f.three_star, f.four_star, f.five_star, f.low_star,
f.bsr_orders,
f.bsr_orders_sale,
f.asin_bought_month,
f.ao_val,
f.zr_counts,
f.sp_counts, f.sb_counts, f.vi_counts, f.bs_counts,
f.ac_counts, f.tr_counts, f.er_counts,
f.zr_flow_proportion,
f.matrix_flow_proportion,
f.matrix_ao_val,
f.one_two_val, f.three_four_val, f.five_six_val, f.eight_val,
f.category_first_id,
f.category_id,
f.first_category_rank,
f.current_category_rank,
f.weight,
f.volume,
f.asin_weight_ratio,
f.color, f.size, f.style, f.material,
COALESCE(uma.package_quantity, f.package_quantity) AS package_quantity,
f.is_package_quantity_abnormal,
f.variation_num,
f.page_inventory,
f.activity_type,
COALESCE(f.launch_time, kp.keepa_launch_time) AS launch_time,
CASE
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(f.asin_crawl_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(f.asin_crawl_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(f.asin_crawl_date, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 1080 THEN 6
ELSE 7
END AS launch_time_type,
f.img_url,
f.img_num,
f.img_type AS img_type_arr,
f.img_info,
f.account_id,
f.account_name,
f.buy_box_seller_type,
f.site_name,
f.follow_sellers_count,
f.asin_lob_info,
f.is_contains_lob_info,
f.asin_lqs_rating,
f.asin_lqs_rating_detail,
f.amazon_label,
f.is_movie_label,
f.is_brand_label,
f.multi_color_flag,
f.multi_color_str,
f.rank_type, f.ao_val_type, f.price_type, f.rating_type,
f.size_type, f.weight_type, f.site_name_type, f.quantity_variation_type,
f.rank_rise, f.rank_mom AS rank_change, f.rank_yoy,
f.ao_rise, f.ao_mom AS ao_change, f.ao_yoy,
f.price_rise, f.price_mom AS price_change, f.price_yoy,
f.rating_rise, f.rating_mom AS rating_change, f.rating_yoy,
f.comments_rise, f.comments_mom AS comments_change, f.comments_yoy,
f.bsr_orders_rise, f.bsr_orders_mom AS bsr_orders_change, f.bsr_orders_yoy,
f.sales_rise, f.sales_mom AS sales_change, f.sales_yoy,
f.variation_rise, f.variation_mom AS variation_change, f.variation_yoy,
f.bought_month_mom, f.bought_month_yoy,
pr.ocean_profit, pr.air_profit,
FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60) AS tracking_since,
CASE
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(f.asin_crawl_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(f.asin_crawl_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(f.asin_crawl_date, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 1080 THEN 6
ELSE 7
END AS tracking_since_type,
kp.package_length,
kp.package_width,
kp.package_height,
CASE WHEN kp.item_weight > 0 THEN kp.item_weight ELSE kp.package_weight END AS item_weight,
CASE WHEN ba.brand_name_norm IS NOT NULL THEN 1 ELSE 0 END AS is_alarm_brand,
CAST(-1 AS SMALLINT) AS bsr_best_orders_type,
COALESCE(cf.asin_cate_flag, ARRAY(0)) AS asin_source_flag,
COALESCE(cf.bsr_latest_date, CAST('1970-01-01' AS DATE)) AS bsr_last_seen_at,
COALESCE(cf.bsr_30day_count, 0) AS bsr_seen_count_30d,
COALESCE(cf.nsr_latest_date, CAST('1970-01-01' AS DATE)) AS nsr_last_seen_at,
COALESCE(cf.nsr_30day_count, 0) AS nsr_seen_count_30d,
CAST(
CASE
WHEN sa.asin IS NOT NULL OR ab.brand_lower IS NOT NULL THEN 1
WHEN (
(COALESCE(chf.is_need_cat, 0) = 1 AND COALESCE(chd.is_need_cat, 0) = 1)
OR f.asin NOT LIKE 'B0%'
) THEN 2
WHEN (COALESCE(chf.is_hide_cat, 0) = 1 OR COALESCE(chc.is_hide_cat, 0) = 1)
AND NOT (COALESCE(chf.is_white_cat, 0) = 1 OR COALESCE(chc.is_white_cat, 0) = 1) THEN 3
ELSE 0
END
AS SMALLINT) AS asin_type,
uma.usr_mask_progress,
COALESCE(uma.usr_mask_type, umc.usr_mask_type) AS usr_mask_type,
COALESCE(aa.auctions_num, 0) AS auctions_num,
COALESCE(aa.auctions_num_all, 0) AS auctions_num_all,
COALESCE(aa.skus_num_creat, 0) AS skus_num_creat,
COALESCE(aa.skus_num_creat_all, 0) AS skus_num_creat_all,
f.title_matching_degree,
bb.brand_badge_reason,
ARRAY_JOIN(ARRAY_SLICE(SPLIT_BY_STRING(stem_en(f.title), ' '), 1, 15), ' ') AS title_stem_15
FROM `dwt`.`{site_name}_flow_asin_month` f
LEFT JOIN `dwd`.`dwd_asin_profit_rate_latest` pr
ON f.asin = pr.asin AND f.price = pr.price AND pr.site_name = '{site_name}'
LEFT JOIN `dwd`.`dwd_keepa_asin_detail` kp
ON f.asin = kp.asin AND kp.site_name = '{site_name}'
LEFT JOIN (
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand_name_norm
FROM `selection`.`brand_alert_erp`
WHERE brand_name IS NOT NULL
) ba ON f.brand = ba.brand_name_norm
LEFT JOIN (
SELECT asin FROM `mysql_selection`.`selection`.`us_self_asin` GROUP BY asin
) sa ON f.asin = sa.asin
LEFT JOIN (
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand_lower
FROM `mysql_selection`.`selection`.`amazon_brand` WHERE brand_type = '1'
) ab ON LOWER(f.brand) = ab.brand_lower
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chf ON f.category_first_id = chf.category_id_base
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chc ON f.category_id = chc.category_id_base
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chd ON f.desc_category_first_id = chd.category_id_base
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 (
SELECT asin, asin_cate_flag, bsr_latest_date, bsr_30day_count, nsr_latest_date, nsr_30day_count
FROM `dwd`.`dwd_asin_source_flag`
WHERE site_name = '{site_name}' AND date_type = 'month' AND date_info = '{date_info}'
) cf ON f.asin = cf.asin
LEFT JOIN `dwd`.`dwd_asin_auction` aa ON f.asin = aa.asin
LEFT JOIN `dwd`.`dwd_st_brand_badge` bb
ON f.brand = bb.brand AND bb.site_name = '{site_name}' AND bb.date_info = '{date_info}'
WHERE f.date_info = '{date_info}'
"""
def modify_mission_record_status(site_name, date_info, result_type):
"""流程记录表更新:仅 month + formal 模式才入库 mysql workflow_everyday
参考 export_es/es_flow_asin.py modify_mission_record_status"""
if result_type != 'formal':
print(f"[Step 5] result_type={result_type},跳过流程记录表更新")
return
record_table = 'workflow_everyday'
record_table_name_field = f'{site_name}_flow_asin_last_month'
record_type = 'month'
cur_date = date_info
engine_mysql = DBUtil.get_db_engine(db_type=DbTypes.mysql.name, site_name='us')
select_sql = (
f"select id from {record_table} where site_name='{site_name}' and date_type='month' "
f"and report_date='{cur_date}' and page='流量选品' and status_val=14 and is_end='是'"
)
df_is_finished = pd.read_sql(select_sql, engine_mysql)
if df_is_finished.empty:
replace_sql = f"""
replace into {record_table} (site_name, report_date, status, status_val, table_name, date_type, page, is_end, remark, export_db_type)
VALUES ('{site_name}', '{cur_date}', '流量选品计算完毕', 14, '{record_table_name_field}', '{record_type}', '流量选品', '是', '流量选品计算完毕', 'doris')
"""
DBUtil.exec_sql('mysql', 'us', replace_sql)
print(f"[Step 5] 流程记录表 workflow_everyday 已写入:{site_name} {cur_date}")
else:
print(f"[Step 5] 流程记录表已存在该记录,跳过")
def main(site_name, date_info, result_type='formal'):
assert site_name in SUPPORTED_SITES, f"不支持的站点:{site_name},仅支持 us/uk/de"
assert result_type in ('formal', 'test'), f"不支持的 result_type:{result_type},仅支持 formal/test"
doris_table = f"{site_name}_flow_asin_month" # dwt 主表,不区分 test/formal
date_info_underscore = date_info.replace('-', '_')
# selection 月物化表:test 模式加 _test 后缀
env_suffix = '_test' if result_type == 'test' else ''
selection_table = f"{site_name}_flow_asin_month_{date_info_underscore}{env_suffix}"
print(f"启动:site={site_name}, date_info={date_info}, result_type={result_type}")
print(f" dwt 主表:dwt.{doris_table}(不区分 test/formal)")
print(f" selection 物化表:selection.{selection_table}")
spark = SparkUtil.get_spark_session(
f"DwtFlowAsinMonth: {site_name} {date_info} {result_type}"
)
# ===== Step 0:Doris 端建 selection 月物化表(IF NOT EXISTS) =====
print(f"[Step 0] Doris 建表 selection.{selection_table}")
_exec_doris_sql([build_create_table_sql(selection_table, date_info)])
# ===== Step 1:读 Hive dwt_flow_asin 月数据 ===== def read_and_normalize_month_data(spark, site_name, date_info):
"""读 Hive dwt_flow_asin 月数据,规范化字段(与 dwt.us_flow_asin_30day DDL 对齐)
返回 cache 好的 df_save,调用方用完需自行 unpersist(见 write_dwt_month_table)"""
sql = f""" sql = f"""
SELECT SELECT
asin, asin,
...@@ -511,7 +338,6 @@ def main(site_name, date_info, result_type='formal'): ...@@ -511,7 +338,6 @@ def main(site_name, date_info, result_type='formal'):
print(f"sql=\n{sql}") print(f"sql=\n{sql}")
df_raw = spark.sql(sqlQuery=sql).repartition(40, 'asin') df_raw = spark.sql(sqlQuery=sql).repartition(40, 'asin')
# ===== Step 2:字段规范化(与 dwt.us_flow_asin_30day DDL 对齐)=====
def _norm_dt(col_name): def _norm_dt(col_name):
"""text → DATETIME 容错:先 yyyy-MM-dd HH:mm:ss,失败退到 yyyy-MM-dd""" """text → DATETIME 容错:先 yyyy-MM-dd HH:mm:ss,失败退到 yyyy-MM-dd"""
c = F.col(col_name) c = F.col(col_name)
...@@ -647,8 +473,15 @@ def main(site_name, date_info, result_type='formal'): ...@@ -647,8 +473,15 @@ def main(site_name, date_info, result_type='formal'):
count = df_save.count() count = df_save.count()
print(f"写入数据量:{count:,}") print(f"写入数据量:{count:,}")
df_save.show(5, truncate=False) df_save.show(5, truncate=False)
return df_save
# ===== Step 3:写入 Doris dwt 主表 ===== # ============================================================
# [Step 3] 写入 Doris dwt 主表
# ============================================================
def write_dwt_month_table(df_save, doris_table):
"""写入 Doris dwt 主表 dwt.{site}_flow_asin_month,写完 unpersist 传入的 df_save"""
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, "
...@@ -677,11 +510,1082 @@ def main(site_name, date_info, result_type='formal'): ...@@ -677,11 +510,1082 @@ def main(site_name, date_info, result_type='formal'):
) )
df_save.unpersist() df_save.unpersist()
# ===== Step 4:Doris 端 INSERT OVERWRITE 到 selection 月物化表 =====
print(f"[Step 4] Doris INSERT OVERWRITE selection.{selection_table}")
_exec_doris_sql([build_insert_overwrite_sql(site_name, selection_table, date_info)])
# ===== Step 5:流程记录表更新(仅 formal 模式)===== # ============================================================
# [Step 4] 维护 dwt.{site}_flow_asin_365day(近12个月ASIN最新快照)
# ============================================================
def build_stg_create_table_sql(site_name):
"""构建基础详情快照中间表 dwt.{site}_flow_asin_365day 建表语句(与 doris_ddl.sql 一致)"""
return f"""
CREATE TABLE IF NOT EXISTS `dwt`.`{site_name}_flow_asin_365day`
(
`asin` VARCHAR(20) NOT NULL,
`date_info` DATE NOT NULL,
`ao_val` DECIMAL(20,4),
`zr_counts` INT,
`sp_counts` INT,
`sb_counts` INT,
`vi_counts` INT,
`bs_counts` INT,
`ac_counts` INT,
`tr_counts` INT,
`er_counts` INT,
`bsr_orders` INT,
`bsr_orders_sale` DECIMAL(20,2),
`title` STRING,
`title_len` INT,
`price` DECIMAL(20,2),
`rating` DECIMAL(10,1),
`total_comments` INT,
`buy_box_seller_type` TINYINT,
`page_inventory` INT,
`volume` STRING,
`weight` DECIMAL(20,4),
`color` STRING,
`size` STRING,
`style` STRING,
`material` STRING,
`launch_time` DATETIME,
`img_num` INT,
`parent_asin` STRING,
`img_type` ARRAY<INT>,
`img_url` STRING,
`activity_type` STRING,
`one_two_val` DECIMAL(20,4) DEFAULT "0.0000",
`three_four_val` DECIMAL(20,4) DEFAULT "0.0000",
`five_six_val` DECIMAL(20,4) DEFAULT "0.0000",
`eight_val` DECIMAL(20,4) DEFAULT "0.0000",
`brand` STRING,
`variation_num` INT,
`one_star` INT,
`two_star` INT,
`three_star` INT,
`four_star` INT,
`five_star` INT,
`low_star` INT,
`together_asin` STRING,
`account_name` STRING,
`account_id` STRING,
`rank_rise` INT,
`rank_mom` DECIMAL(20,4),
`rank_yoy` DECIMAL(20,4),
`ao_rise` DECIMAL(20,4),
`ao_mom` DECIMAL(20,4),
`ao_yoy` DECIMAL(20,4),
`price_rise` DECIMAL(20,2),
`price_mom` DECIMAL(20,4),
`price_yoy` DECIMAL(20,4),
`rating_rise` DECIMAL(20,1),
`rating_mom` DECIMAL(20,4),
`rating_yoy` DECIMAL(20,4),
`comments_rise` INT,
`comments_mom` DECIMAL(20,4),
`comments_yoy` DECIMAL(20,4),
`bsr_orders_rise` INT,
`bsr_orders_mom` DECIMAL(20,4),
`bsr_orders_yoy` DECIMAL(20,4),
`sales_rise` DECIMAL(20,2),
`sales_mom` DECIMAL(20,4),
`sales_yoy` DECIMAL(20,4),
`variation_rise` INT,
`variation_mom` DECIMAL(20,4),
`variation_yoy` DECIMAL(20,4),
`bought_month_mom` DECIMAL(20,4),
`bought_month_yoy` DECIMAL(20,4),
`size_type` TINYINT,
`rating_type` TINYINT,
`site_name_type` TINYINT,
`weight_type` TINYINT,
`ao_val_type` TINYINT,
`rank_type` TINYINT,
`price_type` TINYINT,
`quantity_variation_type` TINYINT DEFAULT "0",
`package_quantity` INT DEFAULT "1",
`is_movie_label` TINYINT DEFAULT "0",
`is_brand_label` TINYINT DEFAULT "0",
`asin_crawl_date` DATETIME,
`category_first_id` STRING,
`category_id` STRING,
`desc_category_first_id` STRING,
`first_category_rank` INT,
`current_category_rank` INT,
`asin_weight_ratio` DECIMAL(20,4),
`site_name` STRING,
`asin_bought_month` INT,
`asin_lqs_rating` DECIMAL(20,1),
`asin_lqs_rating_detail` STRING,
`asin_lob_info` STRING,
`is_contains_lob_info` TINYINT,
`is_package_quantity_abnormal` TINYINT,
`zr_flow_proportion` DECIMAL(20,4),
`matrix_flow_proportion` DECIMAL(20,4),
`matrix_ao_val` DECIMAL(20,4),
`product_features` STRING,
`img_info` STRING,
`collapse_asin` STRING,
`follow_sellers_count` INT DEFAULT "-1",
`asin_describe` STRING,
`fbm_price` DECIMAL(20,2),
`describe_len` INT DEFAULT "0",
`title_matching_degree` DECIMAL(20,4),
`multi_color_flag` TINYINT DEFAULT "0",
`multi_color_str` STRING,
`amazon_label` STRING
)
ENGINE=OLAP
UNIQUE KEY(`asin`)
DISTRIBUTED BY HASH(`asin`) BUCKETS 32
PROPERTIES (
"replication_num" = "3",
"enable_unique_key_merge_on_write" = "true",
"function_column.sequence_col" = "date_info"
)
"""
# 中间表字段列表(不含 date_info,因为 INSERT INTO 时按 SELECT * FROM dwt.month 顺序写入,date_info 在 SELECT 列表末尾)
STG_INSERT_COLUMNS = """
asin, date_info, 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, volume, weight, color, size, style, material, launch_time,
img_num, parent_asin, img_type, img_url, activity_type,
one_two_val, three_four_val, five_six_val, eight_val, brand, variation_num,
one_star, two_star, three_star, four_star, five_star, low_star,
together_asin, account_name, account_id,
rank_rise, rank_mom, rank_yoy, ao_rise, ao_mom, ao_yoy,
price_rise, price_mom, price_yoy, rating_rise, rating_mom, rating_yoy,
comments_rise, comments_mom, comments_yoy,
bsr_orders_rise, bsr_orders_mom, bsr_orders_yoy,
sales_rise, sales_mom, sales_yoy,
variation_rise, variation_mom, variation_yoy,
bought_month_mom, bought_month_yoy,
size_type, rating_type, site_name_type, weight_type, ao_val_type, rank_type, price_type,
quantity_variation_type, package_quantity, is_movie_label, is_brand_label,
asin_crawl_date, category_first_id, category_id, desc_category_first_id,
first_category_rank, current_category_rank, asin_weight_ratio, site_name,
asin_bought_month, asin_lqs_rating, asin_lqs_rating_detail,
asin_lob_info, is_contains_lob_info, is_package_quantity_abnormal,
zr_flow_proportion, matrix_flow_proportion, matrix_ao_val,
product_features, img_info, collapse_asin, follow_sellers_count,
asin_describe, fbm_price, describe_len, title_matching_degree,
multi_color_flag, multi_color_str, amazon_label
""".strip()
def maintain_stg_table(site_name, months, date_info):
"""维护基础详情快照中间表 dwt.{site}_flow_asin_365day:
1) 强制重写入参 date_info 当月数据
2) 探测缺失月份,逐月 INSERT INTO 补齐
3) 清理 < 近12月窗口起点的过期数据
:param months: 近12个月列表,需按 yyyy-MM 升序排列
"""
stg_table = f"{site_name}_flow_asin_365day"
# 注意:中间表 date_info 为 DATE 类型(月首日,如 '2026-05-01'),
# 而 dwt 月表 date_info 为 VARCHAR(yyyy-MM,如 '2026-05')。
# 写入时需用 CAST(CONCAT(date_info, '-01') AS DATE) 把月表的 VARCHAR 转为 DATE。
stg_select_columns = STG_INSERT_COLUMNS.replace(
'date_info',
"CAST(CONCAT(date_info, '-01') AS DATE)",
1 # 只替换 SELECT 列表中第一次出现的 date_info(列名本身)
)
def _build_insert_sql(month_yyyymm):
return f"""
INSERT INTO `dwt`.`{stg_table}` ({STG_INSERT_COLUMNS})
SELECT {stg_select_columns}
FROM `dwt`.`{site_name}_flow_asin_month`
WHERE date_info = '{month_yyyymm}'
"""
print(f"[维护365day] 强制重写入参当月 {date_info} 数据到 dwt.{stg_table}")
_exec_doris_sql([_build_insert_sql(date_info)])
existing_rows = _query_doris(
f"SELECT DISTINCT DATE_FORMAT(date_info, '%Y-%m') FROM `dwt`.`{stg_table}`"
)
existing_months = {row[0] for row in existing_rows}
print(f"[维护365day] 中间表已有月份: {sorted(existing_months)}")
missing_months = [m for m in months if m not in existing_months]
if missing_months:
print(f"[维护365day] 待补齐月份: {missing_months}")
for m in missing_months:
print(f" → INSERT INTO dwt.{stg_table} WHERE date_info='{m}'")
_exec_doris_sql([_build_insert_sql(m)])
else:
print(f"[维护365day] 近12月数据已齐全,无需补齐")
min_keep_month = months[0]
min_keep_date = f"{min_keep_month}-01"
expired_months = [m for m in existing_months if m < min_keep_month]
if expired_months:
print(f"[维护365day] 清理过期数据(date_info < {min_keep_date}): {sorted(expired_months)}")
_exec_doris_sql([f"DELETE FROM `dwt`.`{stg_table}` WHERE date_info < '{min_keep_date}'"])
else:
print(f"[维护365day] 无过期数据需要清理")
# ============================================================
# [Step 5] 同步 Hive dwt_flow_asin_year → Doris dwt.{site}_flow_asin_365day_extra
# ============================================================
EXTRA_COLUMNS = (
"asin, bought_month_total, "
"bought_month_1, bought_month_2, bought_month_3, bought_month_4, "
"bought_month_5, bought_month_6, bought_month_7, bought_month_8, "
"bought_month_9, bought_month_10, bought_month_11, bought_month_12, "
"bought_month_q1, bought_month_q2, bought_month_q3, bought_month_q4, "
"total_appear_month, bought_month_peak, peak_month_arr, "
"is_periodic_flag, is_seasonal_flag, bsr_seen_count_total, nsr_seen_count_total"
)
def build_extra_create_table_sql(table_name):
"""构建年度聚合中间表 dwt.{site}_flow_asin_365day_extra 建表语句(每次运行 TRUNCATE 后全量重算,无需 sequence_col)"""
return f"""
CREATE TABLE IF NOT EXISTS `dwt`.`{table_name}`
(
`asin` VARCHAR(20) NOT NULL,
`bought_month_total` INT,
`bought_month_1` INT,
`bought_month_2` INT,
`bought_month_3` INT,
`bought_month_4` INT,
`bought_month_5` INT,
`bought_month_6` INT,
`bought_month_7` INT,
`bought_month_8` INT,
`bought_month_9` INT,
`bought_month_10` INT,
`bought_month_11` INT,
`bought_month_12` INT,
`bought_month_q1` INT,
`bought_month_q2` INT,
`bought_month_q3` INT,
`bought_month_q4` INT,
`total_appear_month` ARRAY<INT>,
`bought_month_peak` INT,
`peak_month_arr` ARRAY<INT>,
`is_periodic_flag` INT,
`is_seasonal_flag` INT,
`bsr_seen_count_total` INT,
`nsr_seen_count_total` INT
) ENGINE=OLAP
UNIQUE KEY(`asin`)
COMMENT '流量选品年度聚合指标'
DISTRIBUTED BY HASH(`asin`) BUCKETS 32
PROPERTIES (
"replication_num" = "3",
"enable_unique_key_merge_on_write" = "true"
)
"""
def sync_extra_table(spark, site_name, date_info):
"""从 Hive dwt_flow_asin_year 对应 site_name+date_info 分区读取年度聚合指标,
同步到 Doris dwt.{site}_flow_asin_365day_extra(TRUNCATE 后全量重算)
依赖:dwt/dwt_flow_asin_year.py 必须已经算好并写入该分区"""
extra_tb = f"{site_name}_flow_asin_365day_extra"
# total_appear_month / peak_month_arr 在 Hive 那边存的是逗号拼接的 STRING
# (dwt/dwt_flow_asin_year.py 里的说明),这里读出来转回 ARRAY<INT> 给 Doris 用
sql = f"""
SELECT
asin, bought_month_total,
bought_month_1, bought_month_2, bought_month_3, bought_month_4,
bought_month_5, bought_month_6, bought_month_7, bought_month_8,
bought_month_9, bought_month_10, bought_month_11, bought_month_12,
bought_month_q1, bought_month_q2, bought_month_q3, bought_month_q4,
CAST(SPLIT(total_appear_month, ',') AS ARRAY<INT>) AS total_appear_month,
bought_month_peak,
CAST(SPLIT(peak_month_arr, ',') AS ARRAY<INT>) AS peak_month_arr,
is_periodic_flag, is_seasonal_flag, bsr_seen_count_total, nsr_seen_count_total
FROM dwt_flow_asin_year
WHERE site_name = '{site_name}' AND date_info = '{date_info}'
"""
df_extra = spark.sql(sqlQuery=sql).repartition(20, 'asin').cache()
row_count = df_extra.count()
print(f"[同步365day_extra] 读取 Hive dwt_flow_asin_year[site_name={site_name}, date_info={date_info}]:"
f"{row_count} 条")
if row_count == 0:
# 读到0条大概率是 dwt/dwt_flow_asin_year.py 还没跑这个 site_name+date_info,
# 先报错,不要往下 TRUNCATE,避免把上个月还有效的数据清空
raise ValueError(
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"
)
_exec_doris_sql([build_extra_create_table_sql(extra_tb)])
_exec_doris_sql([f"TRUNCATE TABLE `dwt`.`{extra_tb}`"])
DorisHelper.spark_export_with_columns(
df_save=df_extra,
db_name='dwt',
table_name=extra_tb,
table_columns=EXTRA_COLUMNS,
)
print(f"[同步365day_extra] 年度聚合指标同步至 dwt.{extra_tb} 完成")
# ============================================================
# [Step 6] selection 月表 INSERT OVERWRITE
# ============================================================
def build_month_insert_overwrite_sql(site_name, table_name, date_info):
"""构建 INSERT OVERWRITE 到 selection.{table_name} 的 SQL;
table_name 由外层拼接:{site}_flow_asin_month_{yyyy_mm}[_test]"""
return f"""
INSERT OVERWRITE TABLE `selection`.`{table_name}`
SELECT
f.asin,
f.parent_asin,
f.collapse_asin,
f.asin_crawl_date,
f.title,
stem_en(f.title) AS title_stem,
f.title_len,
f.brand,
f.asin_describe,
f.describe_len,
f.product_features,
f.together_asin,
f.price,
f.fbm_price,
f.rating,
f.total_comments,
f.one_star, f.two_star, f.three_star, f.four_star, f.five_star, f.low_star,
f.bsr_orders,
f.bsr_orders_sale,
f.asin_bought_month,
f.ao_val,
f.zr_counts,
f.sp_counts, f.sb_counts, f.vi_counts, f.bs_counts,
f.ac_counts, f.tr_counts, f.er_counts,
f.zr_flow_proportion,
f.matrix_flow_proportion,
f.matrix_ao_val,
f.one_two_val, f.three_four_val, f.five_six_val, f.eight_val,
f.category_first_id,
f.category_id,
f.first_category_rank,
f.current_category_rank,
f.weight,
f.volume,
f.asin_weight_ratio,
f.color, f.size, f.style, f.material,
COALESCE(uma.package_quantity, f.package_quantity) AS package_quantity,
f.is_package_quantity_abnormal,
f.variation_num,
f.page_inventory,
f.activity_type,
COALESCE(f.launch_time, kp.keepa_launch_time) AS launch_time,
CASE
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(f.asin_crawl_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(f.asin_crawl_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(f.asin_crawl_date, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 1080 THEN 6
ELSE 7
END AS launch_time_type,
f.img_url,
f.img_num,
f.img_type AS img_type_arr,
f.img_info,
f.account_id,
f.account_name,
f.buy_box_seller_type,
f.site_name,
f.follow_sellers_count,
f.asin_lob_info,
f.is_contains_lob_info,
f.asin_lqs_rating,
f.asin_lqs_rating_detail,
f.amazon_label,
f.is_movie_label,
f.is_brand_label,
f.multi_color_flag,
f.multi_color_str,
f.rank_type, f.ao_val_type, f.price_type, f.rating_type,
f.size_type, f.weight_type, f.site_name_type, f.quantity_variation_type,
f.rank_rise, f.rank_mom AS rank_change, f.rank_yoy,
f.ao_rise, f.ao_mom AS ao_change, f.ao_yoy,
f.price_rise, f.price_mom AS price_change, f.price_yoy,
f.rating_rise, f.rating_mom AS rating_change, f.rating_yoy,
f.comments_rise, f.comments_mom AS comments_change, f.comments_yoy,
f.bsr_orders_rise, f.bsr_orders_mom AS bsr_orders_change, f.bsr_orders_yoy,
f.sales_rise, f.sales_mom AS sales_change, f.sales_yoy,
f.variation_rise, f.variation_mom AS variation_change, f.variation_yoy,
f.bought_month_mom, f.bought_month_yoy,
pr.ocean_profit, pr.air_profit,
FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60) AS tracking_since,
CASE
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(f.asin_crawl_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(f.asin_crawl_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(f.asin_crawl_date, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 1080 THEN 6
ELSE 7
END AS tracking_since_type,
kp.package_length,
kp.package_width,
kp.package_height,
CASE WHEN kp.item_weight > 0 THEN kp.item_weight ELSE kp.package_weight END AS item_weight,
CASE WHEN ba.brand_name_norm IS NOT NULL THEN 1 ELSE 0 END AS is_alarm_brand,
CAST(-1 AS SMALLINT) AS bsr_best_orders_type,
COALESCE(cf.asin_cate_flag, ARRAY(0)) AS asin_source_flag,
COALESCE(cf.bsr_latest_date, CAST('1970-01-01' AS DATE)) AS bsr_last_seen_at,
COALESCE(cf.bsr_30day_count, 0) AS bsr_seen_count_30d,
COALESCE(cf.nsr_latest_date, CAST('1970-01-01' AS DATE)) AS nsr_last_seen_at,
COALESCE(cf.nsr_30day_count, 0) AS nsr_seen_count_30d,
CAST(
CASE
WHEN sa.asin IS NOT NULL OR ab.brand_lower IS NOT NULL THEN 1
WHEN (
(COALESCE(chf.is_need_cat, 0) = 1 AND COALESCE(chd.is_need_cat, 0) = 1)
OR f.asin NOT LIKE 'B0%'
) THEN 2
WHEN (COALESCE(chf.is_hide_cat, 0) = 1 OR COALESCE(chc.is_hide_cat, 0) = 1)
AND NOT (COALESCE(chf.is_white_cat, 0) = 1 OR COALESCE(chc.is_white_cat, 0) = 1) THEN 3
ELSE 0
END
AS SMALLINT) AS asin_type,
uma.usr_mask_progress,
COALESCE(uma.usr_mask_type, umc.usr_mask_type) AS usr_mask_type,
COALESCE(aa.auctions_num, 0) AS auctions_num,
COALESCE(aa.auctions_num_all, 0) AS auctions_num_all,
COALESCE(aa.skus_num_creat, 0) AS skus_num_creat,
COALESCE(aa.skus_num_creat_all, 0) AS skus_num_creat_all,
f.title_matching_degree,
bb.brand_badge_reason,
ARRAY_JOIN(ARRAY_SLICE(SPLIT_BY_STRING(stem_en(f.title), ' '), 1, 15), ' ') AS title_stem_15,
ai.ai_package_quantity,
ai.ai_package_quantity_arr,
ai.ai_material,
ai.ai_color,
ai.ai_appearance,
ai.ai_size,
ai.ai_shape,
ai.ai_function,
ai.ai_scene_title,
ai.ai_scene_comment,
ai.ai_uses,
ai.ai_theme,
ai.ai_crowd,
ai.ai_short_desc,
addr.seller_province,
addr.seller_city
FROM `dwt`.`{site_name}_flow_asin_month` f
LEFT JOIN `dwd`.`dwd_asin_profit_rate_latest` pr
ON f.asin = pr.asin AND f.price = pr.price AND pr.site_name = '{site_name}'
LEFT JOIN `dwd`.`dwd_keepa_asin_detail` kp
ON f.asin = kp.asin AND kp.site_name = '{site_name}'
LEFT JOIN (
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand_name_norm
FROM `selection`.`brand_alert_erp`
WHERE brand_name IS NOT NULL
) ba ON f.brand = ba.brand_name_norm
LEFT JOIN (
SELECT asin FROM `mysql_selection`.`selection`.`us_self_asin` GROUP BY asin
) sa ON f.asin = sa.asin
LEFT JOIN (
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand_lower
FROM `mysql_selection`.`selection`.`amazon_brand` WHERE brand_type = '1'
) ab ON LOWER(f.brand) = ab.brand_lower
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chf ON f.category_first_id = chf.category_id_base
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chc ON f.category_id = chc.category_id_base
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chd ON f.desc_category_first_id = chd.category_id_base
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 (
SELECT asin, asin_cate_flag, bsr_latest_date, bsr_30day_count, nsr_latest_date, nsr_30day_count
FROM `dwd`.`dwd_asin_source_flag`
WHERE site_name = '{site_name}' AND date_type = 'month' AND date_info = '{date_info}'
) cf ON f.asin = cf.asin
LEFT JOIN `dwd`.`dwd_asin_auction` aa ON f.asin = aa.asin
LEFT JOIN `dwd`.`dwd_st_brand_badge` bb
ON f.brand = bb.brand AND bb.site_name = '{site_name}' AND bb.date_info = '{date_info}'
-- ===== AI分析数据 =====
LEFT JOIN (
SELECT
asin,
package_quantity AS ai_package_quantity,
package_quantity_arr AS ai_package_quantity_arr,
material AS ai_material,
color AS ai_color,
appearance AS ai_appearance,
size AS ai_size,
shape AS ai_shape,
function AS ai_function,
scene_title AS ai_scene_title,
scene_comment AS ai_scene_comment,
uses AS ai_uses,
theme AS ai_theme,
crowd AS ai_crowd,
short_desc AS ai_short_desc
FROM `selection`.`{site_name}_ai_asin_analyze_detail`
) ai ON f.asin = ai.asin
-- ===== 卖家地址(按 site_name 查对应分区) =====
LEFT JOIN (
SELECT seller_id, fb_country_name, seller_province, seller_city
FROM `dim`.`dim_seller_address`
WHERE site_name = '{site_name}'
) addr ON f.account_id = addr.seller_id AND f.site_name = addr.fb_country_name
WHERE f.date_info = '{date_info}'
"""
# ============================================================
# [Step 7] selection 年表建表 + INSERT OVERWRITE
# ============================================================
def build_year_create_table_sql(table_name):
"""构建 selection.{table_name} 建表语句:以 selection 月表(build_month_create_table_sql)
为基准调整而来——
剔除月度专属字段:asin_bought_month / asin_source_flag / bsr_last_seen_at /
bsr_seen_count_30d / nsr_last_seen_at / nsr_seen_count_30d
末尾追加:latest_date_info + dwt.{site}_flow_asin_365day_extra 里的年度聚合字段
旧表已手动清理,直接 CREATE 新表,不再需要 ALTER 兼容旧表结构"""
return f"""
CREATE TABLE IF NOT EXISTS `selection`.`{table_name}`
(
`asin` VARCHAR(20) NOT NULL,
`parent_asin` STRING NULL,
`collapse_asin` STRING NULL,
`asin_crawl_date` DATETIME NULL,
`title` STRING NULL,
`title_stem` STRING NULL,
`title_len` INT NULL,
`brand` STRING NULL,
`asin_describe` STRING NULL,
`describe_len` INT NULL,
`product_features` STRING NULL,
`together_asin` STRING NULL,
`price` DECIMAL(20,2) NULL,
`fbm_price` DECIMAL(20,2) NULL,
`rating` DECIMAL(10,1) NULL,
`total_comments` INT NULL,
`one_star` INT NULL,
`two_star` INT NULL,
`three_star` INT NULL,
`four_star` INT NULL,
`five_star` INT NULL,
`low_star` INT NULL,
`bsr_orders` INT NULL,
`bsr_orders_sale` DECIMAL(20,2) NULL,
`ao_val` DECIMAL(20,4) NULL,
`zr_counts` INT NULL,
`sp_counts` INT NULL,
`sb_counts` INT NULL,
`vi_counts` INT NULL,
`bs_counts` INT NULL,
`ac_counts` INT NULL,
`tr_counts` INT NULL,
`er_counts` INT NULL,
`zr_flow_proportion` DECIMAL(20,4) NULL,
`matrix_flow_proportion` DECIMAL(20,4) NULL,
`matrix_ao_val` DECIMAL(20,4) NULL,
`one_two_val` DECIMAL(20,4) NULL,
`three_four_val` DECIMAL(20,4) NULL,
`five_six_val` DECIMAL(20,4) NULL,
`eight_val` DECIMAL(20,4) NULL,
`category_first_id` STRING NULL,
`category_id` STRING NULL,
`first_category_rank` INT NULL,
`current_category_rank` INT NULL,
`weight` DECIMAL(20,4) NULL,
`volume` STRING NULL,
`asin_weight_ratio` DECIMAL(20,4) NULL,
`color` STRING NULL,
`size` STRING NULL,
`style` STRING NULL,
`material` STRING NULL,
`package_quantity` INT NULL,
`is_package_quantity_abnormal` TINYINT NULL,
`variation_num` INT NULL,
`page_inventory` INT NULL,
`activity_type` STRING NULL,
`launch_time` DATETIME NULL,
`launch_time_type` INT NULL,
`img_url` STRING NULL,
`img_num` INT NULL,
`img_type_arr` ARRAY<INT> NULL,
`img_info` STRING NULL,
`account_id` STRING NULL,
`account_name` STRING NULL,
`buy_box_seller_type` TINYINT NULL,
`site_name` STRING NULL,
`follow_sellers_count` INT NULL,
`asin_lob_info` STRING NULL,
`is_contains_lob_info` TINYINT NULL,
`asin_lqs_rating` DECIMAL(20,1) NULL,
`asin_lqs_rating_detail` STRING NULL,
`amazon_label` STRING NULL,
`is_movie_label` TINYINT NULL,
`is_brand_label` TINYINT NULL,
`multi_color_flag` TINYINT NULL,
`multi_color_str` STRING NULL,
`rank_type` TINYINT NULL,
`ao_val_type` TINYINT NULL,
`price_type` TINYINT NULL,
`rating_type` TINYINT NULL,
`size_type` TINYINT NULL,
`weight_type` TINYINT NULL,
`site_name_type` TINYINT NULL,
`quantity_variation_type` TINYINT NULL,
`rank_rise` INT NULL,
`rank_change` DECIMAL(20,4) NULL,
`rank_yoy` DECIMAL(20,4) NULL,
`ao_rise` DECIMAL(20,4) NULL,
`ao_change` DECIMAL(20,4) NULL,
`ao_yoy` DECIMAL(20,4) NULL,
`price_rise` DECIMAL(20,2) NULL,
`price_change` DECIMAL(20,4) NULL,
`price_yoy` DECIMAL(20,4) NULL,
`rating_rise` DECIMAL(20,1) NULL,
`rating_change` DECIMAL(20,4) NULL,
`rating_yoy` DECIMAL(20,4) NULL,
`comments_rise` INT NULL,
`comments_change` DECIMAL(20,4) NULL,
`comments_yoy` DECIMAL(20,4) NULL,
`bsr_orders_rise` INT NULL,
`bsr_orders_change` DECIMAL(20,4) NULL,
`bsr_orders_yoy` DECIMAL(20,4) NULL,
`sales_rise` DECIMAL(20,2) NULL,
`sales_change` DECIMAL(20,4) NULL,
`sales_yoy` DECIMAL(20,4) NULL,
`variation_rise` INT NULL,
`variation_change` DECIMAL(20,4) NULL,
`variation_yoy` DECIMAL(20,4) NULL,
`bought_month_mom` DECIMAL(20,4) NULL,
`bought_month_yoy` DECIMAL(20,4) NULL,
`ocean_profit` DECIMAL(20,4) NULL,
`air_profit` DECIMAL(20,4) NULL,
`tracking_since` STRING NULL,
`tracking_since_type` INT NULL,
`package_length` INT NULL,
`package_width` INT NULL,
`package_height` INT NULL,
`item_weight` INT NULL,
`is_alarm_brand` INT NULL,
`bsr_best_orders_type` SMALLINT NULL,
`asin_type` SMALLINT NULL,
`usr_mask_progress` STRING NULL,
`usr_mask_type` STRING NULL,
`auctions_num` INT NULL,
`auctions_num_all` INT NULL,
`skus_num_creat` INT NULL,
`skus_num_creat_all` INT NULL,
`title_matching_degree` DECIMAL(20,4) NULL,
`brand_badge_reason` STRING NULL,
`title_stem_15` STRING NULL COMMENT '标题词干前15个词(按空格分词截取)',
`ai_package_quantity` STRING NULL COMMENT 'AI分析-包装数量描述',
`ai_package_quantity_arr` ARRAY<INT> NULL COMMENT 'AI分析-包装数量数组',
`ai_material` STRING NULL COMMENT 'AI分析-材质',
`ai_color` STRING NULL COMMENT 'AI分析-颜色',
`ai_appearance` STRING NULL COMMENT 'AI分析-外观',
`ai_size` STRING NULL COMMENT 'AI分析-尺寸',
`ai_shape` STRING NULL COMMENT 'AI分析-形状',
`ai_function` STRING NULL COMMENT 'AI分析-功能',
`ai_scene_title` STRING NULL COMMENT 'AI分析-使用场景(标题)',
`ai_scene_comment` STRING NULL COMMENT 'AI分析-使用场景(评论)',
`ai_uses` STRING NULL COMMENT 'AI分析-用途',
`ai_theme` STRING NULL COMMENT 'AI分析-主题',
`ai_crowd` STRING NULL COMMENT 'AI分析-目标人群',
`ai_short_desc` STRING NULL COMMENT 'AI分析-简短描述',
`seller_province` STRING NULL COMMENT '卖家所在省份(dim_seller_address关联,仅国内卖家有值)',
`seller_city` STRING NULL COMMENT '卖家所在城市(dim_seller_address关联,仅国内卖家有值)',
`latest_date_info` STRING NULL COMMENT '最新出现月份(yyyy-MM)',
`bought_month_total` INT NULL,
`bought_month_1` INT NULL,
`bought_month_2` INT NULL,
`bought_month_3` INT NULL,
`bought_month_4` INT NULL,
`bought_month_5` INT NULL,
`bought_month_6` INT NULL,
`bought_month_7` INT NULL,
`bought_month_8` INT NULL,
`bought_month_9` INT NULL,
`bought_month_10` INT NULL,
`bought_month_11` INT NULL,
`bought_month_12` INT NULL,
`bought_month_q1` INT NULL,
`bought_month_q2` INT NULL,
`bought_month_q3` INT NULL,
`bought_month_q4` INT NULL,
`total_appear_month` ARRAY<INT> NULL,
`bsr_seen_count_total` INT NULL,
`nsr_seen_count_total` INT NULL,
`bought_month_peak` INT NULL,
`peak_month_arr` ARRAY<INT> NULL,
`is_periodic_flag` INT NULL,
`is_seasonal_flag` INT NULL,
INDEX idx_title (`title`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题倒排索引',
INDEX idx_title_stem (`title_stem`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题词干倒排索引',
INDEX idx_title_stem_15 (`title_stem_15`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '标题词干前15个词倒排索引',
INDEX idx_color (`color`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '颜色倒排索引',
INDEX idx_brand (`brand`) USING INVERTED PROPERTIES("parser" = "english") COMMENT '品牌倒排索引'
) ENGINE=OLAP
UNIQUE KEY(`asin`)
COMMENT '流量选品近 12 个月聚合表'
DISTRIBUTED BY HASH(`asin`) BUCKETS 32
PROPERTIES (
"replication_num" = "3",
"enable_unique_key_merge_on_write" = "true"
)
"""
def build_year_insert_overwrite_sql(site_name, table_name):
"""构造 INSERT OVERWRITE SQL
- 主体: dwt.{site}_flow_asin_365day (每 asin 最新月快照,字段名沿用 dwt 月表)
- 年度聚合: dwt.{site}_flow_asin_365day_extra (Hive 算好经 sync_extra_table 同步来的年度指标,
含周期性/季节性/峰值)
- 外层 LEFT JOIN: profit_rate / keepa / brand_alert / self_asin / category_hide / user_mask / auction 等
"""
return f"""
INSERT OVERWRITE TABLE `selection`.`{table_name}`
SELECT
f.asin,
f.parent_asin,
f.collapse_asin,
f.asin_crawl_date,
f.title,
stem_en(f.title) AS title_stem,
f.title_len,
f.brand,
f.asin_describe,
f.describe_len,
f.product_features,
f.together_asin,
f.price,
f.fbm_price,
f.rating,
f.total_comments,
f.one_star, f.two_star, f.three_star, f.four_star, f.five_star, f.low_star,
f.bsr_orders,
f.bsr_orders_sale,
f.ao_val,
f.zr_counts,
f.sp_counts, f.sb_counts, f.vi_counts, f.bs_counts,
f.ac_counts, f.tr_counts, f.er_counts,
f.zr_flow_proportion,
f.matrix_flow_proportion,
f.matrix_ao_val,
f.one_two_val, f.three_four_val, f.five_six_val, f.eight_val,
f.category_first_id,
f.category_id,
f.first_category_rank,
f.current_category_rank,
f.weight,
f.volume,
f.asin_weight_ratio,
f.color, f.size, f.style, f.material,
COALESCE(uma.package_quantity, f.package_quantity) AS package_quantity,
f.is_package_quantity_abnormal,
f.variation_num,
f.page_inventory,
f.activity_type,
COALESCE(f.launch_time, kp.keepa_launch_time) AS launch_time,
CASE
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(f.asin_crawl_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(f.asin_crawl_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(f.asin_crawl_date, COALESCE(f.launch_time, kp.keepa_launch_time)) <= 1080 THEN 6
ELSE 7
END AS launch_time_type,
f.img_url,
f.img_num,
f.img_type AS img_type_arr,
f.img_info,
f.account_id,
f.account_name,
f.buy_box_seller_type,
f.site_name,
f.follow_sellers_count,
f.asin_lob_info,
f.is_contains_lob_info,
f.asin_lqs_rating,
f.asin_lqs_rating_detail,
f.amazon_label,
f.is_movie_label,
f.is_brand_label,
f.multi_color_flag,
f.multi_color_str,
f.rank_type, f.ao_val_type, f.price_type, f.rating_type,
f.size_type, f.weight_type, f.site_name_type, f.quantity_variation_type,
f.rank_rise, f.rank_mom AS rank_change, f.rank_yoy,
f.ao_rise, f.ao_mom AS ao_change, f.ao_yoy,
f.price_rise, f.price_mom AS price_change, f.price_yoy,
f.rating_rise, f.rating_mom AS rating_change, f.rating_yoy,
f.comments_rise, f.comments_mom AS comments_change, f.comments_yoy,
f.bsr_orders_rise, f.bsr_orders_mom AS bsr_orders_change, f.bsr_orders_yoy,
f.sales_rise, f.sales_mom AS sales_change, f.sales_yoy,
f.variation_rise, f.variation_mom AS variation_change, f.variation_yoy,
f.bought_month_mom, f.bought_month_yoy,
pr.ocean_profit, pr.air_profit,
FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60) AS tracking_since,
CASE
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(f.asin_crawl_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(f.asin_crawl_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(f.asin_crawl_date, FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)) <= 1080 THEN 6
ELSE 7
END AS tracking_since_type,
kp.package_length,
kp.package_width,
kp.package_height,
CASE WHEN kp.item_weight > 0 THEN kp.item_weight ELSE kp.package_weight END AS item_weight,
CASE WHEN ba.brand_name_norm IS NOT NULL THEN 1 ELSE 0 END AS is_alarm_brand,
CAST(-1 AS SMALLINT) AS bsr_best_orders_type,
CAST(
CASE
WHEN sa.asin IS NOT NULL OR ab.brand_lower IS NOT NULL THEN 1
WHEN (
(COALESCE(chf.is_need_cat, 0) = 1 AND COALESCE(chd.is_need_cat, 0) = 1)
OR f.asin NOT LIKE 'B0%'
) THEN 2
WHEN (COALESCE(chf.is_hide_cat, 0) = 1 OR COALESCE(chc.is_hide_cat, 0) = 1)
AND NOT (COALESCE(chf.is_white_cat, 0) = 1 OR COALESCE(chc.is_white_cat, 0) = 1) THEN 3
ELSE 0
END
AS SMALLINT) AS asin_type,
uma.usr_mask_progress,
COALESCE(uma.usr_mask_type, umc.usr_mask_type) AS usr_mask_type,
COALESCE(aa.auctions_num, 0) AS auctions_num,
COALESCE(aa.auctions_num_all, 0) AS auctions_num_all,
COALESCE(aa.skus_num_creat, 0) AS skus_num_creat,
COALESCE(aa.skus_num_creat_all, 0) AS skus_num_creat_all,
f.title_matching_degree,
bb.brand_badge_reason,
ARRAY_JOIN(ARRAY_SLICE(SPLIT_BY_STRING(stem_en(f.title), ' '), 1, 15), ' ') AS title_stem_15,
ai.ai_package_quantity,
ai.ai_package_quantity_arr,
ai.ai_material,
ai.ai_color,
ai.ai_appearance,
ai.ai_size,
ai.ai_shape,
ai.ai_function,
ai.ai_scene_title,
ai.ai_scene_comment,
ai.ai_uses,
ai.ai_theme,
ai.ai_crowd,
ai.ai_short_desc,
addr.seller_province,
addr.seller_city,
-- ===== 年度聚合(dwt/dwt_flow_asin_year.py 在 Hive 预计算,经 sync_extra_table 同步到 Doris)=====
DATE_FORMAT(f.date_info, '%Y-%m') AS latest_date_info,
COALESCE(agg.bought_month_total, 0) AS bought_month_total,
agg.bought_month_1, agg.bought_month_2, agg.bought_month_3, agg.bought_month_4,
agg.bought_month_5, agg.bought_month_6, agg.bought_month_7, agg.bought_month_8,
agg.bought_month_9, agg.bought_month_10, agg.bought_month_11, agg.bought_month_12,
agg.bought_month_q1, agg.bought_month_q2, agg.bought_month_q3, agg.bought_month_q4,
agg.total_appear_month,
COALESCE(agg.bsr_seen_count_total, 0) AS bsr_seen_count_total,
COALESCE(agg.nsr_seen_count_total, 0) AS nsr_seen_count_total,
agg.bought_month_peak,
agg.peak_month_arr,
COALESCE(agg.is_periodic_flag, 0) AS is_periodic_flag,
COALESCE(agg.is_seasonal_flag, 0) AS is_seasonal_flag
FROM `dwt`.`{site_name}_flow_asin_365day` f
LEFT JOIN `dwt`.`{site_name}_flow_asin_365day_extra` agg ON f.asin = agg.asin
-- ===== 利润率 =====
LEFT JOIN `dwd`.`dwd_asin_profit_rate_latest` pr
ON f.asin = pr.asin AND f.price = pr.price AND pr.site_name = '{site_name}'
-- ===== Keepa 详情 =====
LEFT JOIN `dwd`.`dwd_keepa_asin_detail` kp
ON f.asin = kp.asin AND kp.site_name = '{site_name}'
-- ===== 预警品牌 =====
LEFT JOIN (
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand_name_norm
FROM `selection`.`brand_alert_erp`
WHERE brand_name IS NOT NULL
) ba ON f.brand = ba.brand_name_norm
-- ===== 自营 ASIN =====
LEFT JOIN (
SELECT asin FROM `mysql_selection`.`selection`.`us_self_asin` GROUP BY asin
) sa ON f.asin = sa.asin
-- ===== 自有品牌 =====
LEFT JOIN (
SELECT DISTINCT LOWER(TRIM(brand_name)) AS brand_lower
FROM `mysql_selection`.`selection`.`amazon_brand` WHERE brand_type = '1'
) ab ON LOWER(f.brand) = ab.brand_lower
-- ===== 一级分类规则 =====
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chf ON f.category_first_id = chf.category_id_base
-- ===== 当前分类规则 =====
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chc ON f.category_id = chc.category_id_base
-- ===== 描述匹配分类规则 =====
LEFT JOIN (
SELECT category_id_base,
MAX(CASE WHEN hide_type = 'is_need' THEN 1 ELSE 0 END) AS is_need_cat,
MAX(CASE WHEN hide_type = 'is_hide' THEN 1 ELSE 0 END) AS is_hide_cat,
MAX(CASE WHEN hide_type = 'is_white' THEN 1 ELSE 0 END) AS is_white_cat
FROM `mysql_selection`.`selection`.`us_bs_category_hide` GROUP BY category_id_base
) chd ON f.desc_category_first_id = chd.category_id_base
-- ===== 用户标记 ASIN/品类 =====
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
-- ===== 拍卖/SKU =====
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
ON f.brand = bb.brand AND bb.site_name = '{site_name}' AND bb.date_info = DATE_FORMAT(f.date_info, '%Y-%m')
-- ===== AI分析数据 =====
LEFT JOIN (
SELECT
asin,
package_quantity AS ai_package_quantity,
package_quantity_arr AS ai_package_quantity_arr,
material AS ai_material,
color AS ai_color,
appearance AS ai_appearance,
size AS ai_size,
shape AS ai_shape,
function AS ai_function,
scene_title AS ai_scene_title,
scene_comment AS ai_scene_comment,
uses AS ai_uses,
theme AS ai_theme,
crowd AS ai_crowd,
short_desc AS ai_short_desc
FROM `selection`.`{site_name}_ai_asin_analyze_detail`
) ai ON f.asin = ai.asin
-- ===== 卖家地址(按 site_name 查对应分区) =====
LEFT JOIN (
SELECT seller_id, fb_country_name, seller_province, seller_city
FROM `dim`.`dim_seller_address`
WHERE site_name = '{site_name}'
) addr ON f.account_id = addr.seller_id AND f.site_name = addr.fb_country_name
"""
# ============================================================
# [Step 8] 流程记录表更新(月+年各写一条,仅 formal 模式)
# ============================================================
def modify_mission_record_status(site_name, date_info, result_type):
"""流程记录表更新:仅 formal 模式才入库 mysql workflow_everyday,月/年各写一条记录"""
if result_type != 'formal':
print(f"[Step 8] result_type={result_type},跳过流程记录表更新")
return
record_table = 'workflow_everyday'
engine_mysql = DBUtil.get_db_engine(db_type=DbTypes.mysql.name, site_name='us')
def _write_record(date_type, table_name_field, status_remark):
select_sql = (
f"select id from {record_table} where site_name='{site_name}' and date_type='{date_type}' "
f"and report_date='{date_info}' and page='流量选品' and status_val=14 and is_end='是'"
)
df_is_finished = pd.read_sql(select_sql, engine_mysql)
if df_is_finished.empty:
replace_sql = f"""
replace into {record_table} (site_name, report_date, status, status_val, table_name, date_type, page, is_end, remark, export_db_type)
VALUES ('{site_name}', '{date_info}', '{status_remark}', 14, '{table_name_field}', '{date_type}', '流量选品', '是', '{status_remark}', 'doris')
"""
DBUtil.exec_sql('mysql', 'us', replace_sql)
print(f"[Step 8] 流程记录表 workflow_everyday 已写入:{site_name} {date_info} date_type={date_type}")
else:
print(f"[Step 8] 流程记录表已存在该记录,跳过:date_type={date_type}")
_write_record('month', f'{site_name}_flow_asin_last_month', '流量选品计算完毕')
_write_record('year', f'{site_name}_flow_asin_365day', '流量选品近一年计算完毕')
def main(site_name, date_info, result_type='formal'):
assert site_name in SUPPORTED_SITES, f"不支持的站点:{site_name},仅支持 us/uk/de"
assert result_type in ('formal', 'test'), f"不支持的 result_type:{result_type},仅支持 formal/test"
doris_table = f"{site_name}_flow_asin_month" # dwt 主表,不区分 test/formal
date_info_underscore = date_info.replace('-', '_')
env_suffix = '_test' if result_type == 'test' else ''
month_selection_table = f"{site_name}_flow_asin_month_{date_info_underscore}{env_suffix}"
year_selection_table = f"{site_name}_flow_asin_365day{env_suffix}"
last_12_month = sorted(CommonUtil.get_month_offset(date_info, -i) for i in range(0, 12))
print(f"启动:site={site_name}, date_info={date_info}, result_type={result_type}")
print(f" dwt 主表:dwt.{doris_table}(不区分 test/formal)")
print(f" selection 月物化表:selection.{month_selection_table}")
print(f" selection 年表:selection.{year_selection_table}")
print(f" 近12月窗口:{last_12_month}")
spark = SparkUtil.get_spark_session(
f"DwtFlowAsinMonth: {site_name} {date_info} {result_type}"
)
# ===== [Step 1] Doris 建表 selection.{month_selection_table} =====
print(f"[Step 1] Doris 建表 selection.{month_selection_table}")
_exec_doris_sql([build_month_create_table_sql(month_selection_table, date_info)])
# ===== [Step 2] 读 Hive dwt_flow_asin 月数据 + 字段规范化 =====
print(f"[Step 2] 读 Hive dwt_flow_asin 月数据 + 字段规范化")
df_save = read_and_normalize_month_data(spark, site_name, date_info)
# ===== [Step 3] 写入 Doris dwt 主表 =====
write_dwt_month_table(df_save, doris_table)
# ===== [Step 4] 写入 Doris dwt 年表 dwt.{site}_flow_asin_365day(自动聚合+清理过期数据)=====
print(f"[Step 4] 写入 Doris dwt 年表 dwt.{site_name}_flow_asin_365day")
_exec_doris_sql([build_stg_create_table_sql(site_name)])
maintain_stg_table(site_name, last_12_month, date_info)
# ===== [Step 5] 同步年度聚合指标 dwt.{site}_flow_asin_365day_extra =====
print(f"[Step 5] 同步年度聚合指标 dwt.{site_name}_flow_asin_365day_extra")
sync_extra_table(spark, site_name, date_info)
# ===== [Step 6] Doris INSERT OVERWRITE 到 selection 月物化表 =====
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)])
# ===== [Step 7] Doris INSERT OVERWRITE 到 selection 年表 =====
print(f"[Step 7] Doris INSERT OVERWRITE selection.{year_selection_table}")
_exec_doris_sql([build_year_create_table_sql(year_selection_table)])
_exec_doris_sql([build_year_insert_overwrite_sql(site_name, year_selection_table)])
# ===== [Step 8] 流程记录表更新(月+年各一条,仅 formal 模式)=====
modify_mission_record_status(site_name, date_info, result_type) modify_mission_record_status(site_name, date_info, result_type)
print("success!") print("success!")
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment