Commit 1978d698 by hejiangming

新增品类调研计算 读取流量选品数据

parent d059a090
...@@ -14,6 +14,8 @@ from pyspark.sql import functions as F ...@@ -14,6 +14,8 @@ from pyspark.sql import functions as F
sys.path.append(os.path.dirname(sys.path[0])) # 上级目录 sys.path.append(os.path.dirname(sys.path[0])) # 上级目录
from pyspark.storagelevel import StorageLevel from pyspark.storagelevel import StorageLevel
from utils.templates import Templates from utils.templates import Templates
from utils.hdfs_utils import HdfsUtils # 落表前清 HDFS 分区目录(幂等重跑)
from utils.common_util import CommonUtil # 构造 HDFS 分区路径
# from ..utils.templates import Templates # from ..utils.templates import Templates
# from AmazonSpider.pyspark_job.utils.templates_test import Templates # from AmazonSpider.pyspark_job.utils.templates_test import Templates
from pyspark.sql.types import StringType, IntegerType, DoubleType, LongType from pyspark.sql.types import StringType, IntegerType, DoubleType, LongType
...@@ -36,6 +38,9 @@ class DwtBsTop100(Templates): ...@@ -36,6 +38,9 @@ class DwtBsTop100(Templates):
self.df_deduplicated = self.spark.sql("select 1+1;") # DF 对象占位符 self.df_deduplicated = self.spark.sql("select 1+1;") # DF 对象占位符
self.df_deduplicated2 = self.spark.sql("select 1+1;") # DF 对象占位符 self.df_deduplicated2 = self.spark.sql("select 1+1;") # DF 对象占位符
self.df_cate_bs_data = self.spark.sql("select 1+1;") # DF 对象占位符 self.df_cate_bs_data = self.spark.sql("select 1+1;") # DF 对象占位符
self.df_profit_rate = self.spark.sql("select 1+1;") # DF 对象占位符(利润率 asin+price 粒度)
self.band_margin_cols = [] # 区间利润率带闸列名(④价格12+⑤月销12=24),handle_bs_add_column 填充
self.pkg_band_cols = [] # 打包段"带闸月销列"名(⑥,7段),handle_bs_add_column 填充
self.df_cate_tree = self.spark.sql("select 1+1;") # DF 对象占位符 self.df_cate_tree = self.spark.sql("select 1+1;") # DF 对象占位符
self.df_nodes_num_concat = self.spark.sql("select 1+1;") # DF 对象占位符 self.df_nodes_num_concat = self.spark.sql("select 1+1;") # DF 对象占位符
self.df_nodes_num_concat2 = self.spark.sql("select 1+1;") # DF 对象占位符 self.df_nodes_num_concat2 = self.spark.sql("select 1+1;") # DF 对象占位符
...@@ -45,29 +50,36 @@ class DwtBsTop100(Templates): ...@@ -45,29 +50,36 @@ class DwtBsTop100(Templates):
def read_data(self): def read_data(self):
# self.df_bs_data = self.spark.read.csv("/opt/module/spark-3.2.0-bin-hadoop3.2/demo/py_demo/dwt/asin_bsr_us_2025-11.csv",header=True,inferSchema=True).limit(1000000).cache()
sql = f""" sql = f"""
select asin, asin_type, bsr_orders, category_first_id, category_id, first_category_rank, current_category_rank, select asin, asin_type, bsr_orders, category_first_id, category_id, first_category_rank, current_category_rank,
asin_price, asin_rating, asin_buy_box_seller_type, asin_is_new, asin_total_comments, asin_launch_time, asin_launch_time_type, asin_price, asin_rating, asin_buy_box_seller_type, asin_is_new, asin_total_comments, asin_launch_time, asin_launch_time_type,
asin_brand_name, is_brand_label, asin_bought_month as buy_data_bought_month asin_brand_name, is_brand_label, asin_bought_month as buy_data_bought_month, package_quantity, seller_country_name, account_id
from dwt_flow_asin where site_name='{self.site_name}' and date_type='{self.date_type}' and date_info='{self.date_info}' from dwt_flow_asin where site_name='{self.site_name}' and date_type='{self.date_type}' and date_info='{self.date_info}'
""" """
print(f"1. 读取dwt_flow_asin数据: sql -- {sql}") print(f"1. 读取dwt_flow_asin数据: sql -- {sql}")
self.df_bs_data = self.spark.sql(sqlQuery=sql).cache() # 下游线性链只读一次,去掉 cache
self.df_bs_data = self.spark.sql(sqlQuery=sql)
# self.df_bs_data.show(10, truncate=False) # self.df_bs_data.show(10, truncate=False)
print("df_bs_data数量", self.df_bs_data.count()) # print("df_bs_data数量", self.df_bs_data.count())
import pandas as pd import pandas as pd
if pd.__version__ >= "2.0.0" and not hasattr(pd.DataFrame, "iteritems"): if pd.__version__ >= "2.0.0" and not hasattr(pd.DataFrame, "iteritems"):
# 给 DataFrame 动态添加 iteritems 方法,映射到新版 items() # 给 DataFrame 动态添加 iteritems 方法,映射到新版 items()
pd.DataFrame.iteritems = pd.DataFrame.items pd.DataFrame.iteritems = pd.DataFrame.items
engine = get_remote_engine(site_name="us", db_type="mysql") engine = get_remote_engine(site_name="us", db_type="mysql")
sql = f"SELECT nodes_num, category_id, category_parent_id, category_first_id, redirect_flag, redirect_first_id, en_name as cur_category_name from us_bs_category where delete_time is null;" sql = f"SELECT nodes_num, category_id, category_parent_id, category_first_id, en_name as cur_category_name from us_bs_category where delete_time is null;"
df_cate = engine.read_sql(sql) df_cate = engine.read_sql(sql)
self.df_cate_bs_data = self.spark.createDataFrame(df_cate) self.df_cate_bs_data = self.spark.createDataFrame(df_cate)
self.df_cate_bs_data.cache() self.df_cate_bs_data.cache()
# print(self.df_cate_bs_data.show()) # print(self.df_cate_bs_data.show())
print("df_cate_bs_data数量", self.df_cate_bs_data.count()) print("df_cate_bs_data数量", self.df_cate_bs_data.count())
# 读利润率数据(dim_asin_profit_rate_info)——给每个 ASIN 补"海运/空运毛利润率",
sql_pr = f"""select asin, price as asin_price, ocean_profit, air_profit
from dim_asin_profit_rate_info where site_name='{self.site_name}'"""
print(f"2. 读取利润率数据: sql -- {sql_pr}")
self.df_profit_rate = self.spark.sql(sqlQuery=sql_pr).dropDuplicates(['asin', 'asin_price'])
def clean_column(self, df, col_name): def clean_column(self, df, col_name):
"""统一处理 ID 列:去除 .0 后缀,处理空值""" """统一处理 ID 列:去除 .0 后缀,处理空值"""
return df.withColumn( return df.withColumn(
...@@ -84,20 +96,12 @@ class DwtBsTop100(Templates): ...@@ -84,20 +96,12 @@ class DwtBsTop100(Templates):
) )
def handle_bs_redirect_flag(self): def handle_bs_redirect_flag(self):
print("handle_bs_redirect_flag调用 重定向一级分类关系处理") print("handle_bs_redirect_flag调用 分类树按指定列去重(不做重定向)")
# 重定向一级分类关系处理
self.df_cate_bs_data = self.df_cate_bs_data.withColumn("category_first_id",
F.when(self.df_cate_bs_data.redirect_flag == 1,
self.df_cate_bs_data.redirect_first_id).otherwise(
self.df_cate_bs_data.category_first_id))
# 根据指定列去重 # 根据指定列去重
drop_col_list = ['category_id', 'category_first_id', 'category_parent_id', 'nodes_num', 'cur_category_name'] drop_col_list = ['category_id', 'category_first_id', 'category_parent_id', 'nodes_num', 'cur_category_name']
# self.df_cate_bs_data = self.df_cate_bs_data.drop_duplicates(drop_col_list) # 根据指定列去重 # self.df_cate_bs_data = self.df_cate_bs_data.drop_duplicates(drop_col_list) # 根据指定列去重
self.df_cate_bs_data = self.df_cate_bs_data.withColumn("_cate_order", self.df_cate_bs_data = self.df_cate_bs_data.withColumn("_cate_order",
F.monotonically_increasing_id()) # 添加一个顺序列 方便去重操作 F.monotonically_increasing_id()) # 添加一个顺序列 方便去重操作
print("df_cate_bs_data的_cate_order的值是")
print(self.df_cate_bs_data.select("_cate_order").tail(10))
w = Window.partitionBy(drop_col_list).orderBy(F.col("_cate_order").asc()) w = Window.partitionBy(drop_col_list).orderBy(F.col("_cate_order").asc())
# 2. 生成行号并过滤 # 2. 生成行号并过滤
self.df_cate_bs_data = self.df_cate_bs_data.withColumn("rn", F.row_number().over(w)) \ self.df_cate_bs_data = self.df_cate_bs_data.withColumn("rn", F.row_number().over(w)) \
...@@ -142,7 +146,7 @@ class DwtBsTop100(Templates): ...@@ -142,7 +146,7 @@ class DwtBsTop100(Templates):
# order_cols = [c for c in self.df_cate_tree.columns if c.startswith("_cate_order")] # order_cols = [c for c in self.df_cate_tree.columns if c.startswith("_cate_order")]
# for col in order_cols: # for col in order_cols:
# df_cate_tree = df_cate_tree.drop(col) # df_cate_tree = df_cate_tree.drop(col)
print("self.df_cate_tree数量", self.df_cate_tree.count()) # print("self.df_cate_tree数量", self.df_cate_tree.count())
print("self.df_cate_tree列名", self.df_cate_tree.columns) print("self.df_cate_tree列名", self.df_cate_tree.columns)
# print(self.df_cate_tree.show(10, truncate=False)) # print(self.df_cate_tree.show(10, truncate=False))
print("handle_bs_hierarchy_to_row 层级关系处理完毕") print("handle_bs_hierarchy_to_row 层级关系处理完毕")
...@@ -174,20 +178,58 @@ class DwtBsTop100(Templates): ...@@ -174,20 +178,58 @@ class DwtBsTop100(Templates):
# 排除asin_type not in (1, 2), 1:内部asin,2:视频音乐图书等asin # 排除asin_type not in (1, 2), 1:内部asin,2:视频音乐图书等asin
print("排除asin_type") print("排除asin_type")
self.df_bs_data = self.df_bs_data.filter(~(self.df_bs_data.asin_type.isin([1, 2]))) self.df_bs_data = self.df_bs_data.filter(~(self.df_bs_data.asin_type.isin([1, 2])))
# ============================================================
# 下游 handle_bs_nodes_num_concat/concat2 的 cols_select_list、handle_bs_proportion
# 都写死 launch_time_type0_rate ~ type7_rate。若靠 distinct 造列,某分区恰好 没有某个桶值(如冷门小站点/极早月份无"最近30天上架"=1)→ 该列不生成 →
# 下游 select 点到不存在的列直接 AnalysisException,整个任务崩。
# 固定 0~7 造列后,无论数据缺不缺值,8 列恒定存在,缺值时该桶列全为 False → 下游 sum(cast int)=0 → 占比 0/asin_count=0(0%)
# ============================================================
print("处理上架时间分类") print("处理上架时间分类")
# 处理上架时间分类 上架时间类型(默认0,1:最近30天,2:1-3个月,3:3-6个月,4:6-12个月,5:1-2年,6:2-3年,7:3年以上) # 处理上架时间分类 上架时间类型(默认0,1:最近30天,2:1-3个月,3:3-6个月,4:6-12个月,5:1-2年,6:2-3年,7:3年以上)
launch_time_type_list = [x[0] for x in self.df_bs_data.select("asin_launch_time_type").distinct().collect()] # launch_time_type_list = [x[0] for x in self.df_bs_data.select("asin_launch_time_type").distinct().collect()]
for rate in launch_time_type_list: # for rate in launch_time_type_list:
for rate in range(8):
self.df_bs_data = self.df_bs_data.withColumn(f"launch_time_type{rate}_rate", self.df_bs_data = self.df_bs_data.withColumn(f"launch_time_type{rate}_rate",
self.df_bs_data["asin_launch_time_type"] == rate) self.df_bs_data["asin_launch_time_type"] == rate)
# # 处理卖家类型 buybox卖家类型(1:Amazon,2:FBA,3FBM,4无BB卖家)
# asin_buy_box_seller_type_list = [x[0] for x in self.df_bs_data.select("asin_buy_box_seller_type").distinct().collect()]
print("处理卖家类型") print("处理卖家类型")
# 处理卖家类型 buybox卖家类型(1:Amazon,2:FBA,3FBM,4无BB卖家) # for rate in asin_buy_box_seller_type_list:
asin_buy_box_seller_type_list = [x[0] for x in for rate in range(5):
self.df_bs_data.select("asin_buy_box_seller_type").distinct().collect()]
for rate in asin_buy_box_seller_type_list:
self.df_bs_data = self.df_bs_data.withColumn(f"asin_buy_box_seller_type{rate}_rate", self.df_bs_data = self.df_bs_data.withColumn(f"asin_buy_box_seller_type{rate}_rate",
self.df_bs_data["asin_buy_box_seller_type"] == rate) self.df_bs_data["asin_buy_box_seller_type"] == rate)
# ============================================================
# 【卖家所属地占比】 造 3 个布尔列(美/中/其他),走 rate 路径——handle_bs_level 里 rate_cols
# ('rate' in c)自动 sum(cast int)=各国 asin 数,handle_bs_proportion 再 ÷ 分类总 asin 数 = 数量占比。
# seller_country_name 是国家代码(US/CN/HK/IN/...),含 null。
# 【null 归属·必须 null-safe】"其他国家"含 null:
# - us/cn 用 when(==,True).otherwise(False):null 比较得 null → 落 otherwise=False(不进美/中)。
# - other 用 when(US或CN, False).otherwise(True):null → 条件 null → 落 otherwise=True(进其他)。
# 三者互斥且穷尽(任一 asin 恰好命中一个) → 三列计数和 = asin_count → 三个占比相加=1。
# ============================================================
print("处理卖家所属地")
c = F.col("seller_country_name")
self.df_bs_data = self.df_bs_data.withColumn("seller_country_us_rate", F.when(c == "US", True).otherwise(False))
self.df_bs_data = self.df_bs_data.withColumn("seller_country_cn_rate", F.when(c == "CN", True).otherwise(False))
self.df_bs_data = self.df_bs_data.withColumn("seller_country_other_rate",
F.when((c == "US") | (c == "CN"), False).otherwise(True))
# ============================================================
# 【ASIN 月销 50+/100+/200+ 占比】:月销≥阈值的 asin 数 ÷ 分类总 asin 数(数量占比)。
# rate 路径——handle_bs_level 里 rate_cols('rate' in c)自动 sum(cast int)=达标 asin 数,
# handle_bs_proportion 再 ÷ asin_count = 占比。
# 【累计阈值·不互斥】50+ 含 100+ 含 200+,三者不相加=1,各是独立占比。
# 【null-safe + cast】asin_amazon_orders 是 string,cast double 再比;月销 null → 条件 null →
# otherwise(False)→0,算作"不达标",在分母(asin_count)里、不在分子。此时 handle_bse_column_rename
# 已把 buy_data_bought_month 改名为 asin_amazon_orders,故这里可直接用。
# ============================================================
print("处理月销50+/100+/200+")
_mo = F.col("asin_amazon_orders").cast("double")
self.df_bs_data = self.df_bs_data.withColumn("asin_orders_ge_50_rate", F.when(_mo >= 50, True).otherwise(False))
self.df_bs_data = self.df_bs_data.withColumn("asin_orders_ge_100_rate", F.when(_mo >= 100, True).otherwise(False))
self.df_bs_data = self.df_bs_data.withColumn("asin_orders_ge_200_rate", F.when(_mo >= 200, True).otherwise(False))
def handle_bs_asin_deduplicated(self): def handle_bs_asin_deduplicated(self):
print("handle_bs_asin_deduplicated 调用") print("handle_bs_asin_deduplicated 调用")
# # 对每列 转换类型为 string 然后统一将空值 nan null 设置为none 把数据末尾的.0去除 # # 对每列 转换类型为 string 然后统一将空值 nan null 设置为none 把数据末尾的.0去除
...@@ -220,7 +262,8 @@ class DwtBsTop100(Templates): ...@@ -220,7 +262,8 @@ class DwtBsTop100(Templates):
self.df_deduplicated = self.df_deduplicated.withColumn("rn", F.row_number().over(window_spec)) \ self.df_deduplicated = self.df_deduplicated.withColumn("rn", F.row_number().over(window_spec)) \
.filter(F.col("rn") == 1) \ .filter(F.col("rn") == 1) \
.drop("rn") .drop("rn")
self.df_deduplicated .cache() # 下游线性链只读一次,去掉 cache
# self.df_deduplicated .cache()
print("df_deduplicated打印") print("df_deduplicated打印")
# self.df_deduplicated.show(10, truncate=False) # self.df_deduplicated.show(10, truncate=False)
# print(self.df_deduplicated.columns) # print(self.df_deduplicated.columns)
...@@ -228,6 +271,19 @@ class DwtBsTop100(Templates): ...@@ -228,6 +271,19 @@ class DwtBsTop100(Templates):
def handle_bs_add_column(self): def handle_bs_add_column(self):
print("handle_bs_add_column 调用") print("handle_bs_add_column 调用")
# 【利润率 join · 在 ASIN 层做一次】此时 df_deduplicated 是"一个 asin 一行
self.df_deduplicated = self.df_deduplicated.join(self.df_profit_rate, on=['asin', 'asin_price'], how='left')
self.df_deduplicated = self.df_deduplicated \
.withColumnRenamed("ocean_profit", "asin_ocean_freight_gross_margin") \
.withColumnRenamed("air_profit", "asin_air_freight_gross_margin")
# 销量口径切换:所有"销量"类字段的底层,由"BSR排名预估销量"改为"亚马逊前台真实月销"。 BSR预估 = 排名→查表×系数 猜出来的;同款商品多变体常共享同一排名
# 把底层列 asin_bsr_orders 直接重定向为 asin_amazon_orders,下游 new/brand/销售额/sum/均值/占比 自动跟着变(字段名不改 避免动 DDL 和下游)
# 页面不要显示 asin_bsr_orders_sum 这一列(保留字段仅为兼容,避免两列重复展示)
self.df_deduplicated = self.df_deduplicated.withColumn(
"asin_bsr_orders", F.col("asin_amazon_orders")
)
self.df_deduplicated = self.df_deduplicated.withColumn("asin_bsr_orders_new", self.df_deduplicated = self.df_deduplicated.withColumn("asin_bsr_orders_new",
F.when(self.df_deduplicated.asin_is_new == 1, F.when(self.df_deduplicated.asin_is_new == 1,
self.df_deduplicated.asin_bsr_orders).otherwise( self.df_deduplicated.asin_bsr_orders).otherwise(
...@@ -248,56 +304,76 @@ class DwtBsTop100(Templates): ...@@ -248,56 +304,76 @@ class DwtBsTop100(Templates):
self.df_deduplicated = self.df_deduplicated.withColumn("asin_bsr_orders_sales_brand", self.df_deduplicated = self.df_deduplicated.withColumn("asin_bsr_orders_sales_brand",
self.df_deduplicated.asin_bsr_orders_brand * self.df_deduplicated.asin_price) self.df_deduplicated.asin_bsr_orders_brand * self.df_deduplicated.asin_price)
def handle_bs_nodes_num_concat(self):
print("handle_bs_nodes_num_concat 调用")
raw_list = self.df_deduplicated.filter(F.col("nodes_num").isNotNull()).agg(F.collect_set("nodes_num")).first()[
0]
nodes_num_list = sorted([int(x) for x in raw_list])
print("df_node_num_list打印", nodes_num_list)
cols_select_list = ['asin', 'asin_amazon_orders', 'asin_bsr_orders', 'asin_bsr_orders_new',
'asin_bsr_orders_brand',
'asin_bsr_orders_sales', 'asin_bsr_orders_sales_new', 'asin_bsr_orders_sales_brand',
'asin_bs_cate_1_rank', 'asin_price', 'asin_rating', 'asin_is_new', 'is_brand_label',
'asin_total_comments', 'asin_bs_cate_current_id', 'asin_bs_cate_1_id',
'launch_time_type0_rate', 'launch_time_type1_rate', 'launch_time_type2_rate',
'launch_time_type3_rate', 'launch_time_type4_rate',
'launch_time_type5_rate', 'launch_time_type6_rate', 'launch_time_type7_rate',
'asin_buy_box_seller_type0_rate', 'asin_buy_box_seller_type1_rate',
'asin_buy_box_seller_type2_rate', 'asin_buy_box_seller_type3_rate',
'asin_buy_box_seller_type4_rate']
df_node_num_list = []
for idx, num in enumerate(nodes_num_list):
num = int(num)
df_node_num = self.df_deduplicated.filter(self.df_deduplicated["nodes_num"] == num).select(
*cols_select_list)
# print((df_node_num.count(), len(df_node_num.columns)))
df_node_num = df_node_num.withColumnRenamed("asin_bs_cate_current_id", f"asin_bs_cate_{num - 1}_id")
cols_list = [f"asin_bs_cate_{i}_id" for i in
range(1, num)] # 生成当前分类的层级路径 如 num=4 生成 1-3的asin_bs_cate_x_id
# 添加循环顺序列
df_node_num = df_node_num.withColumn("_loop_order", F.lit(idx))
df_first_tree = self.df_cate_tree.select(*cols_list).dropDuplicates()
df_first_tree = df_first_tree.withColumn(
"asin_bs_cate_parent_id_join",
F.concat_ws("&&&&", *[df_first_tree[c].cast("string") for c in cols_list[:-1]])
)
df_node_num = df_node_num.join(df_first_tree, how="left",
on=["asin_bs_cate_1_id", f"asin_bs_cate_{num - 1}_id"])
# 以一级分类id和当前分类ID 为条件 左连接
df_node_num_list.append(df_node_num)
self.df_nodes_num_concat = reduce(lambda a, b: a.unionByName(b, allowMissingColumns=True), df_node_num_list) # ============================================================
# 【区间利润率 · 造"带闸列"】:
# 段内 = 该 asin 利润率、段外 = null;下游 handle_bs_level 对每列 F.avg = 该分类该段的平均利润率。
# 沿用上游 dwt_flow_asin.asin_price_type 的左闭右开 >=下界 AND <上界。
# 价格:>0 排除 0/null(价格 0 是脏数据);10→10-20、80→"80及以上"(>=80)。
# 月销:字段 asin_amazon_orders(前台月销)实测为干净正整数、无 0,null(前台无徽章)条件不成立自动排除。
# 【显式 cast("double")】asin_amazon_orders 是 string 类型(前台月销原始文本),Hive CAST 曾抽风,
# Spark 里显式转 double 再比大小更稳;asin_price 同样转一下统一。
# 【列名】与 整体均值列同源前缀 asin_{ocean|air}_freight_gross_margin + 段后缀;
# 不能含 "rate",否则被 handle_bs_level 的 rate_cols('rate' in c)当计数 sum 抓走。
# ============================================================
price_col = F.col("asin_price").cast("double")
orders_col = F.col("asin_amazon_orders").cast("double")
band_defs = [
# 价格区间(后缀 price_*)
("price_0_10", (price_col > 0) & (price_col < 10)),
("price_10_20", (price_col >= 10) & (price_col < 20)),
("price_20_30", (price_col >= 20) & (price_col < 30)),
("price_30_50", (price_col >= 30) & (price_col < 50)),
("price_50_80", (price_col >= 50) & (price_col < 80)),
("price_80", (price_col >= 80)),
# 月销区间(后缀 orders_*)
("orders_0_50", (orders_col >= 0) & (orders_col < 50)),
("orders_50_100", (orders_col >= 50) & (orders_col < 100)),
("orders_100_200", (orders_col >= 100) & (orders_col < 200)),
("orders_200_500", (orders_col >= 200) & (orders_col < 500)),
("orders_500_1000", (orders_col >= 500) & (orders_col < 1000)),
("orders_1000", (orders_col >= 1000)),
]
margin_srcs = [
("ocean", "asin_ocean_freight_gross_margin"),
("air", "asin_air_freight_gross_margin"),
]
# 存全部带闸列名到 self(12 + 12 = 24 列),供 cols_select_list / 聚合 / round 三处复用同一份,避免名字漂移
self.band_margin_cols = []
for suffix, cond in band_defs:
for mtype, src in margin_srcs:
col = f"asin_{mtype}_freight_gross_margin_{suffix}"
self.df_deduplicated = self.df_deduplicated.withColumn(col, F.when(cond, F.col(src)))
self.band_margin_cols.append(col)
# ============================================================
# 【打包段月销占比 · 造"带闸月销列"】:按打包件数分段,展示"该段月销 ÷ 品类总月销"。
# 每个打包段造一列:打包件数落本段 = 该 asin 月销(asin_amazon_orders),否则 null。
# 下游 handle_bs_level 对每列 F.sum = 该分类该打包段的月销合计;calculate_column 里再 ÷ 品类总月销。
# 打包件数是整数,左闭右开:1 / [2,5) / [5,10) / [10,20) / [20,50) / [50,100) / [100,+)。
# package_quantity 上游 na.fill 到 1、无 null 无 0,故最小段就是 =1。显式 cast("double") 防类型坑。
# 【命名】中缀 pkg_{段}、列名不含 "rate"(带闸列是月销求和的中间量,不能被 rate_cols 当计数);
# 最终占比列在 calculate_column 里生成 _rate 后缀。
# ============================================================
pkg_col = F.col("package_quantity").cast("double")
pkg_bands = [
("1", (pkg_col >= 1) & (pkg_col < 2)),
("2_4", (pkg_col >= 2) & (pkg_col < 5)),
("5_10", (pkg_col >= 5) & (pkg_col < 10)),
("10_20", (pkg_col >= 10) & (pkg_col < 20)),
("20_50", (pkg_col >= 20) & (pkg_col < 50)),
("50_100", (pkg_col >= 50) & (pkg_col < 100)),
("100", (pkg_col >= 100)),
]
self.pkg_band_cols = []
for suffix, cond in pkg_bands:
col = f"asin_pkg_{suffix}_orders_sum"
self.df_deduplicated = self.df_deduplicated.withColumn(col, F.when(cond, F.col("asin_amazon_orders")))
self.pkg_band_cols.append(col)
window_spec = Window.partitionBy("asin").orderBy(F.col("_loop_order").asc()) # 按 asin 去重,保留第一条 # 注:handle_bs_nodes_num_concat(非2版,2026-07-02 删除)——早期草稿,缺 fillna(nodes_num,-1) 会丢 null 分类的 asin、
# 生成行号并过滤 # 用 dropDuplicates 去重不确定,且结果 df_nodes_num_concat 从未被下游消费(空跑浪费)。正式逻辑见下方 concat2。
self.df_nodes_num_concat = self.df_nodes_num_concat.withColumn("rn", F.row_number().over(window_spec)) \ # 如需找回:开发母本 spark_hjm/bs_top100/dwt_bs_top100.py 或 git 历史。
.filter(F.col("rn") == 1) \
.drop("rn", "_loop_order")
print("df_nodes_num_concat打印") # 这个后续没用到
print(self.df_nodes_num_concat.count())
print(self.df_nodes_num_concat.columns)
def handle_bs_nodes_num_concat2(self): def handle_bs_nodes_num_concat2(self):
print("handle_bs_nodes_num_concat2 打印") print("handle_bs_nodes_num_concat2 打印")
...@@ -318,14 +394,37 @@ class DwtBsTop100(Templates): ...@@ -318,14 +394,37 @@ class DwtBsTop100(Templates):
'launch_time_type5_rate', 'launch_time_type6_rate', 'launch_time_type7_rate', 'launch_time_type5_rate', 'launch_time_type6_rate', 'launch_time_type7_rate',
'asin_buy_box_seller_type0_rate', 'asin_buy_box_seller_type1_rate', 'asin_buy_box_seller_type0_rate', 'asin_buy_box_seller_type1_rate',
'asin_buy_box_seller_type2_rate', 'asin_buy_box_seller_type3_rate', 'asin_buy_box_seller_type2_rate', 'asin_buy_box_seller_type3_rate',
'asin_buy_box_seller_type4_rate', 'nodes_num'] 'asin_buy_box_seller_type4_rate',
# 卖家所属地 3 个布尔列,带进拆行 select 才能活到 rate 路径求占比
'seller_country_us_rate', 'seller_country_cn_rate', 'seller_country_other_rate',
# 月销50+/100+/200+ 3 个布尔列,同走 rate 路径求数量占比
'asin_orders_ge_50_rate', 'asin_orders_ge_100_rate', 'asin_orders_ge_200_rate',
# 集中度要用的品牌名、卖家id(在 handle_bs_level 里做 top10 排序)
'asin_brand_name', 'account_id',
# 利润率两列(asin 层已 join),带进拆行 select 才能活到分层聚合求平均
'asin_ocean_freight_gross_margin', 'asin_air_freight_gross_margin',
'nodes_num']
# 区间利润率的 24 个带闸列(价格12+月销12)同样带进拆行 select,才能活到 handle_bs_level 求段内平均
cols_select_list = cols_select_list + self.band_margin_cols
# 打包段的 7 个带闸月销列同样带进拆行 select,才能活到 handle_bs_level 求段内月销合计
cols_select_list = cols_select_list + self.pkg_band_cols
df_node_num_list2 = [] df_node_num_list2 = []
for idx,num in enumerate(nodes_num_list2) : for idx,num in enumerate(nodes_num_list2) :
num = int(num) num = int(num)
df_node_num = self.df_deduplicated2.filter(self.df_deduplicated2.nodes_num == num).select(*cols_select_list) df_node_num = self.df_deduplicated2.filter(self.df_deduplicated2.nodes_num == num).select(*cols_select_list)
df_node_num = df_node_num.withColumn("_loop_order", F.lit(idx)) df_node_num = df_node_num.withColumn("_loop_order", F.lit(idx))
# print(num,df_node_num.count(),len(df_node_num.columns)) # print(num,df_node_num.count(),len(df_node_num.columns))
if num > 1: # ============================================================
# num==2 = 这个 ASIN 最深的分类就是"一级"本身(nodes_num:根=1、一级=2、二级=3…)。
# 此时 current_id 和 cols_select_list 里已选入的 asin_bs_cate_1_id 指向同一个分类。
# 若走通用逻辑(下面 elif)把 current_id 改名成 asin_bs_cate_{num-1}_id:num-1=1 →
# 变成 asin_bs_cate_1_id,和已有的那列撞名 → 出现两列同名 → 后面 unionByName报 "duplicate column asin_bs_cate_1_id" 直接崩。
# 一级没有更上层祖先、不需要 join 分类树补中间层,把多余的 current_id drop 掉即可;该 ASIN 靠已有的 asin_bs_cate_1_id 正常进入 level=1 聚合。
# 例:某 ASIN 最深就挂在 electronics(nodes_num=2)→ 保留 asin_bs_cate_1_id=electronics、drop current_id。
# ============================================================
if num == 2:
df_node_num = df_node_num.drop("asin_bs_cate_current_id")
elif num > 2:
df_node_num = df_node_num.withColumnRenamed("asin_bs_cate_current_id", f"asin_bs_cate_{num - 1}_id") df_node_num = df_node_num.withColumnRenamed("asin_bs_cate_current_id", f"asin_bs_cate_{num - 1}_id")
cols_list = [f"asin_bs_cate_{i}_id" for i in cols_list = [f"asin_bs_cate_{i}_id" for i in
range(1, num)] # 生成当前分类的层级路径 如 num=4 生成 1-3的asin_bs_cate_x_id range(1, num)] # 生成当前分类的层级路径 如 num=4 生成 1-3的asin_bs_cate_x_id
...@@ -363,15 +462,34 @@ class DwtBsTop100(Templates): ...@@ -363,15 +462,34 @@ class DwtBsTop100(Templates):
self.df_nodes_num_concat2 = df_nodes_num_concat2.withColumn("rn", F.row_number().over(window_spec)) \ self.df_nodes_num_concat2 = df_nodes_num_concat2.withColumn("rn", F.row_number().over(window_spec)) \
.filter(F.col("rn") == 1) \ .filter(F.col("rn") == 1) \
.drop("rn") .drop("rn")
self.df_nodes_num_concat2.cache() # 原缓存位置在此(clean_column 循环之前),循环会重新赋值成新 DF 使缓存失效 → 挪到循环之后
# self.df_nodes_num_concat2.cache()
# 处理 asin_bs_cate_ 空值 null # 处理 asin_bs_cate_ 空值 null
id_cols = [c for c in self.df_nodes_num_concat2.columns if 'asin_bs_cate_' in c and '_id' in c] id_cols = [c for c in self.df_nodes_num_concat2.columns if 'asin_bs_cate_' in c and '_id' in c]
for id_col in id_cols: for id_col in id_cols:
self.df_nodes_num_concat2 = self.clean_column(self.df_nodes_num_concat2, id_col) self.df_nodes_num_concat2 = self.clean_column(self.df_nodes_num_concat2, id_col)
# 缓存最终版:handle_bs_level 每层 filter + 3 组集中度都读它
self.df_nodes_num_concat2.cache()
print("df_nodes_num_concat2打印") print("df_nodes_num_concat2打印")
# self.df_nodes_num_concat2.show(10,truncate=False) # self.df_nodes_num_concat2.show(10,truncate=False)
print(self.df_nodes_num_concat2.count()) # print(self.df_nodes_num_concat2.count())
# print(self.df_nodes_num_concat2.columns) # print(self.df_nodes_num_concat2.columns)
def _cr10_top_sum(self, df, group_keys, sales_col, out_col, salt_n=50):
"""
两段式加盐取每个分类的 top10 销量之和(集中度分子),解决大分类窗口排序倾斜。
原理:全局 top10 必然落在它所属盐桶的 top10 里 → 先按(分类,盐桶)各取 top10 候选,
再从候选(每分类 ≤10*salt_n 行)取真正 top10 求和,结果与直接全量 top10 完全等价,
但把"一个 task 排几百万行"拆成"salt_n 个 task 各排几万行",消除 level=1 大类的倾斜。
"""
df_s = df.withColumn("_salt", (F.rand() * salt_n).cast("int"))
# 第一段:每个(分类,盐桶)内取 top10 候选
w1 = Window.partitionBy(*group_keys, "_salt").orderBy(F.col(sales_col).desc_nulls_last())
cand = df_s.withColumn("_rn_c", F.row_number().over(w1)).filter(F.col("_rn_c") <= 10)
# 第二段:候选里取真正 top10 求和
w2 = Window.partitionBy(*group_keys).orderBy(F.col(sales_col).desc_nulls_last())
top = cand.withColumn("_rn_c2", F.row_number().over(w2)).filter(F.col("_rn_c2") <= 10)
return top.groupBy(*group_keys).agg(F.sum(sales_col).alias(out_col))
def handle_bs_level(self): def handle_bs_level(self):
print("handle_bs_level 调用") print("handle_bs_level 调用")
import re import re
...@@ -398,7 +516,18 @@ class DwtBsTop100(Templates): ...@@ -398,7 +516,18 @@ class DwtBsTop100(Templates):
if level == 1: if level == 1:
group_keys = [curr_col] group_keys = [curr_col]
else: else:
group_keys = [curr_col, parent_col] # ============================================================
# 【多归属拆行 · 用当前层的祖先列分组】
# 同一分类可能挂在多个父级/一级下(如 Cables 同时在 electronics、pc 下)。
# 只按(当前,直接父级)分组会把不同一级来的 ASIN 合成一行、一级归属靠 min 随机挑、数量两支相加(超额)。
# 【为什么不用 concat2 里现成的 asin_bs_cate_parent_id_join】那个是每个 ASIN "最深路径"的拼接,
# 在中间层不代表"到当前层的路径"——用它当分组键会把某节点按各 ASIN 的更深路径过度打散
# (level=2 的行却带 level=6 的面包屑、裂成上百行)。
# 用"一级..当前层"这几列祖先列现场分组:asin_bs_cate_1_id..asin_bs_cate_{level}_id,
# 这才是"到当前层为止"的完整路径,唯一标识"这条路径下的这个节点",多归属才按真实路径裂行。
# 【例】三级 X:(electronics,Cables,X)、(pc,Cables,X) 各一行,各算各的。
# ============================================================
group_keys = [f"asin_bs_cate_{i}_id" for i in range(1, level + 1)]
aggs = [ aggs = [
F.sum("asin_amazon_orders").alias("asin_amazon_orders"), F.sum("asin_amazon_orders").alias("asin_amazon_orders"),
...@@ -415,30 +544,101 @@ class DwtBsTop100(Templates): ...@@ -415,30 +544,101 @@ class DwtBsTop100(Templates):
F.avg("asin_rating").alias("asin_rating"), F.avg("asin_rating").alias("asin_rating"),
F.avg("asin_bs_cate_1_rank").alias("asin_bs_cate_1_rank"), F.avg("asin_bs_cate_1_rank").alias("asin_bs_cate_1_rank"),
F.avg("asin_total_comments").alias("asin_total_comments"), F.avg("asin_total_comments").alias("asin_total_comments"),
# 类目下海运/空运平均毛利润率 利润率为空的产品不计入
# 只对有利润率数据的 ASIN 求平均。整组全空 → avg 返回 null
F.avg("asin_ocean_freight_gross_margin").alias("asin_ocean_freight_gross_margin_avg"),
F.avg("asin_air_freight_gross_margin").alias("asin_air_freight_gross_margin_avg"),
] ]
# 区间利润率(价格+月销):对 24 个带闸列各求 F.avg = 该分类各区间段的平均利润率(段外行是 null 自动跳过)。
# 别名保持与带闸列同名(不加 _avg),避免多套命名;下游 round/落表按这批列名统一处理。
for pbc in self.band_margin_cols:
aggs.append(F.avg(F.col(pbc)).alias(pbc))
# 打包段):对 7 个带闸月销列各求 F.sum = 该分类该打包段的月销合计(段外 null 不计)。
# 这里只求和;占比(÷品类总月销)在 calculate_column 里算。别名同名。
for pkc in self.pkg_band_cols:
aggs.append(F.sum(F.col(pkc)).alias(pkc))
#添加 rate 列的聚合 #添加 rate 列的聚合
for rc in rate_cols: for rc in rate_cols:
# aggs.append(F.sum(F.col(rc)).alias(rc)) # aggs.append(F.sum(F.col(rc)).alias(rc))
aggs.append(F.sum(F.col(rc).cast("int")).alias(rc)) aggs.append(F.sum(F.col(rc).cast("int")).alias(rc))
if level > 1 and "asin_bs_cate_parent_id_join" in self.df_nodes_num_concat2.columns: # 注:parent_id_join 已作为分组键(见上方),不再用 F.min 聚合——每组本就只有一个路径值
aggs.append(F.min("asin_bs_cate_parent_id_join").alias("asin_bs_cate_parent_id_join"))
df_res = df_level.groupBy(group_keys).agg(*aggs) df_res = df_level.groupBy(group_keys).agg(*aggs)
# ============================================================
# 【集中度 CR10】:每个分类算 top10 实体销量 ÷ 品类总月销。
# 前面都是"整体聚合"(sum/avg),集中度要"分类内排序取前10再求和",是窗口排序,
# 必须在按 group_keys(当前层完整路径)分好组后、rename 之前,用 Window.partitionBy(group_keys) 做。
# 【月销要 cast double】asin_amazon_orders 是 string,orderBy 不转会按字典序("9">"1000")排错;
# sum 也转 double,与分母 asin_bsr_orders_sum(Spark 对 string 求 sum 亦转 double)口径一致。
# 【null 处理】月销 null → desc_nulls_last 排最后、不进 top10,sum 也跳过;
# 品牌/卖家 null → 不是一个真实品牌/卖家,先 filter 掉再排名(它们的销量仍在分母品类总月销里,
# 故品牌/卖家集中度可能 <1,口径=A:分母统一品类总月销,见设计确认)。
# 【不足10个】全取 → 商品集中度=1、品牌/卖家≤1,正常。
# 三个 top10 求和列 _c_asin/_c_brand/_c_seller 先 join 回 df_res,占比在 calculate_column 里 ÷ 总月销。
# ============================================================
orders_d = F.col("asin_amazon_orders").cast("double")
# ===== 原集中度逻辑(单窗口全量排序,level=1 大类几百万行挤一个 task → 倾斜卡住)。保留备查,换成下面的加盐版 =====
# # 商品集中度:asin 直接按月销排序取 top10
# w_asin = Window.partitionBy(*group_keys).orderBy(orders_d.desc_nulls_last())
# df_c_asin = df_level.withColumn("_rn_c", F.row_number().over(w_asin)) \
# .filter(F.col("_rn_c") <= 10) \
# .groupBy(group_keys).agg(F.sum(orders_d).alias("_c_asin_sales"))
# df_res = df_res.join(df_c_asin, on=group_keys, how="left")
# # 品牌集中度:先按(路径,品牌)汇总月销(排除无品牌),再排序取 top10 品牌求和
# df_brand_sales = df_level.filter(F.col("asin_brand_name").isNotNull() & (F.col("asin_brand_name") != "")) \
# .groupBy(*group_keys, "asin_brand_name").agg(F.sum(orders_d).alias("_brand_sales"))
# w_brand = Window.partitionBy(*group_keys).orderBy(F.col("_brand_sales").desc_nulls_last())
# df_c_brand = df_brand_sales.withColumn("_rn_c", F.row_number().over(w_brand)) \
# .filter(F.col("_rn_c") <= 10) \
# .groupBy(group_keys).agg(F.sum("_brand_sales").alias("_c_brand_sales"))
# df_res = df_res.join(df_c_brand, on=group_keys, how="left")
# # 卖家集中度:同品牌,把 asin_brand_name 换成 account_id(排除无卖家)
# df_seller_sales = df_level.filter(F.col("account_id").isNotNull() & (F.col("account_id") != "")) \
# .groupBy(*group_keys, "account_id").agg(F.sum(orders_d).alias("_seller_sales"))
# w_seller = Window.partitionBy(*group_keys).orderBy(F.col("_seller_sales").desc_nulls_last())
# df_c_seller = df_seller_sales.withColumn("_rn_c", F.row_number().over(w_seller)) \
# .filter(F.col("_rn_c") <= 10) \
# .groupBy(group_keys).agg(F.sum("_seller_sales").alias("_c_seller_sales"))
# df_res = df_res.join(df_c_seller, on=group_keys, how="left")
# ===== 新:两段式加盐取 top10(见 _cr10_top_sum),结果等价、消除倾斜 =====
# 先把月销 cast 成一列,三处集中度复用
df_level_o = df_level.withColumn("_orders_d", orders_d)
# 商品集中度:每行一个 asin,按月销取 top10 asin 求和
df_c_asin = self._cr10_top_sum(df_level_o, group_keys, "_orders_d", "_c_asin_sales")
df_res = df_res.join(df_c_asin, on=group_keys, how="left")
# 品牌集中度:先按(路径,品牌)汇总月销(排除无品牌),再取 top10 品牌求和
df_brand_sales = df_level_o.filter(F.col("asin_brand_name").isNotNull() & (F.col("asin_brand_name") != "")) \
.groupBy(*group_keys, "asin_brand_name").agg(F.sum("_orders_d").alias("_brand_sales"))
df_c_brand = self._cr10_top_sum(df_brand_sales, group_keys, "_brand_sales", "_c_brand_sales")
df_res = df_res.join(df_c_brand, on=group_keys, how="left")
# 卖家集中度:把品牌换成 account_id(排除无卖家)
df_seller_sales = df_level_o.filter(F.col("account_id").isNotNull() & (F.col("account_id") != "")) \
.groupBy(*group_keys, "account_id").agg(F.sum("_orders_d").alias("_seller_sales"))
df_c_seller = self._cr10_top_sum(df_seller_sales, group_keys, "_seller_sales", "_c_seller_sales")
df_res = df_res.join(df_c_seller, on=group_keys, how="left")
df_res = df_res.withColumnRenamed(curr_col, "asin_bs_cate_current_id") df_res = df_res.withColumnRenamed(curr_col, "asin_bs_cate_current_id")
if level == 1: if level == 1:
df_res = df_res.withColumn("asin_bs_cate_parent_id", F.lit("0")) df_res = df_res.withColumn("asin_bs_cate_parent_id", F.lit("0"))
df_res = df_res.withColumn("asin_bs_cate_parent_id_join", F.lit("0")) df_res = df_res.withColumn("asin_bs_cate_parent_id_join", F.lit("0"))
else: else:
df_res = df_res.withColumnRenamed(parent_col, "asin_bs_cate_parent_id") # 面包屑 = 0&&&&一级&&&&…&&&&直接父级,用"一级..上一层"祖先列现场拼。
if "asin_bs_cate_parent_id_join" in df_res.columns: # 必须在 rename parent_col 之前拼(此时 asin_bs_cate_{level-1}_id 还在),拼完再 rename。
ancestor_join_cols = [f"asin_bs_cate_{i}_id" for i in range(1, level)] # 1..level-1
df_res = df_res.withColumn( df_res = df_res.withColumn(
"asin_bs_cate_parent_id_join", "asin_bs_cate_parent_id_join",
F.concat(F.lit("0&&&&"), F.col("asin_bs_cate_parent_id_join")) F.concat(F.lit("0&&&&"),
F.concat_ws("&&&&", *[F.col(c).cast("string") for c in ancestor_join_cols]))
) )
else: df_res = df_res.withColumnRenamed(parent_col, "asin_bs_cate_parent_id")
df_res = df_res.withColumn("asin_bs_cate_parent_id_join", F.lit("0&&&&&nan")) # 删掉分组用剩的中间祖先列(1..level-2),只保留 current_id/parent_id/parent_id_join,
# 否则它们会随 unionByName 混进最终结果,污染列结构。
for _c in [f"asin_bs_cate_{i}_id" for i in range(1, level - 1)]:
if _c in df_res.columns:
df_res = df_res.drop(_c)
df_res = df_res.withColumn("level", F.lit(level)) df_res = df_res.withColumn("level", F.lit(level))
...@@ -446,7 +646,8 @@ class DwtBsTop100(Templates): ...@@ -446,7 +646,8 @@ class DwtBsTop100(Templates):
# print(f"Level {level} Count: {df_res.count()}") # print(f"Level {level} Count: {df_res.count()}")
df_level_list.append(df_res) df_level_list.append(df_res)
self.df_level_concat = reduce(lambda a, b: a.unionByName(b, allowMissingColumns=True), df_level_list) self.df_level_concat = reduce(lambda a, b: a.unionByName(b, allowMissingColumns=True), df_level_list)
self.df_level_concat.cache() # join 后即被重新赋值、从未被用的死缓存,去掉
# self.df_level_concat.cache()
# print("self.df_level_concat",self.df_level_concat.count()) # print("self.df_level_concat",self.df_level_concat.count())
# print(self.df_level_concat.columns) # print(self.df_level_concat.columns)
# 增加分类名称 # 增加分类名称
...@@ -454,7 +655,10 @@ class DwtBsTop100(Templates): ...@@ -454,7 +655,10 @@ class DwtBsTop100(Templates):
print("增加分类名称") print("增加分类名称")
df_cate_name = self.df_cate_bs_data.select(*['asin_bs_cate_current_id', 'cur_category_name','_cate_order']) df_cate_name = self.df_cate_bs_data.select(*['asin_bs_cate_current_id', 'cur_category_name','_cate_order'])
window_spec = Window.partitionBy('asin_bs_cate_current_id').orderBy(F.col("_cate_order").asc_nulls_last()) window_spec = Window.partitionBy('asin_bs_cate_current_id').orderBy(F.col("_cate_order").asc_nulls_last())
df_cate_name = df_cate_name.withColumn("_rn",F.row_number().over(window_spec)).filter(F.col("_rn") == 1).drop("_rn") # _cate_order 仅用于上面 window 去重排序,用完 drop
df_cate_name = df_cate_name.withColumn("_rn",F.row_number().over(window_spec)).filter(F.col("_rn") == 1).drop("_rn", "_cate_order")
# inner join 补分类名:current_id 不在 us_bs_category 的行会被丢。实测仅丢 2 个非标准一级
# (private-brands / book-special-features 之类,树里无对应 category_id、页面也不查) → 保留 inner,无害。
self.df_level_concat =self.df_level_concat.join(df_cate_name, on='asin_bs_cate_current_id',how="inner") self.df_level_concat =self.df_level_concat.join(df_cate_name, on='asin_bs_cate_current_id',how="inner")
print("增加分类名称 join后的数据") print("增加分类名称 join后的数据")
# self.df_level_concat.show(10,truncate=False ) # self.df_level_concat.show(10,truncate=False )
...@@ -466,6 +670,8 @@ class DwtBsTop100(Templates): ...@@ -466,6 +670,8 @@ class DwtBsTop100(Templates):
print("handle_bs_level_calculate_column 调用 字段占比计算") print("handle_bs_level_calculate_column 调用 字段占比计算")
rename_columns = { rename_columns = {
"asin_amazon_orders": "asin_amazon_orders_sum", "asin_amazon_orders": "asin_amazon_orders_sum",
# asin_bsr_orders 底层已在 handle_bs_add_column 重定向为前台月销 →
# asin_bsr_orders_sum 实为"前台月销总销量",数值等同 asin_amazon_orders_sum,页面不显示此列
"asin_bsr_orders": "asin_bsr_orders_sum", "asin_bsr_orders": "asin_bsr_orders_sum",
"asin_bsr_orders_new": "asin_bsr_orders_sum_new", "asin_bsr_orders_new": "asin_bsr_orders_sum_new",
"asin_bsr_orders_brand": "asin_bsr_orders_sum_brand", "asin_bsr_orders_brand": "asin_bsr_orders_sum_brand",
...@@ -488,9 +694,12 @@ class DwtBsTop100(Templates): ...@@ -488,9 +694,12 @@ class DwtBsTop100(Templates):
"asin_count_rate_new", "asin_count_rate_new",
F.col("asin_count_new") / F.col("asin_count") F.col("asin_count_new") / F.col("asin_count")
) )
# 分母=当前分类总销量,为 0/null(该分类无销量)→ -1 占位(无法计算)
# 分子=新品销量 sum,null 表示无可加行 → coalesce 按真实 0 计
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_sum_rate_new", "asin_bsr_orders_sum_rate_new",
F.col("asin_bsr_orders_sum_new") / F.col("asin_bsr_orders_sum") F.when(F.col("asin_bsr_orders_sum").isNull() | (F.col("asin_bsr_orders_sum") == 0), F.lit(-1))
.otherwise(F.coalesce(F.col("asin_bsr_orders_sum_new"), F.lit(0)) / F.col("asin_bsr_orders_sum"))
) )
# 品牌占比 # 品牌占比
...@@ -498,9 +707,11 @@ class DwtBsTop100(Templates): ...@@ -498,9 +707,11 @@ class DwtBsTop100(Templates):
"asin_count_rate_brand", "asin_count_rate_brand",
F.col("asin_count_brand") / F.col("asin_count") F.col("asin_count_brand") / F.col("asin_count")
) )
# 分母=当前分类总销量,为 0/null → -1 占位;分子=品牌销量 sum,null 按真实 0 计
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_sum_rate_brand", "asin_bsr_orders_sum_rate_brand",
F.col("asin_bsr_orders_sum_brand") / F.col("asin_bsr_orders_sum") F.when(F.col("asin_bsr_orders_sum").isNull() | (F.col("asin_bsr_orders_sum") == 0), F.lit(-1))
.otherwise(F.coalesce(F.col("asin_bsr_orders_sum_brand"), F.lit(0)) / F.col("asin_bsr_orders_sum"))
) )
# 平均销售额 # 平均销售额
...@@ -508,13 +719,17 @@ class DwtBsTop100(Templates): ...@@ -508,13 +719,17 @@ class DwtBsTop100(Templates):
"asin_bsr_orders_sales_mean", "asin_bsr_orders_sales_mean",
F.col("asin_bsr_orders_sales_sum") / F.col("asin_count") F.col("asin_bsr_orders_sales_sum") / F.col("asin_count")
) )
# 分母=新品数,为 0/null(无新品)→ -1 占位;分子=新品销售额 sum,null 按真实 0 计
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_sales_mean_new", "asin_bsr_orders_sales_mean_new",
F.col("asin_bsr_orders_sales_sum_new") / F.col("asin_count_new") F.when(F.col("asin_count_new").isNull() | (F.col("asin_count_new") == 0), F.lit(-1))
.otherwise(F.coalesce(F.col("asin_bsr_orders_sales_sum_new"), F.lit(0)) / F.col("asin_count_new"))
) )
# 分母=品牌数,为 0/null(无品牌)→ -1 占位;分子=品牌销售额 sum,null 按真实 0 计
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_sales_mean_brand", "asin_bsr_orders_sales_mean_brand",
F.col("asin_bsr_orders_sales_sum_brand") / F.col("asin_count_brand") F.when(F.col("asin_count_brand").isNull() | (F.col("asin_count_brand") == 0), F.lit(-1))
.otherwise(F.coalesce(F.col("asin_bsr_orders_sales_sum_brand"), F.lit(0)) / F.col("asin_count_brand"))
) )
# 平均月销量 # 平均月销量
...@@ -522,13 +737,17 @@ class DwtBsTop100(Templates): ...@@ -522,13 +737,17 @@ class DwtBsTop100(Templates):
"asin_bsr_orders_mean", "asin_bsr_orders_mean",
F.col("asin_bsr_orders_sum") / F.col("asin_count") F.col("asin_bsr_orders_sum") / F.col("asin_count")
) )
# 分母=新品数,为 0/null(无新品)→ -1 占位;分子=新品销量 sum,null 按真实 0 计
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_mean_new", "asin_bsr_orders_mean_new",
F.col("asin_bsr_orders_sum_new") / F.col("asin_count_new") F.when(F.col("asin_count_new").isNull() | (F.col("asin_count_new") == 0), F.lit(-1))
.otherwise(F.coalesce(F.col("asin_bsr_orders_sum_new"), F.lit(0)) / F.col("asin_count_new"))
) )
# 分母=品牌数,为 0/null(无品牌)→ -1 占位;分子=品牌销量 sum,null 按真实 0 计
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_mean_brand", "asin_bsr_orders_mean_brand",
F.col("asin_bsr_orders_sum_brand") / F.col("asin_count_brand") F.when(F.col("asin_count_brand").isNull() | (F.col("asin_count_brand") == 0), F.lit(-1))
.otherwise(F.coalesce(F.col("asin_bsr_orders_sum_brand"), F.lit(0)) / F.col("asin_count_brand"))
) )
# ========== 3. 保留小数位数 ========== # ========== 3. 保留小数位数 ==========
...@@ -548,11 +767,26 @@ class DwtBsTop100(Templates): ...@@ -548,11 +767,26 @@ class DwtBsTop100(Templates):
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_count_brand", F.round("asin_count_brand", 0) "asin_count_brand", F.round("asin_count_brand", 0)
) )
# 月均销量属"件数"→取整(2026-07-02 补 round(0),统一口径:销量=整数)。
# 这三列上面是除法(sum/count)得 double,之前漏了取整会带小数;mean_new/mean_brand 的 -1 占位 round 后仍是 -1,不受影响。
self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_mean", F.round("asin_bsr_orders_mean", 0)
)
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bs_cate_1_rank_mean", F.round("asin_bs_cate_1_rank_mean", 0) "asin_bsr_orders_mean_new", F.round("asin_bsr_orders_mean_new", 0)
) )
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_total_comments_mean", F.round("asin_total_comments_mean", 0) "asin_bsr_orders_mean_brand", F.round("asin_bsr_orders_mean_brand", 0)
)
# 均值列 F.avg 遇"整组该字段全 null"(如某冷门分类下所有 asin 都无排名)会返回 NULL;
# 按项目约定 Java 不存 NULL → coalesce 填 -1 占位(排名/价格/评分均无负数,-1 不会和真实值混淆,Java 转 null)。
# 已实测该分区 rank/price/rating 均存在整组全空的分类(comments 恒非空故不守)。
self.df_level_concat = self.df_level_concat.withColumn(
"asin_bs_cate_1_rank_mean", F.coalesce(F.round("asin_bs_cate_1_rank_mean", 0), F.lit(-1))
)
# comments 当前上游 fillna(0) 恒非空,avg 不会 NULL;仍一并兜底,防上游哪天去掉 fillna(-1 占位同上)
self.df_level_concat = self.df_level_concat.withColumn(
"asin_total_comments_mean", F.coalesce(F.round("asin_total_comments_mean", 0), F.lit(-1))
) )
# round(2) 保留2位小数 # round(2) 保留2位小数
...@@ -574,11 +808,12 @@ class DwtBsTop100(Templates): ...@@ -574,11 +808,12 @@ class DwtBsTop100(Templates):
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_sales_mean_brand", F.round("asin_bsr_orders_sales_mean_brand", 2) "asin_bsr_orders_sales_mean_brand", F.round("asin_bsr_orders_sales_mean_brand", 2)
) )
# 同上:价格/评分均值整组全空 → NULL,coalesce 填 -1 占位
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_price_mean", F.round("asin_price_mean", 2) "asin_price_mean", F.coalesce(F.round("asin_price_mean", 2), F.lit(-1))
) )
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_rating_mean", F.round("asin_rating_mean", 2) "asin_rating_mean", F.coalesce(F.round("asin_rating_mean", 2), F.lit(-1))
) )
# round(4) 保留4位小数 # round(4) 保留4位小数
...@@ -594,68 +829,89 @@ class DwtBsTop100(Templates): ...@@ -594,68 +829,89 @@ class DwtBsTop100(Templates):
self.df_level_concat = self.df_level_concat.withColumn( self.df_level_concat = self.df_level_concat.withColumn(
"asin_bsr_orders_sum_rate_brand", F.round("asin_bsr_orders_sum_rate_brand", 4) "asin_bsr_orders_sum_rate_brand", F.round("asin_bsr_orders_sum_rate_brand", 4)
) )
# 海运/空运平均毛利润率:比例,round(4)。整组全空(该分类下无任一 ASIN 有利润率) → avg 返回 null。
# 【和其它均值列不同】此列无值时直接落 null(不填 -1/-1000 占位)——已与后端确认利润率列可收 null。
# 好处:null 天然不参与 同比环比的算术;坏处/注意:落表 save_data 时这批利润率列不能 na.fill。
self.df_level_concat = self.df_level_concat.withColumn(
"asin_ocean_freight_gross_margin_avg",
F.round("asin_ocean_freight_gross_margin_avg", 4)
)
self.df_level_concat = self.df_level_concat.withColumn(
"asin_air_freight_gross_margin_avg",
F.round("asin_air_freight_gross_margin_avg", 4)
)
# 区间利润率 24 列(价格12+月销12):同利润率均值,比例 round(4),整段全空 → avg 为 null → 直接落 null(同上)。
for pbc in self.band_margin_cols:
self.df_level_concat = self.df_level_concat.withColumn(
pbc, F.round(pbc, 4)
)
def handle_bs_level_proportion(self): # 打包段月销占比 7 列:段月销合计 ÷ 品类总月销(asin_bsr_orders_sum) = 该段占比。
# 被handle_bs_level_proportion2代替 # 分母=品类总月销,为 0/null(该分类无月销) → -1 占位(无法计算,与其它占比列口径一致,Java 转 null)。
print("handle_bs_level_proportion 调用 计算层级占比") # 分子=段月销 sum,null(该段无月销 asin) → coalesce 按真实 0 计。生成 _rate 后缀的最终占比列,round(4)。
row_list = self.df_level_concat.select("level").distinct().collect() # 注:占比列名含 "rate" 无妨——rate_cols 在 handle_bs_level 已用完,此处是更后面的 calculate_column。
level_list = sorted( [row.level for row in row_list]) for pkc in self.pkg_band_cols:
print("level_list",level_list) rate_col = pkc.replace("_orders_sum", "_orders_rate")
self.df_level_concat.cache() # 先缓存输入数据,防止重复计算 self.df_level_concat = self.df_level_concat.withColumn(
df_level_rate_list = [] rate_col,
for level in level_list: F.when(F.col("asin_bsr_orders_sum").isNull() | (F.col("asin_bsr_orders_sum") == 0), F.lit(-1))
df_cur_lever = self.df_level_concat.filter(self.df_level_concat.level == level).select(*['asin_bs_cate_current_id', 'asin_bs_cate_parent_id', .otherwise(F.round(F.coalesce(F.col(pkc), F.lit(0)) / F.col("asin_bsr_orders_sum"), 4))
'asin_bsr_orders_sum', 'level']) )
# print(level, df_cur_lever.count(),len(df_cur_lever.columns))
df_level_parent = df_cur_lever.groupBy("asin_bs_cate_parent_id").agg( # 集中度 CR10 3 列:top10 实体销量(_c_*_sales,handle_bs_level 里算好) ÷ 品类总月销(asin_bsr_orders_sum)。
F.sum("asin_bsr_orders_sum").alias("asin_bsr_orders_sum")) # 分母=品类总月销,为 0/null → -1 占位(同其它占比列口径)。分子 null(该分类无对应实体) → coalesce 0。round(4)。
df_level_parent = df_level_parent.withColumnRenamed("asin_bsr_orders_sum","asin_bsr_orders_sum_parent") # 商品集中度分母含全部月销、分子是 top10 asin → 通常接近满值;品牌/卖家分子排除了无品牌/无卖家 → 可能 <1(口径A)。
for raw_col, out_col in [
df_cur_lever = df_cur_lever.join(df_level_parent, on=["asin_bs_cate_parent_id"],how="inner") ("_c_asin_sales", "asin_concentration"),
("_c_brand_sales", "brand_concentration"),
df_cur_lever = df_cur_lever.withColumn("asin_bsr_orders_sum_rate",F.col("asin_bsr_orders_sum") / F.col("asin_bsr_orders_sum_parent")) ("_c_seller_sales", "seller_concentration"),
df_cur_lever = df_cur_lever.withColumn("asin_bsr_orders_sum_rate",F.round("asin_bsr_orders_sum_rate",4)) ]:
df_level_rate_list.append(df_cur_lever) self.df_level_concat = self.df_level_concat.withColumn(
print("循环完毕") out_col,
df_level_rate_concat = reduce(lambda a, b: a.unionByName(b, allowMissingColumns=True), df_level_rate_list) F.when(F.col("asin_bsr_orders_sum").isNull() | (F.col("asin_bsr_orders_sum") == 0), F.lit(-1))
print("df_level_rate_concat 打印") .otherwise(F.round(F.coalesce(F.col(raw_col), F.lit(0)) / F.col("asin_bsr_orders_sum"), 4))
# df_level_rate_concat.show(10,truncate=False) )
# print(df_level_rate_concat.count())
# print(df_level_rate_concat.columns) # 注:handle_bs_level_proportion(非2版,2026-07-02 删除)——groupBy+inner join 逐层直译版,慢且 inner join 可能丢行,
# self.df_save生成 # 由下方 handle_bs_level_proportion2(窗口函数一次算出父级总销量,等价且更快、不丢行)取代。
print("self.df_save 保存对象生成") # 如需找回:开发母本 spark_hjm/bs_top100/dwt_bs_top100.py 或 git 历史。
self.df_save = self.df_level_concat.join(df_level_rate_concat, how="inner", on=['asin_bs_cate_current_id', 'asin_bs_cate_parent_id', 'asin_bsr_orders_sum', 'level'])
print("df_save的值为")
self.df_save.show(10,truncate=False)
print(self.df_save.count())
print(self.df_save.columns)
def handle_bs_level_proportion2(self): def handle_bs_level_proportion2(self):
print("handle_bs_level_proportion 调用 计算层级占比") print("handle_bs_level_proportion 调用 计算层级占比")
self.df_level_concat.cache() # 下游到写出是线性链只读一次,去掉
# 2. 定义窗口:按父分类ID (asin_bs_cate_parent_id) 分组 # self.df_level_concat.cache()
# 这一步相当于原来的: groupBy("asin_bs_cate_parent_id") # 2. 定义窗口:按"完整父级路径 parent_id_join"分组(不是只按直接父级 id)
window_spec = Window.partitionBy("asin_bs_cate_parent_id") # 直接父级 id 在多归属时不唯一——如 Cables 同时挂在 electronics、pc 下,只按 Cables 分组会把
# 两个一级支的子级销量全加进分母 → 占比偏小、跨一级混算(分母 220、实际同支应为 160)。
# parent_id_join 形如 "0&&&&electronics&&&&Cables",唯一标识"某条路径下的这个父级",只有同路径
# 的兄弟才进同一分母。level=1(值"0")、level=2(值"0&&&&一级")的分组与按直接父级 id 完全等价,
# 只有 level≥3 的多归属分类被修正——和 handle_bs_level 的拆行改动对齐。
window_spec = Window.partitionBy("asin_bs_cate_parent_id_join")
# 3. 计算父节点总订单数 # 3. 计算父节点总订单数
# 这一步相当于原来的: agg(sum("asin_bsr_orders_sum")) + join # 这一步相当于原来的: agg(sum("asin_bsr_orders_sum")) + join
# 效果:每一行都会多出一列 'asin_bsr_orders_sum_parent', # 效果:每一行都会多出一列 'asin_bsr_orders_sum_parent',
# 这一列的值等于该行所属父分类下所有行的 orders_sum 之和。 # 这一列的值等于该行所属父分类下所有行的 orders_sum 之和。
# 父级总销量也属"件数"→取整(2026-07-02 补 round(0),统一口径:销量=整数)。
# 它是占比分母(本级/父级),取整后与分子 asin_bsr_orders_sum 同为整数,口径一致。
self.df_save = self.df_level_concat.withColumn( self.df_save = self.df_level_concat.withColumn(
"asin_bsr_orders_sum_parent", "asin_bsr_orders_sum_parent",
F.sum("asin_bsr_orders_sum").over(window_spec) F.round(F.sum("asin_bsr_orders_sum").over(window_spec), 0)
) )
self.df_save.cache() # 下游到写出是线性链只读一次,去掉
# self.df_save.cache()
# 4. 计算占比 # 4. 计算占比
# 逻辑:自身订单 / 父节点总订单 # 逻辑:自身订单 / 父节点总订单
# 分母=父级总销量,为 0/null(整支无销量)→ -1 占位(无法计算,Java 转 null),与其余占比列口径统一
self.df_save = self.df_save.withColumn( self.df_save = self.df_save.withColumn(
"asin_bsr_orders_sum_rate", "asin_bsr_orders_sum_rate",
F.when(F.col("asin_bsr_orders_sum_parent") != 0, F.when(F.col("asin_bsr_orders_sum_parent") != 0,
F.round(F.col("asin_bsr_orders_sum") / F.col("asin_bsr_orders_sum_parent"), 4) F.round(F.col("asin_bsr_orders_sum") / F.col("asin_bsr_orders_sum_parent"), 4)
).otherwise(0) ).otherwise(F.lit(-1))
) )
# 释放缓存 # 释放缓存
self.df_level_concat.unpersist() # 上面已不缓存 df_level_concat,unpersist 无对象,去掉
# self.df_level_concat.unpersist()
def handle_bs_proportion(self): def handle_bs_proportion(self):
...@@ -668,12 +924,18 @@ class DwtBsTop100(Templates): ...@@ -668,12 +924,18 @@ class DwtBsTop100(Templates):
"launch_time_type6_rate", "launch_time_type7_rate", "launch_time_type6_rate", "launch_time_type7_rate",
"asin_buy_box_seller_type0_rate", "asin_buy_box_seller_type1_rate", "asin_buy_box_seller_type0_rate", "asin_buy_box_seller_type1_rate",
"asin_buy_box_seller_type2_rate", "asin_buy_box_seller_type3_rate", "asin_buy_box_seller_type2_rate", "asin_buy_box_seller_type3_rate",
"asin_buy_box_seller_type4_rate" "asin_buy_box_seller_type4_rate",
# 卖家所属地 3 列:与上架时间/卖家类型同为"数量占比",同样 ÷ asin_count。
# 三列 null-safe 互斥穷尽(见 handle_bs_categorical_data),计数和=asin_count → 三占比相加=1。
"seller_country_us_rate", "seller_country_cn_rate", "seller_country_other_rate",
# 月销50+/100+/200+ 3 列:达标 asin 数 ÷ asin_count。累计阈值不互斥、不相加=1。
"asin_orders_ge_50_rate", "asin_orders_ge_100_rate", "asin_orders_ge_200_rate"
] ]
# 使用循环批量更新列 # 使用循环批量更新列:桶计数 ÷ 该分类总数量 = 占比,round(4) 对齐目标表 double(10,4)
# 分母 asin_count 恒 ≥1、分子(桶计数 sum)上游无 null 恒非空 → 除法不会 null,只需取整,无需 -1 守卫
for col_name in rate_cols: for col_name in rate_cols:
self.df_save = self.df_save.withColumn(col_name, F.col(col_name) / F.col("asin_count")) self.df_save = self.df_save.withColumn(col_name, F.round(F.col(col_name) / F.col("asin_count"), 4))
def handle_data(self): def handle_data(self):
self.handle_bs_redirect_flag() self.handle_bs_redirect_flag()
...@@ -682,328 +944,88 @@ class DwtBsTop100(Templates): ...@@ -682,328 +944,88 @@ class DwtBsTop100(Templates):
self.handle_bs_categorical_data() self.handle_bs_categorical_data()
self.handle_bs_asin_deduplicated() self.handle_bs_asin_deduplicated()
self.handle_bs_add_column() self.handle_bs_add_column()
self.handle_bs_nodes_num_concat()
self.handle_bs_nodes_num_concat2() self.handle_bs_nodes_num_concat2()
self.handle_bs_level() self.handle_bs_level()
self.handle_bs_level_calculate_column() self.handle_bs_level_calculate_column()
# self.handle_bs_level_proportion()
self.handle_bs_level_proportion2() self.handle_bs_level_proportion2()
self.handle_bs_proportion() self.handle_bs_proportion()
self.handle_bs_to_csv_ceil() # 不转csv的话可以不用加 为了对齐mysql格式(mysql设置了约束 数据会清洗 截断)
# print("df_save 计算完成")
# self.df_save.show(10, truncate=False)
print(f" df_save总行数: {self.df_save.count()}")
# print(self.df_save.columns)
def handle_bs_to_csv_ceil(self):
df = self.df_save
# ========== 所有数值字段先填充NULL为0 ==========
all_numeric_cols = [
# int类型
"asin_amazon_orders_sum", "asin_bsr_orders_sum", "asin_bsr_orders_sum_new",
"asin_bsr_orders_sum_brand", "asin_bsr_orders_sales_sum",
"asin_bsr_orders_sales_sum_new", "asin_bsr_orders_sales_sum_brand",
"asin_count_new", "asin_count_brand", "asin_count",
"asin_total_comments_mean", "level",
"asin_bsr_orders_mean_new", "asin_bsr_orders_sales_mean_new",
"asin_bsr_orders_sales_mean_brand", "asin_bsr_orders_mean_brand",
"asin_bs_cate_1_rank_mean",
# double类型
"asin_bsr_orders_mean", "asin_bsr_orders_sum_parent",
"asin_bsr_orders_sales_mean", "asin_price_mean", "asin_rating_mean",
"asin_bsr_orders_sum_rate", "asin_bsr_orders_sum_rate_new",
"asin_bsr_orders_sum_rate_brand", "asin_count_rate_new", "asin_count_rate_brand",
"launch_time_type0_rate", "launch_time_type1_rate", "launch_time_type2_rate",
"launch_time_type3_rate", "launch_time_type4_rate", "launch_time_type5_rate",
"launch_time_type6_rate", "launch_time_type7_rate",
"asin_buy_box_seller_type0_rate", "asin_buy_box_seller_type1_rate",
"asin_buy_box_seller_type2_rate", "asin_buy_box_seller_type3_rate",
"asin_buy_box_seller_type4_rate",
]
for c in all_numeric_cols:
if c in df.columns:
df = df.withColumn(c, F.coalesce(F.col(c), F.lit(0)))
# ========== 1. 纯int类型字段(输出纯整数,如 "2480")==========
# 这些字段在SQL中是int,且通常不会有NULL值
pure_int_cols = [
"asin_amazon_orders_sum",
"asin_bsr_orders_sum",
"asin_bsr_orders_sum_new",
"asin_bsr_orders_sum_brand",
"asin_bsr_orders_sales_sum",
"asin_bsr_orders_sales_sum_new",
"asin_bsr_orders_sales_sum_brand",
"asin_count_new",
"asin_count_brand",
"asin_count",
"asin_total_comments_mean",
"level",
]
for c in pure_int_cols:
if c in df.columns:
df = df.withColumn(c, F.round(F.col(c), 0).cast("long"))
# ========== 2. int类型但pandas读取后会带.0的字段(有NULL值)==========
# 这些字段输出为 "2480.0" 格式,与pandas读取SQL后的格式一致
int_with_decimal_cols = [
"asin_bsr_orders_mean_new",
"asin_bsr_orders_sales_mean_new",
"asin_bsr_orders_sales_mean_brand",
"asin_bsr_orders_mean_brand",
"asin_bs_cate_1_rank_mean",
]
for c in int_with_decimal_cols:
if c in df.columns:
df = df.withColumn(c, F.round(F.col(c), 0).cast("double"))
# ========== 3. double(20,0) - 整数精度的double ==========
double_int_cols = [
"asin_bsr_orders_mean",
"asin_bsr_orders_sum_parent",
]
for c in double_int_cols:
if c in df.columns:
df = df.withColumn(c, F.round(F.col(c), 0))
# ========== 4. double(10,2) 字段 ==========
double2_cols = [
"asin_price_mean",
"asin_rating_mean",
"asin_bsr_orders_sales_mean",
"asin_bsr_orders_sum_rate_new",
"asin_bsr_orders_sum_rate_brand",
]
for c in double2_cols:
if c in df.columns:
df = df.withColumn(c, F.round(F.col(c), 2))
# ========== 5. double(10,4) 比例字段 ==========
double4_cols = [
"asin_bsr_orders_sum_rate",
"asin_count_rate_new",
"asin_count_rate_brand",
"launch_time_type0_rate",
"launch_time_type1_rate",
"launch_time_type2_rate",
"launch_time_type3_rate",
"launch_time_type4_rate",
"launch_time_type5_rate",
"launch_time_type6_rate",
"launch_time_type7_rate",
"asin_buy_box_seller_type0_rate",
"asin_buy_box_seller_type1_rate",
"asin_buy_box_seller_type2_rate",
"asin_buy_box_seller_type3_rate",
"asin_buy_box_seller_type4_rate",
]
for c in double4_cols:
if c in df.columns:
df = df.withColumn(c, F.round(F.col(c), 4))
self.df_save = df
if "_cate_order" in self.df_save.columns:
self.df_save = self.df_save.drop("_cate_order")
# self.df_save = self.df_save.drop("_cate_order") def save_data(self):
def save_pandas_to_csv(self): print("save_data 调用:落 Hive dwt_bs_top100")
"""
将Spark DataFrame转为Pandas并导出CSV,格式与MySQL读取后的pandas格式一致
"""
import pandas as pd
import numpy as np
p_df = self.df_save.toPandas()
# ========== 字段分类定义(根据SQL表结构和pandas读取行为)==========
# 纯int类型:输出纯整数,如 "1234"(不带.0)
pure_int_cols = [
"asin_amazon_orders_sum",
"asin_bsr_orders_sum",
"asin_bsr_orders_sum_new",
"asin_bsr_orders_sum_brand",
"asin_bsr_orders_sales_sum",
"asin_bsr_orders_sales_sum_new",
"asin_bsr_orders_sales_sum_brand",
"asin_count_new",
"asin_count_brand",
"asin_count",
"asin_total_comments_mean",
"level",
]
# int类型但pandas读取后带.0的字段:输出 "1234.0" 格式
int_with_decimal_cols = [
"asin_bsr_orders_mean_new",
"asin_bsr_orders_sales_mean_new",
"asin_bsr_orders_sales_mean_brand",
"asin_bsr_orders_mean_brand",
"asin_bs_cate_1_rank_mean",
]
# double(20,0)类型:输出 "1234.0" 格式
double_int_cols = [
"asin_bsr_orders_mean",
"asin_bsr_orders_sum_parent",
]
# double(10,2)类型:最多2位小数,去尾部0但保留.0
double2_cols = [
"asin_price_mean",
"asin_rating_mean",
"asin_bsr_orders_sales_mean",
"asin_bsr_orders_sum_rate_new",
"asin_bsr_orders_sum_rate_brand",
]
# double(10,4)类型:最多4位小数
double4_cols = [
"asin_bsr_orders_sum_rate",
"asin_count_rate_new",
"asin_count_rate_brand",
"launch_time_type0_rate",
"launch_time_type1_rate",
"launch_time_type2_rate",
"launch_time_type3_rate",
"launch_time_type4_rate",
"launch_time_type5_rate",
"launch_time_type6_rate",
"launch_time_type7_rate",
"asin_buy_box_seller_type0_rate",
"asin_buy_box_seller_type1_rate",
"asin_buy_box_seller_type2_rate",
"asin_buy_box_seller_type3_rate",
"asin_buy_box_seller_type4_rate",
]
# 字符串类型 # ============================================================
str_cols = [ # 44 业务列 + 利润率均值2 + 区间利润率24(self.band_margin_cols)
# + 打包占比7 + 卖家地占比3 + 月销阈值占比3 + 集中度3 = 86 列,再 + 3 分区列。
# ============================================================
business_cols = [
"asin_bs_cate_current_id", "asin_bs_cate_current_id",
"asin_bs_cate_parent_id", "asin_amazon_orders_sum", "asin_bsr_orders_sum", "asin_bsr_orders_sum_new", "asin_bsr_orders_sum_brand",
"asin_bs_cate_parent_id_join", "asin_bsr_orders_sales_sum", "asin_bsr_orders_sales_sum_new", "asin_bsr_orders_sales_sum_brand",
"cur_category_name", "asin_count_new", "asin_count_brand", "asin_count",
"asin_price_mean", "asin_rating_mean", "asin_bs_cate_1_rank_mean", "asin_total_comments_mean",
"launch_time_type0_rate", "launch_time_type1_rate", "launch_time_type2_rate", "launch_time_type3_rate",
"launch_time_type4_rate", "launch_time_type5_rate", "launch_time_type6_rate", "launch_time_type7_rate",
"asin_buy_box_seller_type0_rate", "asin_buy_box_seller_type1_rate", "asin_buy_box_seller_type2_rate",
"asin_buy_box_seller_type3_rate", "asin_buy_box_seller_type4_rate",
"asin_bs_cate_parent_id", "asin_bs_cate_parent_id_join", "level", "cur_category_name",
"asin_count_rate_new", "asin_bsr_orders_sum_rate_new", "asin_count_rate_brand", "asin_bsr_orders_sum_rate_brand",
"asin_bsr_orders_sales_mean", "asin_bsr_orders_sales_mean_new", "asin_bsr_orders_sales_mean_brand",
"asin_bsr_orders_mean", "asin_bsr_orders_mean_new", "asin_bsr_orders_mean_brand",
"asin_bsr_orders_sum_parent", "asin_bsr_orders_sum_rate",
] ]
profit_avg_cols = ["asin_ocean_freight_gross_margin_avg", "asin_air_freight_gross_margin_avg"]
# ========== 格式化函数 ========== pkg_rate_cols = [c.replace("_orders_sum", "_orders_rate") for c in self.pkg_band_cols]
def format_pure_int(v): country_cols = ["seller_country_us_rate", "seller_country_cn_rate", "seller_country_other_rate"]
"""纯int类型:空值为空字符串,否则为纯整数字符串(不带.0)""" orders_ge_cols = ["asin_orders_ge_50_rate", "asin_orders_ge_100_rate", "asin_orders_ge_200_rate"]
if pd.isna(v): conc_cols = ["asin_concentration", "brand_concentration", "seller_concentration"]
return '' final_cols = business_cols + profit_avg_cols + list(self.band_margin_cols) \
try: + pkg_rate_cols + country_cols + orders_ge_cols + conc_cols
return str(int(round(float(v))))
except: # ============================================================
return '' # 【销量/金额 sum 列】某分类全无月销时 sum=null → coalesce(0)(0=无销量,有意义,区别于利润率的"无数据")。
# 【利润率 26 列】保留 null(已与后端确认可收 null,不填占位)。
def format_int_with_decimal(v): # 【均值/占比/集中度列】计算时已填 -1 占位,不动。
"""int类型但需要带.0:如 "1234.0"(与pandas读取SQL后格式一致)""" # 【计数/其它占比】分母 asin_count≥1、分子非空,不会 null,不处理。
if pd.isna(v): # ============================================================
return '' zero_fill = {
try: "asin_amazon_orders_sum", "asin_bsr_orders_sum", "asin_bsr_orders_sum_new", "asin_bsr_orders_sum_brand",
return f"{round(float(v)):.1f}" "asin_bsr_orders_sales_sum", "asin_bsr_orders_sales_sum_new", "asin_bsr_orders_sales_sum_brand",
except: "asin_bsr_orders_mean", "asin_bsr_orders_sum_parent",
return '' # 补漏:这三列同为"分子可能整组 null"却之前未兜底,会给后端落 null(违反不存 null 铁律):
# sales_mean = sales_sum/count,sales_sum 整组月销或价格全 null 时为 null(同 bsr_orders_mean 却漏填);
def format_double_int(v): # count_new = sum(asin_is_new),asin_is_new 上游 dwt_flow_asin 未 na.fill、整组全 null 时为 null;
"""double(20,0)类型:整数但带.0,如 "1234.0" """ # count_rate_new = count_new/count 随之 null。(count_brand 侧 is_brand_label 上游填了 0,安全,无需补)
if pd.isna(v): "asin_bsr_orders_sales_mean", "asin_count_new", "asin_count_rate_new",
return '' }
try: select_exprs = []
return f"{round(float(v)):.1f}" for c in final_cols:
except: if c in zero_fill:
return '' select_exprs.append(F.coalesce(F.col(c), F.lit(0)).alias(c))
def format_double2(v):
"""double(10,2)类型:最多2位小数,去尾部0但整数保留.0"""
if pd.isna(v):
return ''
try:
val = round(float(v), 2)
if val == int(val):
# 整数情况,保留.0 (如 1.0, 0.0)
return f"{int(val)}.0"
else: else:
# 有小数,去尾部0 select_exprs.append(F.col(c))
s = f"{val:.2f}".rstrip('0') # 补 3 个分区列(放最后,与建表 PARTITIONED BY 一致)
return s select_exprs += [
except: F.lit(self.site_name).alias("site_name"),
return '' F.lit(self.date_type).alias("date_type"),
F.lit(self.date_info).alias("date_info"),
def format_double4(v): ]
"""double(10,4)类型:最多4位小数,去尾部0""" df_save = self.df_save.select(*select_exprs)
if pd.isna(v):
return ''
try: partition_dict = {"site_name": self.site_name, "date_type": self.date_type, "date_info": self.date_info}
val = round(float(v), 4) hdfs_path = CommonUtil.build_hdfs_path(self.db_save, partition_dict=partition_dict)
s = f"{val:.4f}".rstrip('0').rstrip('.') print(f"清除 hdfs 目录:{hdfs_path}")
return s HdfsUtils.delete_file_in_folder(hdfs_path)
except:
return '' print(f"落表 {self.db_save},分区 {self.partitions_by},date_info {date_info}")
self.save_data_common(df_save=df_save, db_save=self.db_save,
def format_str(v): partitions_num=self.partitions_num, partitions_by=self.partitions_by)
"""字符串类型""" # 落表后回查该分区实际条数(读已写分区,不重算 df_save 整条链路)
if pd.isna(v): written_count = self.spark.sql(
return '' f"select count(1) as cnt from {self.db_save} "
return str(v) f"where site_name='{self.site_name}' and date_type='{self.date_type}' and date_info='{self.date_info}'"
).collect()[0]["cnt"]
# ========== 应用格式化 ========== print(f"save_data 成功,落表条数={written_count}")
for c in pure_int_cols:
if c in p_df.columns:
p_df[c] = p_df[c].apply(format_pure_int)
for c in int_with_decimal_cols:
if c in p_df.columns:
p_df[c] = p_df[c].apply(format_int_with_decimal)
for c in double_int_cols:
if c in p_df.columns:
p_df[c] = p_df[c].apply(format_double_int)
for c in double2_cols:
if c in p_df.columns:
p_df[c] = p_df[c].apply(format_double2)
for c in double4_cols:
if c in p_df.columns:
p_df[c] = p_df[c].apply(format_double4)
for c in str_cols:
if c in p_df.columns:
p_df[c] = p_df[c].apply(format_str)
# 处理其他未明确分类的列(兜底)
processed_cols = set(pure_int_cols + int_with_decimal_cols + double_int_cols +
double2_cols + double4_cols + str_cols)
for c in p_df.columns:
if c not in processed_cols:
p_df[c] = p_df[c].apply(format_str)
# 3. 最后保存
p_df.to_csv(f"/home/hejiangming/bs_top100_pandas_{self.site_name}_{self.date_info}.csv", index=False, encoding='utf-8',quoting = 1)
print("pandas方法保存csv成功 ")
# def save_spark_to_csv(self):
# (self.df_save.repartition(1) # 确保只生成一个 CSV 文件
# .write
# .mode("overwrite")
# .option("header", "true")
# .option("sep", ",") # 指定分隔符
# .csv(f"/home/hejiangming/bs_top100_spark_{self.site_name}_{self.date_info}.csv"))
# print("spark方法保存csv成功 ") # 没权限
def save_data(self):
print("save_data调用")
# self.df_save = self.df_save.toPandas()
# # quoting = 1 所有字段都加双引号 escapechar="\\" 转义字符 当字段里本身有 双引号 时,需要转义,否则 CSV 解析器会报错
# self.df_save.to_csv(f"/home/hejiangming/bs_top100_{self.site_name}_{self.date_info}.csv", index=False,encoding="utf-8-sig",quoting = 1,escapechar="\\")
self.save_pandas_to_csv()
print("success保存成功")
if __name__ == '__main__': if __name__ == '__main__':
......
"""
@Author : hejiangming
@Description : 品类调研 同比(YoY)/环比(MoM) 计算。
读 dwt_bs_top100 的"本期 / 上月 / 去年同月"三个月分区,按分类路径对齐,
算 14 指标 × 2 = 28 个变化率列(12 业务指标 + 2 利润率指标)。
设计文档:dwt_bs_top100_概览.md 第七节、dwt_bs_top100_字段精度核对.md 第一节。
@SourceTable : dwt_bs_top100 (本期 date_info、上月 M-1、去年同月 M-12 三个分区)
@SinkTable : dwt_bs_top100_change_rate
@CreateTime : 2026/07/14 10:00
@UpdateTime : 2026/07/14 10:00
"""
import os
import sys
sys.path.append(os.path.dirname(sys.path[0]))
from utils.hdfs_utils import HdfsUtils
from utils.common_util import CommonUtil
from utils.spark_util import SparkUtil
from pyspark.sql import functions as F
class DwtBsTop100ChangeRate(object):
"""
在整个品类调研链路里的位置:
dwt_flow_asin → dwt_bs_top100(各月落表,本表的上游)→ 本任务(跨月比)→ 导出 MySQL 给后端/前端。
同比环比本质 = "拿本表历史期数值比",所以必须先有历史月分区(M-1、M-12)才能算。
"""
# ============================================================
# 14 个参与同比环比的指标 + 计算类型(列名/顺序对齐 字段精度核对.md 29-40 行)
# 'up' = 越大越好,普通公式 (本期-对比期)/对比期,正数=变好
# 'rank' = 平均排名,唯一"越小越好",分子调换 (对比期-本期)/对比期,使正数=排名改善
# 'profit' = 利润率,用差值(本期率−对比期率, 涨跌几个百分点),任一期缺失→null,不走 -1000/1000 那套
# ============================================================
METRICS = [
('asin_amazon_orders_sum', 'up'), # 亚马逊月总销量
('asin_bsr_orders_sum', 'up'), # BSR月总销量
('asin_bsr_orders_mean', 'up'), # BSR月均销量
('asin_bsr_orders_sum_new', 'up'), # 新品BSR总销量
('asin_bsr_orders_sum_brand', 'up'), # 品牌BSR总销量
('asin_count', 'up'), # 总产品数
('asin_count_new', 'up'), # 新品数
('asin_count_brand', 'up'), # 品牌数
('asin_bs_cate_1_rank_mean', 'rank'), # 平均一级分类排名(越小越好)
('asin_price_mean', 'up'), # 平均价格
('asin_rating_mean', 'up'), # 平均评分
('asin_total_comments_mean', 'up'), # 平均评论数
('asin_ocean_freight_gross_margin_avg', 'profit'), # 海运平均毛利润率(需求3新增)
('asin_air_freight_gross_margin_avg', 'profit'), # 空运平均毛利润率(需求3新增)
]
# 跨月对齐主键:分类多归属(同 current_id 挂多个父级路径=多行),只用 current_id 会 1:N 串行。
# current_id + parent_id_join(根→直接父级的完整路径)= 每行唯一(已实测 2026-05 分区无重复)。
JOIN_KEYS = ['asin_bs_cate_current_id', 'asin_bs_cate_parent_id_join']
def __init__(self, site_name, date_type, date_info):
self.site_name = site_name
self.date_type = date_type
self.date_info = date_info
self.hive_tb = "dwt_bs_top100_change_rate"
self.source_tb = "dwt_bs_top100"
app_name = f"{self.hive_tb}:{site_name}:{date_type}:{date_info}"
self.spark = SparkUtil.get_spark_session(app_name)
self.partitions_num = CommonUtil.reset_partitions(site_name, 2)
# 环比 = 上月(M-1),同比 = 去年同月(M-12)。月粒度用 get_month_offset 直接推,确定性强。
self.last_month_info = CommonUtil.get_month_offset(self.date_info, -1)
self.last_year_info = CommonUtil.get_month_offset(self.date_info, -12)
print(f"本期={self.date_info},环比对比期(上月)={self.last_month_info},同比对比期(去年同月)={self.last_year_info}")
# 落表前清 HDFS 分区目录(幂等重跑,先删后写不残留旧文件)
hdfs_path = CommonUtil.build_hdfs_path(
self.hive_tb,
partition_dict={"site_name": site_name, "date_type": date_type, "date_info": date_info},
)
print(f"清除 hdfs 目录:{hdfs_path}")
HdfsUtils.delete_hdfs_file(hdfs_path)
self.df_base = self.spark.sql("select 1+1;") # 本期
self.df_last = self.spark.sql("select 1+1;") # 上月(环比对比期)
self.df_year = self.spark.sql("select 1+1;") # 去年同月(同比对比期)
self.df_save = self.spark.sql("select 1+1;")
def _read_one_period(self, date_info, prefix):
"""
读 dwt_bs_top100 一个月分区,只取 join 键 + 14 个指标列。
prefix 非空时把指标列改名加前缀(lm_/ym_),避免和本期同名列 join 后冲突。
"""
metric_cols = [m for m, _ in self.METRICS]
select_cols = self.JOIN_KEYS + metric_cols
sql = f"""
select {', '.join(select_cols)}
from {self.source_tb}
where site_name='{self.site_name}' and date_type='{self.date_type}' and date_info='{date_info}'
"""
df = self.spark.sql(sql)
if prefix:
for m in metric_cols:
df = df.withColumnRenamed(m, f"{prefix}{m}")
return df
def read_data(self):
# 本期额外带 level(落表用,便于后端/排查),对比期只要指标值
metric_cols = [m for m, _ in self.METRICS]
base_sql = f"""
select {', '.join(self.JOIN_KEYS + ['level'] + metric_cols)}
from {self.source_tb}
where site_name='{self.site_name}' and date_type='{self.date_type}' and date_info='{self.date_info}'
"""
self.df_base = self.spark.sql(base_sql).repartition(self.JOIN_KEYS[0]).cache()
print("df_base 行数", self.df_base.count())
self.df_last = self._read_one_period(self.last_month_info, "lm_").repartition(self.JOIN_KEYS[0]).cache()
self.df_year = self._read_one_period(self.last_year_info, "ym_").repartition(self.JOIN_KEYS[0]).cache()
def _rate_expr(self, cur_col, cmp_col, kind):
"""
单个指标的变化率表达式(cur=本次/本期分子,cmp=上次/对比期分母)。真值表见文件说明/概览.md:
空/空→null 空/0→0 空/>0→-1000(下降) 0/空、0/0→0 >0/空、>0/0→1000(上升) 其余正常算。
利润率(profit)单独:用差值(本期率−对比期率, 即涨跌几个百分点),任一期缺失→null(后端确认收 null)。
"""
cur = F.col(cur_col)
cmp = F.col(cmp_col)
# 利润率:改用【差值】= 本期率 − 对比期率(涨跌几个百分点),不做除法。
# 为什么不用相对变化率:对比期为负(亏损基期)时 (本-对)/对 的符号会翻转、误导;差值无此问题,也无需防除0。
# 任一期无数据(null) → null。base 无数据本就是 null。
if kind == 'profit':
# 原相对变化率(负基期符号会反,已弃用):
# return (F.when(cur.isNull() | cmp.isNull() | (cmp == 0), F.lit(None))
# .otherwise(F.round((cur - cmp) / cmp, 4)))
return (F.when(cur.isNull() | cmp.isNull(), F.lit(None))
.otherwise(F.round(cur - cmp, 4)))
# 【关键边界】业务 12 指标在 base 里"空"的编码 = null 或 -1 占位
# (均值列无数据存 -1,销量/计数列 0 兜底不会是 -1、也不会 null)。统一把 (null 或 -1) 当"空"。
# base 业务列不会 null(销量/计数 0 兜底、均值列无数据=-1),"空"只需判 -1;
# 只有对比期(cmp)会因 join 未命中变 null,故 cmp 才判 isNull。
# cur_empty = cur.isNull() | (cur == -1) # 原写法:cur.isNull() 对 base 永不触发,去掉避免误解
cur_empty = (cur == -1)
cmp_empty = cmp.isNull() | (cmp == -1)
cur_zero = (cur == 0)
cmp_zero = (cmp == 0)
# 正常计算:排名越小越好→分子调换让"正数=排名改善";其余越大越好 (本-上)/上
if kind == 'rank':
normal = F.round((cmp - cur) / cmp, 4)
else:
normal = F.round((cur - cmp) / cmp, 4)
# 严格按真值表分支(顺序不能乱:先定本次空/0/>0,再看对比期)
return (
F.when(cur_empty & cmp_empty, F.lit(None)) # 空/空 → 空
.when(cur_empty & cmp_zero, F.lit(0)) # 空/0 → 0
.when(cur_empty, F.lit(-1000)) # 空/>0 → -1000(本期没了=下降)
.when(cur_zero & (cmp_empty | cmp_zero), F.lit(0)) # 0/空、0/0 → 0(无方向)
.when(cur_zero, normal) # 0/>0 → 正常(=-1.0)
.when(cmp_empty | cmp_zero, F.lit(1000)) # >0/空、>0/0 → 1000(对比期无=上升)
.otherwise(normal) # >0/>0 → 正常
)
def handle_data(self):
# 本期为 base 左连对比期:本期在榜的分类才产出行;对比期缺失的行 cmp 列为 null,走上面 den_invalid 分支
df = self.df_base.join(self.df_last, on=self.JOIN_KEYS, how='left') \
.join(self.df_year, on=self.JOIN_KEYS, how='left')
# 逐指标生成 <指标>_lm_rate(环比) 和 <指标>_ym_rate(同比)
for m, kind in self.METRICS:
df = df.withColumn(f"{m}_lm_rate", self._rate_expr(m, f"lm_{m}", kind))
df = df.withColumn(f"{m}_ym_rate", self._rate_expr(m, f"ym_{m}", kind))
self.df_save = df
self.df_base.unpersist()
self.df_last.unpersist()
self.df_year.unpersist()
def save_data(self):
# 落表 = join 键 + level + 28 个 rate 列(每指标 ym 同比、lm 环比)+ 3 分区列
rate_cols = []
for m, _ in self.METRICS:
rate_cols.append(f"{m}_ym_rate")
rate_cols.append(f"{m}_lm_rate")
out_cols = self.JOIN_KEYS + ['level'] + rate_cols
self.df_save = self.df_save.select(*out_cols) \
.withColumn("site_name", F.lit(self.site_name)) \
.withColumn("date_type", F.lit(self.date_type)) \
.withColumn("date_info", F.lit(self.date_info))
# 不用 auto_transfer_type:它见 Hive double 列会强制 cast 成 decimal(10,3),把 4 位砍成 3 位,
# 且极端大变化率(对比期是极小正数时)会超 decimal(10,3) 上限、cast 成 null(哪怕列是 double 也救不了,卡在中间那道 decimal)。
# 表本身是 double,直接把 28 个 rate 列统一 cast 成 double 落表 —— 保留 round(4) 的 4 位精度、double 值域大不溢出。
for c in rate_cols:
self.df_save = self.df_save.withColumn(c, F.col(c).cast("double"))
self.df_save = self.df_save.repartition(self.partitions_num)
partition_by = ["site_name", "date_type", "date_info"]
print(f"当前存储的表名为:{self.hive_tb},分区为 {partition_by}")
self.df_save.write.saveAsTable(name=self.hive_tb, format='hive', mode='append', partitionBy=partition_by)
print("success")
def run(self):
self.read_data()
self.handle_data()
self.save_data()
if __name__ == '__main__':
site_name = CommonUtil.get_sys_arg(1, None) # 站点 us
date_type = CommonUtil.get_sys_arg(2, None) # month
date_info = CommonUtil.get_sys_arg(3, None) # 年-月,如 2026-05
obj = DwtBsTop100ChangeRate(site_name, date_type, date_info)
obj.run()
"""
@Author : hejiangming
@Description : 品类调研主表导出 PG 集群:Hive dwt_bs_top100 → PG {site}_category_top_analysis。
PG 侧是"单基表 + 按月 RANGE(date_info) 分区"(趋势图跨月/跨年查一张表, 详见 us_category_top_analysis.sql)。
每月导数:建当月空分区 + copy 表 → sqoop 灌 copy 表 → 分区交换换入(exchange_pg_part_tb), 幂等重导、无空窗期。
@SourceTable : dwt_bs_top100 (Hive, 分区 site_name/date_type/date_info)
@SinkTable : {site}_category_top_analysis (PG 集群, 按 date_info 月分区)
@CreateTime : 2026/07/15
@UpdateTime : 2026/07/15
"""
import os
import sys
sys.path.append(os.path.dirname(sys.path[0]))
from utils.ssh_util import SSHUtil
from utils.common_util import CommonUtil, DateTypes
from utils.db_util import DBUtil
from utils.hdfs_utils import HdfsUtils
if __name__ == '__main__':
site_name = CommonUtil.get_sys_arg(1, None)
date_type = CommonUtil.get_sys_arg(2, None)
date_info = CommonUtil.get_sys_arg(3, None)
# 最后一个参数判断是否导测试库
test_flag = CommonUtil.get_sys_arg(len(sys.argv) - 1, None)
print(f"执行参数为{sys.argv}")
# 品类调研只有月粒度
assert date_type == DateTypes.month.name, f"仅支持 month, 传入 date_type={date_type}"
if test_flag == 'test':
db_type = 'postgresql_test'
print("导出到PG测试库中")
else:
CommonUtil.judge_is_work_hours(site_name=site_name, date_type=date_type, date_info=date_info,
principal='hejiangming', priority=2, export_tools_type=1,
belonging_to_process=f'品类调研流程_{date_type}')
db_type = 'postgresql_cluster'
print("导出到PG集群库中")
# 获取数据库连接
engine = DBUtil.get_db_engine(db_type, site_name)
# 导出前校验 Hive 分区有数据:空分区会导致灌进空 copy 表、分区交换后把 PG 当月换成空
hive_path = CommonUtil.build_hdfs_path(
"dwt_bs_top100",
partition_dict={"site_name": site_name, "date_type": date_type, "date_info": date_info},
)
if not HdfsUtils.read_list(hive_path):
print(f"[ERROR] Hive 分区无数据文件:{hive_path},跳过导出,请先检查 DWT 计算任务!")
engine.dispose()
sys.exit(1)
print(f"Hive 分区有数据:{hive_path},继续导出")
# 用"copy 表 + 分区交换"方式导数(比先 truncate 再灌更稳:先把数据灌进 copy 表,成功后再原子换入分区,
# sqoop 中途失败也不会清空线上当月数据、无空窗期;机制同 dwt_aba_st_analytics.py 的 month 分支)。
base_tb = f"{site_name}_category_top_analysis"
date_partition = str(date_info).replace("-", "_")
export_table_original = f"{base_tb}_{date_partition}" # 当月分区子表, 如 us_category_top_analysis_2026_05
export_tb_copy = f"{export_table_original}_copy" # 灌数用的临时普通表
next_val = CommonUtil.get_next_val(date_type, date_info) # 下月, 作分区右边界(不含)
# 1) 建当月分区(交换时需 detach 它, 首次跑先建空的);2) 建 copy 表(形如分区子表, 含索引)并清空
sql = f"""
create table if not exists {export_table_original} partition of {base_tb} for values from ('{date_info}') to ('{next_val}');
create table if not exists {export_tb_copy} (like {export_table_original} including all);
truncate table {export_tb_copy};
"""
DBUtil.engine_exec_sql(engine, sql)
# 导出列 = 86 业务列 + 3 分区列(列名与 Hive dwt_bs_top100 一致, sqoop 按名映射 PG 同名列)
export_cols = [
# 维度 / 键
"asin_bs_cate_current_id", "asin_bs_cate_parent_id", "asin_bs_cate_parent_id_join",
"level", "cur_category_name",
# 销量类
"asin_amazon_orders_sum", "asin_bsr_orders_sum", "asin_bsr_orders_sum_new", "asin_bsr_orders_sum_brand",
"asin_bsr_orders_mean", "asin_bsr_orders_mean_new", "asin_bsr_orders_mean_brand", "asin_bsr_orders_sum_parent",
# 销售额类
"asin_bsr_orders_sales_sum", "asin_bsr_orders_sales_sum_new", "asin_bsr_orders_sales_sum_brand",
"asin_bsr_orders_sales_mean", "asin_bsr_orders_sales_mean_new", "asin_bsr_orders_sales_mean_brand",
# 数量类
"asin_count", "asin_count_new", "asin_count_brand",
"asin_count_rate_new", "asin_count_rate_brand",
"asin_bsr_orders_sum_rate_new", "asin_bsr_orders_sum_rate_brand", "asin_bsr_orders_sum_rate",
# 价格/评分/排名/评论 均值
"asin_price_mean", "asin_rating_mean", "asin_bs_cate_1_rank_mean", "asin_total_comments_mean",
# 上架时间分布
"launch_time_type0_rate", "launch_time_type1_rate", "launch_time_type2_rate", "launch_time_type3_rate",
"launch_time_type4_rate", "launch_time_type5_rate", "launch_time_type6_rate", "launch_time_type7_rate",
# 卖家类型分布
"asin_buy_box_seller_type0_rate", "asin_buy_box_seller_type1_rate", "asin_buy_box_seller_type2_rate",
"asin_buy_box_seller_type3_rate", "asin_buy_box_seller_type4_rate",
# 卖家所属地占比
"seller_country_us_rate", "seller_country_cn_rate", "seller_country_other_rate",
# 月销阈值占比
"asin_orders_ge_50_rate", "asin_orders_ge_100_rate", "asin_orders_ge_200_rate",
# 利润率均值
"asin_ocean_freight_gross_margin_avg", "asin_air_freight_gross_margin_avg",
# 区间利润率(海运12)
"asin_ocean_freight_gross_margin_price_0_10", "asin_ocean_freight_gross_margin_price_10_20",
"asin_ocean_freight_gross_margin_price_20_30", "asin_ocean_freight_gross_margin_price_30_50",
"asin_ocean_freight_gross_margin_price_50_80", "asin_ocean_freight_gross_margin_price_80",
"asin_ocean_freight_gross_margin_orders_0_50", "asin_ocean_freight_gross_margin_orders_50_100",
"asin_ocean_freight_gross_margin_orders_100_200", "asin_ocean_freight_gross_margin_orders_200_500",
"asin_ocean_freight_gross_margin_orders_500_1000", "asin_ocean_freight_gross_margin_orders_1000",
# 区间利润率(空运12)
"asin_air_freight_gross_margin_price_0_10", "asin_air_freight_gross_margin_price_10_20",
"asin_air_freight_gross_margin_price_20_30", "asin_air_freight_gross_margin_price_30_50",
"asin_air_freight_gross_margin_price_50_80", "asin_air_freight_gross_margin_price_80",
"asin_air_freight_gross_margin_orders_0_50", "asin_air_freight_gross_margin_orders_50_100",
"asin_air_freight_gross_margin_orders_100_200", "asin_air_freight_gross_margin_orders_200_500",
"asin_air_freight_gross_margin_orders_500_1000", "asin_air_freight_gross_margin_orders_1000",
# 打包段月销占比
"asin_pkg_1_orders_rate", "asin_pkg_2_4_orders_rate", "asin_pkg_5_10_orders_rate",
"asin_pkg_10_20_orders_rate", "asin_pkg_20_50_orders_rate", "asin_pkg_50_100_orders_rate",
"asin_pkg_100_orders_rate",
# 集中度
"asin_concentration", "brand_concentration", "seller_concentration",
# 分区列(也落成普通列, 取分区值)
"site_name", "date_type", "date_info",
]
sh = CommonUtil.build_export_sh(
site_name=site_name,
db_type=db_type,
hive_tb="dwt_bs_top100",
export_tb=export_tb_copy,
col=export_cols,
partition_dict={"site_name": site_name, "date_type": date_type, "date_info": date_info},
)
client = SSHUtil.get_ssh_client()
SSHUtil.exec_command_async(client, sh, ignore_err=False)
client.close()
# 分区交换:把灌好的 copy 表原子换入当月分区(旧分区在函数内转到 copy 名下),避免 truncate 的空窗期。
# 分区表主键/索引由父表继承, cp_index_flag=False 不额外建。
DBUtil.exchange_pg_part_tb(
engine,
source_tb_name=export_tb_copy,
part_master_tb=base_tb,
part_target_tb=export_table_original,
cp_index_flag=False,
part_val={"from": [date_info], "to": [next_val]},
)
# 交换后 copy 表里是旧数据, 删掉省空间(同 aba)
DBUtil.engine_exec_sql(engine, f"drop table if exists {export_tb_copy};")
# 关闭链接
engine.dispose()
# 导出完成状态更新:主表 + 同比环比表两个导出脚本共用同一 belonging_to_process(品类调研流程_{date_type})。
# modify_export_workflow_status 会把本脚本置完成(status=3)、统计同组未完成数——只有最后一个跑完的
# (未完成数=0)才执行下面的 update_workflow_sql,即"两个都导完,页面才展示"(机制同 dwt_aba_last_change_rate.py)。
# 两脚本传同一条 update SQL,谁最后完成都能正确触发。本地/test 运行时该函数自动跳过。
# 本功能不像 aba 会在流程启动时预写一条 workflow_everyday,故用 replace INTO 直接插入/覆盖该行(不存在则建、存在则替换)。
# 只有主表+同比环比两个导出都完成时,modify_export_workflow_status 才会执行这条(即页面此时才出现"完成"行)。
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}');
"""
CommonUtil.modify_export_workflow_status(update_workflow_sql, site_name, date_type, date_info)
print("success")
"""
@Author : hejiangming
@Description : 品类调研同比/环比表导出 PG 集群:Hive dwt_bs_top100_change_rate → PG {site}_category_top_analysis_change_rate。
PG 侧同主表:单基表 + 按月 RANGE(date_info) 分区(详见 us_category_top_analysis.sql)。
每月导数:建当月空分区 + copy 表 → sqoop 灌 copy 表 → 分区交换换入(exchange_pg_part_tb), 幂等重导、无空窗期。
@SourceTable : dwt_bs_top100_change_rate (Hive, 分区 site_name/date_type/date_info)
@SinkTable : {site}_category_top_analysis_change_rate (PG 集群, 按 date_info 月分区)
@CreateTime : 2026/07/15
@UpdateTime : 2026/07/15
"""
import os
import sys
sys.path.append(os.path.dirname(sys.path[0]))
from utils.ssh_util import SSHUtil
from utils.common_util import CommonUtil, DateTypes
from utils.db_util import DBUtil
from utils.hdfs_utils import HdfsUtils
if __name__ == '__main__':
site_name = CommonUtil.get_sys_arg(1, None)
date_type = CommonUtil.get_sys_arg(2, None)
date_info = CommonUtil.get_sys_arg(3, None)
# 最后一个参数判断是否导测试库
test_flag = CommonUtil.get_sys_arg(len(sys.argv) - 1, None)
print(f"执行参数为{sys.argv}")
# 品类调研只有月粒度
assert date_type == DateTypes.month.name, f"仅支持 month, 传入 date_type={date_type}"
if test_flag == 'test':
db_type = 'postgresql_test'
print("导出到PG测试库中")
else:
# 非上班时段导出约束(本地/test 会自动跳过)
CommonUtil.judge_is_work_hours(site_name=site_name, date_type=date_type, date_info=date_info,
principal='hejiangming', priority=2, export_tools_type=1,
belonging_to_process=f'品类调研流程_{date_type}')
db_type = 'postgresql_cluster'
print("导出到PG集群库中")
# 获取数据库连接
engine = DBUtil.get_db_engine(db_type, site_name)
# 导出前校验 Hive 分区有数据:空分区会导致灌进空 copy 表、分区交换后把 PG 当月换成空
hive_path = CommonUtil.build_hdfs_path(
"dwt_bs_top100_change_rate",
partition_dict={"site_name": site_name, "date_type": date_type, "date_info": date_info},
)
if not HdfsUtils.read_list(hive_path):
print(f"[ERROR] Hive 分区无数据文件:{hive_path},跳过导出,请先检查 DWT 计算任务!")
engine.dispose()
sys.exit(1)
print(f"Hive 分区有数据:{hive_path},继续导出")
# 用"copy 表 + 分区交换"方式导数(比先 truncate 再灌更稳:先把数据灌进 copy 表,成功后再原子换入分区,
# sqoop 中途失败也不会清空线上当月数据、无空窗期;机制同 dwt_aba_st_analytics.py 的 month 分支)。
base_tb = f"{site_name}_category_top_analysis_change_rate"
date_partition = str(date_info).replace("-", "_")
export_table_original = f"{base_tb}_{date_partition}" # 当月分区子表, 如 us_category_top_analysis_change_rate_2026_05
export_tb_copy = f"{export_table_original}_copy" # 灌数用的临时普通表
next_val = CommonUtil.get_next_val(date_type, date_info) # 下月, 作分区右边界(不含)
# 1) 建当月分区(交换时需 detach 它, 首次跑先建空的);2) 建 copy 表(形如分区子表, 含索引)并清空
sql = f"""
create table if not exists {export_table_original} partition of {base_tb} for values from ('{date_info}') to ('{next_val}');
create table if not exists {export_tb_copy} (like {export_table_original} including all);
truncate table {export_tb_copy};
"""
DBUtil.engine_exec_sql(engine, sql)
# 导出列 = join键2 + level + 28 变化率列 + 3 分区列(列名与 Hive dwt_bs_top100_change_rate 一致)
export_cols = [
"asin_bs_cate_current_id", "asin_bs_cate_parent_id_join", "level",
"asin_amazon_orders_sum_ym_rate", "asin_amazon_orders_sum_lm_rate",
"asin_bsr_orders_sum_ym_rate", "asin_bsr_orders_sum_lm_rate",
"asin_bsr_orders_mean_ym_rate", "asin_bsr_orders_mean_lm_rate",
"asin_bsr_orders_sum_new_ym_rate", "asin_bsr_orders_sum_new_lm_rate",
"asin_bsr_orders_sum_brand_ym_rate", "asin_bsr_orders_sum_brand_lm_rate",
"asin_count_ym_rate", "asin_count_lm_rate",
"asin_count_new_ym_rate", "asin_count_new_lm_rate",
"asin_count_brand_ym_rate", "asin_count_brand_lm_rate",
"asin_bs_cate_1_rank_mean_ym_rate", "asin_bs_cate_1_rank_mean_lm_rate",
"asin_price_mean_ym_rate", "asin_price_mean_lm_rate",
"asin_rating_mean_ym_rate", "asin_rating_mean_lm_rate",
"asin_total_comments_mean_ym_rate", "asin_total_comments_mean_lm_rate",
"asin_ocean_freight_gross_margin_avg_ym_rate", "asin_ocean_freight_gross_margin_avg_lm_rate",
"asin_air_freight_gross_margin_avg_ym_rate", "asin_air_freight_gross_margin_avg_lm_rate",
# 分区列(也落成普通列, 取分区值)
"site_name", "date_type", "date_info",
]
sh = CommonUtil.build_export_sh(
site_name=site_name,
db_type=db_type,
hive_tb="dwt_bs_top100_change_rate",
export_tb=export_tb_copy,
col=export_cols,
partition_dict={"site_name": site_name, "date_type": date_type, "date_info": date_info},
)
client = SSHUtil.get_ssh_client()
SSHUtil.exec_command_async(client, sh, ignore_err=False)
client.close()
# 分区交换:把灌好的 copy 表原子换入当月分区(旧分区在函数内转到 copy 名下),避免 truncate 的空窗期。
# 分区表主键/索引由父表继承, cp_index_flag=False 不额外建。
DBUtil.exchange_pg_part_tb(
engine,
source_tb_name=export_tb_copy,
part_master_tb=base_tb,
part_target_tb=export_table_original,
cp_index_flag=False,
part_val={"from": [date_info], "to": [next_val]},
)
# 交换后 copy 表里是旧数据, 删掉省空间(同 aba)
DBUtil.engine_exec_sql(engine, f"drop table if exists {export_tb_copy};")
# 关闭链接
engine.dispose()
# 导出完成状态更新:主表 + 同比环比表两个导出脚本共用同一 belonging_to_process(品类调研流程_{date_type})。
# modify_export_workflow_status 会把本脚本置完成(status=3)、统计同组未完成数——只有最后一个跑完的
# (未完成数=0)才执行下面的 update_workflow_sql,即"两个都导完,页面才展示"(机制同 dwt_aba_last_change_rate.py)。
# 两脚本传同一条 update SQL,谁最后完成都能正确触发。本地/test 运行时该函数自动跳过。
# 本功能不像 aba 会在流程启动时预写一条 workflow_everyday,故用 replace INTO 直接插入/覆盖该行(不存在则建、存在则替换)。
# 只有主表+同比环比两个导出都完成时,modify_export_workflow_status 才会执行这条(即页面此时才出现"完成"行)。
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}');
"""
CommonUtil.modify_export_workflow_status(update_workflow_sql, site_name, date_type, date_info)
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