Commit 0e6388f9 by chenyuanjie

流量选品月-修改缓存类型

parent bf1e0aa8
...@@ -415,17 +415,19 @@ class DwtFlowAsin(Templates): ...@@ -415,17 +415,19 @@ class DwtFlowAsin(Templates):
).withColumn("asin_ao_val_type", F.expr("""CASE WHEN asin_ao_val BETWEEN 0 AND 0.1 THEN 1 ).withColumn("asin_ao_val_type", F.expr("""CASE WHEN asin_ao_val BETWEEN 0 AND 0.1 THEN 1
WHEN asin_ao_val BETWEEN 0.1 AND 0.2 THEN 2 WHEN asin_ao_val BETWEEN 0.2 AND 0.4 THEN 3 WHEN asin_ao_val BETWEEN 0.1 AND 0.2 THEN 2 WHEN asin_ao_val BETWEEN 0.2 AND 0.4 THEN 3
WHEN asin_ao_val BETWEEN 0.4 AND 0.8 THEN 4 WHEN asin_ao_val BETWEEN 0.8 AND 1.2 THEN 5 WHEN asin_ao_val BETWEEN 0.4 AND 0.8 THEN 4 WHEN asin_ao_val BETWEEN 0.8 AND 1.2 THEN 5
WHEN asin_ao_val BETWEEN 1.2 AND 2 THEN 6 WHEN asin_ao_val >= 2 THEN 7 ELSE 0 END""")).drop("asin_amazon_orders").cache() WHEN asin_ao_val BETWEEN 1.2 AND 2 THEN 6 WHEN asin_ao_val >= 2 THEN 7 ELSE 0 END""")).drop("asin_amazon_orders").persist(StorageLevel.DISK_ONLY)
# 关联流量asin由于来自广告位,无法计算ao值,填充母体ao值和母体zr占比 # 关联流量asin由于来自广告位,无法计算ao值,填充母体ao值和母体zr占比
df_parent_asin_ao_val = self.df_asin_detail.select('parent_asin', 'matrix_ao_val', 'matrix_flow_proportion').groupby(['parent_asin']).agg( df_ao_stage = self.df_asin_detail
df_parent_asin_ao_val = df_ao_stage.select('parent_asin', 'matrix_ao_val', 'matrix_flow_proportion').groupby(['parent_asin']).agg(
F.first('matrix_ao_val').alias('parent_matrix_ao_val'), F.first('matrix_flow_proportion').alias('parent_matrix_flow_proportion') F.first('matrix_ao_val').alias('parent_matrix_ao_val'), F.first('matrix_flow_proportion').alias('parent_matrix_flow_proportion')
) )
self.df_asin_detail = self.df_asin_detail.join(df_parent_asin_ao_val, on=['parent_asin'], how='left').withColumn( self.df_asin_detail = df_ao_stage.join(df_parent_asin_ao_val, on=['parent_asin'], how='left').withColumn(
"matrix_ao_val", F.coalesce(F.col("matrix_ao_val"), F.col("parent_matrix_ao_val")) "matrix_ao_val", F.coalesce(F.col("matrix_ao_val"), F.col("parent_matrix_ao_val"))
).withColumn( ).withColumn(
"matrix_flow_proportion", F.coalesce(F.col("matrix_flow_proportion"), F.col("parent_matrix_flow_proportion")) "matrix_flow_proportion", F.coalesce(F.col("matrix_flow_proportion"), F.col("parent_matrix_flow_proportion"))
).drop("parent_matrix_ao_val", "parent_matrix_flow_proportion") ).drop("parent_matrix_ao_val", "parent_matrix_flow_proportion")
self.df_asin_measure.unpersist() self.df_asin_measure.unpersist()
df_ao_stage.unpersist()
def handle_parent_asin_variation(self): def handle_parent_asin_variation(self):
"""处理父ASIN变体聚合数据,结果存入 self.df_parent_asin_variat_agg""" """处理父ASIN变体聚合数据,结果存入 self.df_parent_asin_variat_agg"""
...@@ -922,6 +924,12 @@ class DwtFlowAsin(Templates): ...@@ -922,6 +924,12 @@ class DwtFlowAsin(Templates):
self.handle_asin_detail_all_type() self.handle_asin_detail_all_type()
self.handle_asin_category_info() self.handle_asin_category_info()
self.handle_asin_measure() self.handle_asin_measure()
# 打断 lazy 链:measure 阶段含自连接,落盘后释放,避免后续 action 触发时计划树过长导致 OOM
df_measure_done = self.df_asin_detail
self.df_asin_detail = self.df_asin_detail.repartition(60).persist(StorageLevel.DISK_ONLY)
cnt = self.df_asin_detail.count()
print(f"measure 阶段处理完成,数据量: {cnt}")
df_measure_done.unpersist()
self.handle_parent_asin_variation() self.handle_parent_asin_variation()
self.handle_seller_country() self.handle_seller_country()
self.handle_asin_lqs_rating() self.handle_asin_lqs_rating()
......
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