Commit 662adc43 by chenyuanjie

流量选品30天-筛选asin信息库增加过滤

parent 9558457a
...@@ -116,6 +116,7 @@ class KafkaFlowAsinDetail(Templates): ...@@ -116,6 +116,7 @@ class KafkaFlowAsinDetail(Templates):
self.df_self_brand = self.spark.sql("select 1+1;") self.df_self_brand = self.spark.sql("select 1+1;")
self.df_hide_category = self.spark.sql("select 1+1;") self.df_hide_category = self.spark.sql("select 1+1;")
self.df_ai_hide_category = self.spark.sql("select 1+1;") self.df_ai_hide_category = self.spark.sql("select 1+1;")
self.df_ai_asin_exists = self.spark.sql("select 1+1;")
self.color_set = set() self.color_set = set()
self.save_parent_asin_latest_detail_cache = None self.save_parent_asin_latest_detail_cache = None
self.save_asin_latest_detail_cache = None self.save_asin_latest_detail_cache = None
...@@ -1040,6 +1041,20 @@ class KafkaFlowAsinDetail(Templates): ...@@ -1040,6 +1041,20 @@ class KafkaFlowAsinDetail(Templates):
username=mysql_con['username'], query=sql username=mysql_con['username'], query=sql
).repartition(self.repartition_num).persist(StorageLevel.DISK_ONLY) ).repartition(self.repartition_num).persist(StorageLevel.DISK_ONLY)
self.df_asin_auction.show(10, truncate=False) self.df_asin_auction.show(10, truncate=False)
print("13. 读取信息库已存在ASIN(Doris analyze_detail ∪ PG ai_asin_detail),用于导出去重")
# 仅latest+normal模式才会执行save_ai_asin_detail导出,其余模式跳过这两次全表IO
if self.test_flag == 'normal' and self.consumer_type == 'latest':
sql_analyzed = f"SELECT asin FROM `{self.doris_db_selection}`.`{self.site_name}_ai_asin_analyze_detail`"
df_analyzed_asin = DorisHelper.spark_import_with_sql(self.spark, sql_analyzed)
pg_con = DBUtil.get_connection_info("postgresql", "us")
sql_pg_exists = f"SELECT asin FROM {self.site_name}_ai_asin_detail"
df_pg_exists_asin = SparkUtil.read_jdbc_query(
session=self.spark, url=pg_con['url'], pwd=pg_con['pwd'],
username=pg_con['username'], query=sql_pg_exists
)
self.df_ai_asin_exists = df_analyzed_asin.select("asin").unionByName(df_pg_exists_asin.select("asin")) \
.dropDuplicates(['asin']).repartition(self.repartition_num, 'asin').persist(StorageLevel.DISK_ONLY)
self.df_ai_asin_exists.show(10, truncate=False)
# 字段处理逻辑综合 # 字段处理逻辑综合
def handle_all_field(self, df): def handle_all_field(self, df):
...@@ -1221,9 +1236,9 @@ class KafkaFlowAsinDetail(Templates): ...@@ -1221,9 +1236,9 @@ class KafkaFlowAsinDetail(Templates):
df_info_base = df.filter( df_info_base = df.filter(
"asin_type in (0, 1) and asin_bought_month >= 50" "asin_type in (0, 1) and asin_bought_month >= 50"
).join( ).join(
self.df_ai_hide_category, self.df_ai_hide_category, df["asin_bs_cate_current_id"] == self.df_ai_hide_category["category_id"], 'left_anti'
df["asin_bs_cate_current_id"] == self.df_ai_hide_category["category_id"], ).join(
'left_anti' self.df_ai_asin_exists, 'asin', 'left_anti'
) )
df_info_base = df_info_base.select( df_info_base = df_info_base.select(
F.lit(self.site_name).alias('site_name'), F.lit(self.site_name).alias('site_name'),
......
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