Commit e331ff5f by chenyuanjie

流量选品30天-增加竞卖指标计算

parent 81c19a7f
......@@ -395,7 +395,8 @@ class KafkaFlowAsinDetail(Templates):
# 8. 处理变体相关(ao及母体相关,自然占比及母体自然占比,各类型数量,月销信息等)
def handle_asin_measure(self, df):
df = CommonUtil.get_asin_variant_attribute(df_asin_detail=df, df_asin_measure=self.df_asin_measure,
partition_num=self.repartition_num, use_type=1)
partition_num=self.repartition_num, use_type=1,
df_asin_auction=self.df_asin_auction)
# 是否数量变体类型和ao的类型
df = df.withColumn("quantity_variation_type", F.expr("""
CASE WHEN size is not null and size != '' and lower(size) like '%quantity%' THEN 1 ELSE 0 END""")).withColumn(
......@@ -410,7 +411,7 @@ class KafkaFlowAsinDetail(Templates):
.drop("asin_st_counts", "asin_adv_counts")
# 获取parent_asin下最新ASIN信息,导出到 doris 父ASIN最新详情表(仅 latest+normal 模式)
if self.test_flag == 'normal':
df_parent_asin_info = df.filter("parent_asin is not null").select("parent_asin", "asin_vartion_list", "asinUpdateTime", "matrix_flow_proportion", "matrix_ao_val")
df_parent_asin_info = df.filter("parent_asin is not null").select("parent_asin", "asin_vartion_list", "asinUpdateTime", "matrix_flow_proportion", "matrix_ao_val", "auctions_num_all", "skus_num_creat_all")
parent_asin_window = Window.partitionBy(['parent_asin']).orderBy(F.desc_nulls_last("asinUpdateTime"))
df_parent_asin_info = df_parent_asin_info.withColumn("u_rank", F.row_number().over(window=parent_asin_window))
df_parent_asin_info = df_parent_asin_info.repartition(self.repartition_num)
......@@ -425,8 +426,8 @@ class KafkaFlowAsinDetail(Templates):
F.concat_ws(',', F.collect_list("asin")).alias("variation_info"),
F.to_json(F.collect_list(F.struct(F.col("color"), F.col("size"), F.col("style")))).alias("attr_info")
)
# 关联母体 AO 值和自然流量占比(取 parent_asin 维度最新记录的值)
df_parent_matrix = df_parent_asin_info.select("parent_asin", "matrix_flow_proportion", "matrix_ao_val")
# 关联母体 AO 值、自然流量占比、竞卖数据(取 parent_asin 维度最新记录的值)
df_parent_matrix = df_parent_asin_info.select("parent_asin", "matrix_flow_proportion", "matrix_ao_val", "auctions_num_all", "skus_num_creat_all")
df_asin_variat_agg = df_asin_variat_agg.join(df_parent_matrix, on="parent_asin", how="left")
print("导出父ASIN最新变体信息到doris:")
df_doris = df_asin_variat_agg.select(
......@@ -437,8 +438,9 @@ class KafkaFlowAsinDetail(Templates):
"variation_info", "attr_info",
F.current_timestamp().alias("updated_at"),
F.round(F.col("matrix_flow_proportion"), 4).alias("matrix_flow_proportion"),
F.round(F.col("matrix_ao_val"), 4).alias("matrix_ao_val"))
table_columns = "parent_asin, date_info, asin_crawl_date, variation_info, attr_info, updated_at, matrix_flow_proportion, matrix_ao_val"
F.round(F.col("matrix_ao_val"), 4).alias("matrix_ao_val"),
"auctions_num_all", "skus_num_creat_all")
table_columns = "parent_asin, date_info, asin_crawl_date, variation_info, attr_info, updated_at, matrix_flow_proportion, matrix_ao_val, auctions_num_all, skus_num_creat_all"
DorisHelper.spark_export_with_columns(df_save=df_doris, db_name=self.doris_db_selection, table_name=self.parent_asin_latest_detail_table, table_columns=table_columns)
return df
......@@ -853,6 +855,9 @@ class KafkaFlowAsinDetail(Templates):
"sales_rise", "sales_mom", "sales_yoy",
"variation_rise", "variation_mom", "variation_yoy",
"bought_month_mom", "bought_month_yoy", # 月销无绝对值,仅环比/同比
# ── 竞卖/SKU创建 ───────────────────────────────────────────────
"auctions_num", "auctions_num_all", # get_asin_variant_attribute,ASIN竞卖数/父asin累计竞卖数
"skus_num_creat", "skus_num_creat_all", # get_asin_variant_attribute,ASIN创建SKU数/父asin累计创建SKU数
)
df_save = df_save.na.fill(
{"zr_counts": 0, "sp_counts": 0, "sb_counts": 0, "vi_counts": 0, "bs_counts": 0, "ac_counts": 0,
......@@ -1021,6 +1026,20 @@ class KafkaFlowAsinDetail(Templates):
username=mysql_con['username'], query=sql
).persist(StorageLevel.MEMORY_ONLY))
self.df_ai_hide_category.show(10, truncate=False)
print("12. 读取product_audit_asin_sku,得到asin竞卖数、sku创建数")
sql = """
select asin,
count(asin) as auctions_num,
count((case when sku != '' then sku end)) as skus_num_creat
from product_audit_asin_sku
WHERE asin REGEXP '^[0-9A-Z]{10}$'
group by asin
"""
self.df_asin_auction = SparkUtil.read_jdbc_query(
session=self.spark, url=mysql_con['url'], pwd=mysql_con['pwd'],
username=mysql_con['username'], query=sql
).repartition(self.repartition_num).persist(StorageLevel.DISK_ONLY)
self.df_asin_auction.show(10, truncate=False)
# 字段处理逻辑综合
def handle_all_field(self, df):
......@@ -1138,6 +1157,7 @@ class KafkaFlowAsinDetail(Templates):
"best_sellers_herf",
"asin_type",
"is_amazon_new",
"auctions_num", "auctions_num_all", "skus_num_creat", "skus_num_creat_all",
)
table_columns = """asin, asin_ao_val, asin_title, asin_title_len, asin_category_desc, asin_volume,
asin_weight, asin_launch_time, asin_brand_name, one_star, two_star, three_star, four_star, five_star, low_star,
......@@ -1145,7 +1165,8 @@ class KafkaFlowAsinDetail(Templates):
asin_crawl_date, asin_price, asin_rating, asin_total_comments, matrix_ao_val, zr_flow_proportion, matrix_flow_proportion,
date_info, img_url, category_current_id, category_first_rank, category_current_rank, bsr_orders, bsr_orders_sale,
page_inventory, asin_bought_month, seller_json, buy_box_seller_type, asin_describe, asin_fbm_price, asin_describe_len,
asin_weight_str, best_sellers_rank, best_sellers_herf, asin_type, is_amazon_new"""
asin_weight_str, best_sellers_rank, best_sellers_herf, asin_type, is_amazon_new,
auctions_num, auctions_num_all, skus_num_creat, skus_num_creat_all"""
DorisHelper.spark_export_with_columns(
df_save=df_asin_latest_detail, db_name=self.doris_db_selection,
table_name=self.asin_latest_detail_table, table_columns=table_columns
......@@ -1180,7 +1201,8 @@ class KafkaFlowAsinDetail(Templates):
site_name, asin_type, 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, multi_color_flag, multi_color_str, amazon_label, is_amazon_new"""
fbm_price, describe_len, multi_color_flag, multi_color_str, amazon_label, is_amazon_new,
auctions_num, auctions_num_all, skus_num_creat, skus_num_creat_all"""
print(f"写入Doris {self.doris_30day_table}")
DorisHelper.spark_export_with_columns(
df_save=df, db_name=self.doris_db_dwt,
......
......@@ -1776,7 +1776,7 @@ outputformat 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat'
return None
@classmethod
def get_asin_variant_attribute(cls, df_asin_detail: DataFrame, df_asin_measure: DataFrame, partition_num: int=80, use_type: int=0):
def get_asin_variant_attribute(cls, df_asin_detail: DataFrame, df_asin_measure: DataFrame, partition_num: int=80, use_type: int=0, df_asin_auction: DataFrame=None):
"""
Param df_asin_detail: asin详情DataFrame(
字段要求:
......@@ -1792,6 +1792,8 @@ outputformat 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat'
);
Param partition_num: 运行并行度(根据脚本运行资源设置)
Param use_type: 使用类型(0:默认,插件; 1:流量选品)
Param df_asin_auction: asin竞卖数据DataFrame(可选,字段要求:asin, auctions_num, skus_num_creat);
传入时才计算 auctions_num/auctions_num_all/skus_num_creat/skus_num_creat_all,不传则不处理这4个字段
return :
1. dwd_asin_measure必须携带的:
asin、asin_zr_counts, asin_adv_counts, asin_st_counts, asin_amazon_orders, asin_zr_flow_proportion, asin_ao_val
......@@ -1799,6 +1801,7 @@ outputformat 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat'
matrix_ao_val, matrix_flow_proportion, asin_amazon_orders, variant_info(变体asin列表)
3. 流量选品特定的: color, size, style
4. dwd_asin_measure自行携带的字段
5. 传入df_asin_auction时才有:auctions_num, auctions_num_all, skus_num_creat, skus_num_creat_all
"""
# 1.关联获取ao、各类型数量、流量占比信息、月销信息等
df_asin_detail = df_asin_detail.repartition(partition_num)
......@@ -1847,7 +1850,31 @@ outputformat 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat'
).withColumn(
"matrix_flow_proportion", F.coalesce(F.col("matrix_flow_proportion"), F.col("asin_zr_flow_proportion"))
)
# 4.解析变体属性信息(颜色、 尺寸、 风格等)
# 4.竞卖数据(auctions_num/skus_num_creat)及母体竞卖数据(_all):仅当传入df_asin_auction时计算
if df_asin_auction is not None:
df_asin_auction = df_asin_auction.repartition(partition_num)
df_asin_detail = df_asin_detail.join(df_asin_auction, on=['asin'], how='left')
df_asin_detail = df_asin_detail.withColumn(
"auctions_num", F.col("auctions_num").cast("int")
).withColumn(
"skus_num_creat", F.col("skus_num_creat").cast("int")
).na.fill({"auctions_num": 0, "skus_num_creat": 0})
df_variant_asin_auction = df_asin_auction.select(F.col("asin").alias("variant_asin"), "auctions_num", "skus_num_creat")
df_explode_variant_auction_detail = df_explode_variant_attribute.join(
df_variant_asin_auction, on=["variant_asin"], how="inner"
)
df_explode_variant_auction_agg = df_explode_variant_auction_detail.groupby(['asin']).agg(
F.sum("auctions_num").cast("int").alias("sum_auctions_num"),
F.sum("skus_num_creat").cast("int").alias("sum_skus_num_creat")
)
df_asin_detail = df_asin_detail.join(
df_explode_variant_auction_agg, on=['asin'], how='left'
).withColumn(
"auctions_num_all", F.coalesce(F.col("sum_auctions_num"), F.col("auctions_num"))
).withColumn(
"skus_num_creat_all", F.coalesce(F.col("sum_skus_num_creat"), F.col("skus_num_creat"))
).drop("sum_auctions_num", "sum_skus_num_creat")
# 5.解析变体属性信息(颜色、 尺寸、 风格等)
if use_type == 1:
df_asin_attribute = df_explode_variant_attribute.filter(F.col("asin") == F.col("variant_asin")).drop("variant_asin")
df_asin_detail = df_asin_detail.join(
......
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