Commit 0f7e622c by chenyuanjie

月度信息库-增加过滤

parent 662adc43
...@@ -189,20 +189,29 @@ class DwtAiAsin(object): ...@@ -189,20 +189,29 @@ class DwtAiAsin(object):
partitionBy=['site_name', 'date_type', 'date_info'] partitionBy=['site_name', 'date_type', 'date_info']
) )
# ── PG 写入 {site}_ai_asin_detail,只导出增量:不在 Doris selection.{site}_ai_asin_analyze_detail 里的才需要── # ── PG 写入 {site}_ai_asin_detail,只导出增量:不在 Doris selection.{site}_ai_asin_analyze_detail 里
# (已完成AI分析)、且不在PG本表里(之前已经导出过,避免同一asin反复导出产生重复行)的才需要──
pg_con = DBUtil.get_connection_info("postgresql", "us")
export_pg_tb = f"{self.site_name}_ai_asin_detail"
sql_analyzed = f"SELECT asin FROM `selection`.`{self.site_name}_ai_asin_analyze_detail`" sql_analyzed = f"SELECT asin FROM `selection`.`{self.site_name}_ai_asin_analyze_detail`"
df_analyzed_asin = DorisHelper.spark_import_with_sql(self.spark, sql_analyzed).repartition(40, 'asin') df_analyzed_asin = DorisHelper.spark_import_with_sql(self.spark, sql_analyzed).repartition(40, 'asin')
sql_exported = f"SELECT asin FROM {export_pg_tb}"
df_exported_asin = SparkUtil.read_jdbc_query(
session=self.spark, url=pg_con["url"], pwd=pg_con["pwd"], username=pg_con["username"], query=sql_exported
).repartition(40, 'asin')
# persist:下面先count()校验数量、再write()真正落库,两次action共用同一份物化结果, # persist:下面先count()校验数量、再write()真正落库,两次action共用同一份物化结果,
# 不然anti join(含一次Doris JDBC读取)会被重新算一遍 # 不然anti join(含两次JDBC读取)会被重新算一遍
df_to_export = self.df_save.drop( df_to_export = self.df_save.drop(
'date_type', 'date_info', 'account_addr', 'launch_time_type', 'is_new_flag', 'is_ascending_flag' 'date_type', 'date_info', 'account_addr', 'launch_time_type', 'is_new_flag', 'is_ascending_flag'
).join(df_analyzed_asin, 'asin', 'left_anti') \ ).join(df_analyzed_asin, 'asin', 'left_anti') \
.join(df_exported_asin, 'asin', 'left_anti') \
.persist(StorageLevel.DISK_ONLY) .persist(StorageLevel.DISK_ONLY)
export_count = df_to_export.count() export_count = df_to_export.count()
print(f"待导出PG的增量ASIN数量:{export_count}") print(f"待导出PG的增量ASIN数量:{export_count}")
pg_con = DBUtil.get_connection_info("postgresql", "us")
export_pg_tb = f"{self.site_name}_ai_asin_detail"
try: try:
if export_count == 0: if export_count == 0:
print("没有需要导出到PG的增量ASIN,跳过写入") print("没有需要导出到PG的增量ASIN,跳过写入")
......
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