Commit 890cf95f by hejiangming

增加搜索词市场筛选维度字段

parent 7ef1fd6b
import os import os
import sys import sys
import time import time
import json
sys.path.append(os.path.dirname(sys.path[0])) # 上级目录 sys.path.append(os.path.dirname(sys.path[0])) # 上级目录
from utils.templates import Templates from utils.templates import Templates
from pyspark.sql.window import Window from pyspark.sql.window import Window
from pyspark.sql import functions as F from pyspark.sql import functions as F
from pyspark.sql.types import IntegerType from pyspark.sql.types import IntegerType, StringType
from utils.db_util import DBUtil from utils.db_util import DBUtil
from utils.spark_util import SparkUtil from utils.spark_util import SparkUtil
from utils.common_util import CommonUtil from utils.common_util import CommonUtil
from yswg_utils.common_udf import udf_detect_phrase_reg from yswg_utils.common_udf import udf_detect_phrase_reg
# ============================================================
# filter_refinements 维度黑名单(全部小写)
# 【这是什么】filter_refinements 是爬虫抓的"该搜索词市场的筛选维度"(JSON:维度名 -> 值列表),
# 里面混了品牌/平台/物流/交易类维度(如 Local Stores、Subscribe & Save),前端只想展示
# 材质/风格/功能/场景等商品属性维度,所以要按黑名单把非属性维度整条丢掉。
# 【匹配规则】维度名转小写后,只要"包含"下列任一关键词 → 整个维度连同其值一起丢弃(业务确认按包含匹配)。
# 【'brand' 这条】单独一条覆盖所有带品牌的写法(Brands / Sub-Brand / Premium Brands /
# From Our Brands / All Top Brands / Top Brands / Top Brands in Health & Household 等),
# 等价于业务原始 22 项里那 6 个 brand 相关项,故不再逐条重复列。
# 【注意】'color'/'price'/'seller'/'condition' 也按包含丢(业务确认),会连带丢掉
# Frame Color / Skin Condition 这类含该词的维度——这是预期行为。
# 【怎么加词】后续要新增过滤维度,只在本列表加一行小写关键词即可。
# ============================================================
FILTER_REFINEMENTS_BLACKLIST = [
'brand', # 含 brand 全丢(覆盖原 22 项里 6 个 brand 相关项)
'eligible for free shipping',
'delivery day',
'color',
'price',
'deals & discounts',
'customer reviews',
'amazon fashion',
'care instructions',
'seller',
'customizable products',
'subscribe & save',
'handmade products',
'condition',
'amazon certified',
'amazon beauty',
'local stores',
]
class DwtAbaStAnalytics(Templates): class DwtAbaStAnalytics(Templates):
def __init__(self, site_name="us", date_type="week", date_info="2022-40"): def __init__(self, site_name="us", date_type="week", date_info="2022-40"):
...@@ -61,9 +96,16 @@ class DwtAbaStAnalytics(Templates): ...@@ -61,9 +96,16 @@ class DwtAbaStAnalytics(Templates):
# 自身历史 11 个月分区的 rank(峰值月 peak_month / 常年可卖 all_year_text_flag 计算用) # 自身历史 11 个月分区的 rank(峰值月 peak_month / 常年可卖 all_year_text_flag 计算用)
# 仅 month 流程在 read_data 阶段真正读取数据;非 month 流程下游统一用占位值填充 # 仅 month 流程在 read_data 阶段真正读取数据;非 month 流程下游统一用占位值填充
self.df_history_rank = self.spark.sql(f"select 1+1;") self.df_history_rank = self.spark.sql(f"select 1+1;")
# 搜索词市场筛选维度 df,来源 ods_st_quantity_being_sold.filter_refinements(清洗后 JSON)
# 仅 month 流程在 read_data 阶段真正读取数据;非 month 流程下游统一用 lit("{}") 填充
self.df_st_filter = self.spark.sql(f"select 1+1;")
# 自定义udf函数注册 # 自定义udf函数注册
self.u_contains = self.spark.udf.register('u_contains', self.udf_contains, IntegerType()) self.u_contains = self.spark.udf.register('u_contains', self.udf_contains, IntegerType())
# filter_refinements 清洗 UDF:解析 JSON→按黑名单删维度→重序列化,返回干净 JSON 串
self.u_clean_filter = self.spark.udf.register(
'u_clean_filter', self.udf_clean_filter_refinements, StringType()
)
self.u_judge_color = self.spark.udf.register('u_judge_color', self.udf_judge_color, IntegerType()) self.u_judge_color = self.spark.udf.register('u_judge_color', self.udf_judge_color, IntegerType())
self.u_judge_title_color = self.spark.udf.register( self.u_judge_title_color = self.spark.udf.register(
'u_judge_title_color', self.udf_judge_title_color, IntegerType() 'u_judge_title_color', self.udf_judge_title_color, IntegerType()
...@@ -151,6 +193,53 @@ class DwtAbaStAnalytics(Templates): ...@@ -151,6 +193,53 @@ class DwtAbaStAnalytics(Templates):
return 1 return 1
return 0 return 0
@staticmethod
def udf_clean_filter_refinements(s):
"""
清洗 filter_refinements:解析原始 JSON → 按黑名单丢维度 → 重序列化成干净 JSON 串。
空 / 异常 / 过滤后无维度剩余 → 返回 None(下游 na.fill 补 "{}")。
"""
# ============================================================
# Step1 解析:null / 空串 / 非法 JSON / 不是 dict 一律判空
# 为什么全判空:这些都取不出"维度名->值列表"结构,直接交给下游补 "{}"
# 例:s=None → None;s='' → None;s='abc'(非JSON) → None;s='[1,2]'(非dict) → None
# ============================================================
if s is None:
return None
s = s.strip()
if s == '':
return None
try:
obj = json.loads(s)
except Exception:
return None
if not isinstance(obj, dict):
return None
# ============================================================
# Step2 过滤:维度名转小写后"包含"任一黑名单词 → 整条维度(连同值)丢弃
# 为什么转小写再判断:黑名单匹配忽略大小写,'Top Brands'/'top brands' 都要命中
# 例:{'Flavor':[...],'Top Brands':[...],'Frame Color':[...]}
# → 'top brands' 含 'brand' 丢、'frame color' 含 'color' 丢 → 只留 {'Flavor':[...]}
# ============================================================
kept = {}
for k, v in obj.items():
kl = str(k).lower()
if any(bad in kl for bad in FILTER_REFINEMENTS_BLACKLIST): # 含任一黑名单词=非商品属性维度
continue
kept[k] = v
# Step3 过滤后一个维度都不剩 → 判空(如整条 JSON 全是品牌/平台维度)
if not kept:
return None
# ============================================================
# Step4 重序列化成 JSON 串
# ensure_ascii=False:保留原字符(维度名/值可能含非 ASCII),否则会变成 \uXXXX 转义
# 不加 sort_keys:保持爬虫原始维度顺序(PG 端为 jsonb,键顺序本就会被 jsonb 重排,排序无意义)
# ============================================================
return json.dumps(kept, ensure_ascii=False)
def read_data(self): def read_data(self):
# 一些不涵盖month_old的分区,重定义成month,其他正常 # 一些不涵盖month_old的分区,重定义成month,其他正常
spe_date_type = 'month' if 'month_old' == self.date_type else self.date_type spe_date_type = 'month' if 'month_old' == self.date_type else self.date_type
...@@ -166,7 +255,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -166,7 +255,7 @@ class DwtAbaStAnalytics(Templates):
self.df_st_key = self.spark.sql(sqlQuery=sql) self.df_st_key = self.spark.sql(sqlQuery=sql)
self.df_st_key = self.df_st_key.repartition(80, 'search_term').cache() self.df_st_key = self.df_st_key.repartition(80, 'search_term').cache()
print("self.df_st_key:") print("self.df_st_key:")
self.df_st_key.show(10, truncate=True) # self.df_st_key.show(10, truncate=True)
# 获取dwd_st_measure 事实表 # 获取dwd_st_measure 事实表
sql = f""" sql = f"""
...@@ -217,7 +306,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -217,7 +306,7 @@ class DwtAbaStAnalytics(Templates):
self.df_st_asin_measure = self.spark.sql(sqlQuery=sql) self.df_st_asin_measure = self.spark.sql(sqlQuery=sql)
self.df_st_asin_measure = self.df_st_asin_measure.repartition(80, 'asin').cache() self.df_st_asin_measure = self.df_st_asin_measure.repartition(80, 'asin').cache()
print("self.df_st_asin_measure:") print("self.df_st_asin_measure:")
self.df_st_asin_measure.show(10, truncate=True) # self.df_st_asin_measure.show(10, truncate=True)
# 获取dwd_asin_measure 事实表 # 获取dwd_asin_measure 事实表
sql = f""" sql = f"""
...@@ -234,7 +323,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -234,7 +323,7 @@ class DwtAbaStAnalytics(Templates):
self.df_asin_measure = self.spark.sql(sqlQuery=sql) self.df_asin_measure = self.spark.sql(sqlQuery=sql)
self.df_asin_measure = self.df_asin_measure.repartition(80, 'asin').cache() self.df_asin_measure = self.df_asin_measure.repartition(80, 'asin').cache()
print("self.df_asin_measure:") print("self.df_asin_measure:")
self.df_asin_measure.show(10, truncate=True) # self.df_asin_measure.show(10, truncate=True)
# 获取dim_asin_detail表 # 获取dim_asin_detail表
sql = f""" sql = f"""
...@@ -274,7 +363,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -274,7 +363,7 @@ class DwtAbaStAnalytics(Templates):
self.df_asin_detail = self.spark.sql(sqlQuery=sql) self.df_asin_detail = self.spark.sql(sqlQuery=sql)
self.df_asin_detail = self.df_asin_detail.repartition(80, 'asin').cache() self.df_asin_detail = self.df_asin_detail.repartition(80, 'asin').cache()
print("self.df_asin_detail:") print("self.df_asin_detail:")
self.df_asin_detail.show(10, truncate=True) # self.df_asin_detail.show(10, truncate=True)
# 仅获取 asin和country_name,对country_name进行了聚合处理 # 仅获取 asin和country_name,对country_name进行了聚合处理
sql = f""" sql = f"""
...@@ -288,7 +377,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -288,7 +377,7 @@ class DwtAbaStAnalytics(Templates):
self.df_seller_asin_country = self.spark.sql(sqlQuery=sql) self.df_seller_asin_country = self.spark.sql(sqlQuery=sql)
self.df_seller_asin_country = self.df_seller_asin_country.repartition(80, 'asin').cache() self.df_seller_asin_country = self.df_seller_asin_country.repartition(80, 'asin').cache()
print("self.df_seller_asin_country:") print("self.df_seller_asin_country:")
self.df_seller_asin_country.show(10, truncate=True) # self.df_seller_asin_country.show(10, truncate=True)
# 获取 dim_fd_asin_info 表 # 获取 dim_fd_asin_info 表
sql = f""" sql = f"""
...@@ -302,7 +391,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -302,7 +391,7 @@ class DwtAbaStAnalytics(Templates):
self.df_seller_asin_info = self.spark.sql(sqlQuery=sql) self.df_seller_asin_info = self.spark.sql(sqlQuery=sql)
self.df_seller_asin_info = self.df_seller_asin_info.drop_duplicates(['asin']).repartition(80, 'asin').cache() self.df_seller_asin_info = self.df_seller_asin_info.drop_duplicates(['asin']).repartition(80, 'asin').cache()
print("self.df_seller_asin_info:") print("self.df_seller_asin_info:")
self.df_seller_asin_info.show(10, truncate=True) # self.df_seller_asin_info.show(10, truncate=True)
# 获取 dim_st_detail asin1-3共享点击信息表 # 获取 dim_st_detail asin1-3共享点击信息表
sql = f""" sql = f"""
...@@ -352,7 +441,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -352,7 +441,7 @@ class DwtAbaStAnalytics(Templates):
self.df_st_detail = self.df_st_detail.withColumn(col, F.regexp_replace(F.col(col), '\x00', '')) self.df_st_detail = self.df_st_detail.withColumn(col, F.regexp_replace(F.col(col), '\x00', ''))
self.df_st_detail = self.df_st_detail.repartition(80, 'search_term').cache() self.df_st_detail = self.df_st_detail.repartition(80, 'search_term').cache()
print("self.df_st_detail:") print("self.df_st_detail:")
self.df_st_detail.show(10, truncate=True) # self.df_st_detail.show(10, truncate=True)
# 获取dws_st_num_stats表 取max_num、most_proportion # 获取dws_st_num_stats表 取max_num、most_proportion
sql = f""" sql = f"""
...@@ -371,7 +460,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -371,7 +460,7 @@ class DwtAbaStAnalytics(Templates):
self.df_st_num_stats = self.spark.sql(sqlQuery=sql) self.df_st_num_stats = self.spark.sql(sqlQuery=sql)
self.df_st_num_stats = self.df_st_num_stats.repartition(80, 'search_term').cache() self.df_st_num_stats = self.df_st_num_stats.repartition(80, 'search_term').cache()
print("self.df_st_num_stats:") print("self.df_st_num_stats:")
self.df_st_num_stats.show(10, truncate=True) # self.df_st_num_stats.show(10, truncate=True)
# 获取dwt_st_market表 取market_cycle_type # 获取dwt_st_market表 取market_cycle_type
sql = f""" sql = f"""
...@@ -386,7 +475,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -386,7 +475,7 @@ class DwtAbaStAnalytics(Templates):
self.df_st_market = self.spark.sql(sqlQuery=sql) self.df_st_market = self.spark.sql(sqlQuery=sql)
self.df_st_market = self.df_st_market.repartition(80, 'search_term').cache() self.df_st_market = self.df_st_market.repartition(80, 'search_term').cache()
print("self.df_st_market:") print("self.df_st_market:")
self.df_st_market.show(10, truncate=True) # self.df_st_market.show(10, truncate=True)
# 获取dwd_st_volume_fba 取gross_profit_fee_air 和 gross_profit_fee_sea # 获取dwd_st_volume_fba 取gross_profit_fee_air 和 gross_profit_fee_sea
# sql = f""" # sql = f"""
...@@ -417,7 +506,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -417,7 +506,7 @@ class DwtAbaStAnalytics(Templates):
self.df_asin_label = self.spark.sql(sqlQuery=sql) self.df_asin_label = self.spark.sql(sqlQuery=sql)
self.df_asin_label = self.df_asin_label.repartition(80, 'asin').cache() self.df_asin_label = self.df_asin_label.repartition(80, 'asin').cache()
print("self.df_asin_label:") print("self.df_asin_label:")
self.df_asin_label.show(10, truncate=True) # self.df_asin_label.show(10, truncate=True)
# 获取品牌词库 # 获取品牌词库
sql = f""" sql = f"""
...@@ -433,7 +522,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -433,7 +522,7 @@ class DwtAbaStAnalytics(Templates):
self.df_st_brand = self.spark.sql(sqlQuery=sql) self.df_st_brand = self.spark.sql(sqlQuery=sql)
self.df_st_brand = self.df_st_brand.repartition(80, 'search_term').cache() self.df_st_brand = self.df_st_brand.repartition(80, 'search_term').cache()
print("self.df_st_brand:") print("self.df_st_brand:")
self.df_st_brand.show(10, truncate=True) # self.df_st_brand.show(10, truncate=True)
# 从pgsql获取特殊字符匹配字典表:match_character_dict # 从pgsql获取特殊字符匹配字典表:match_character_dict
pg_sql = f""" pg_sql = f"""
...@@ -469,7 +558,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -469,7 +558,7 @@ class DwtAbaStAnalytics(Templates):
self.df_is_hidden_cate = self.spark.sql(sqlQuery=sql) self.df_is_hidden_cate = self.spark.sql(sqlQuery=sql)
self.df_is_hidden_cate = self.df_is_hidden_cate.repartition(80).cache() self.df_is_hidden_cate = self.df_is_hidden_cate.repartition(80).cache()
print("self.df_is_hidden_cate:") print("self.df_is_hidden_cate:")
self.df_is_hidden_cate.show(10, truncate=True) # self.df_is_hidden_cate.show(10, truncate=True)
# asin利润率 # asin利润率
sql = f""" sql = f"""
...@@ -478,7 +567,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -478,7 +567,7 @@ class DwtAbaStAnalytics(Templates):
""" """
self.df_asin_profit_rate = self.spark.sql(sqlQuery=sql).repartition(80, 'asin').dropDuplicates(['asin', 'asin_price']).cache() self.df_asin_profit_rate = self.spark.sql(sqlQuery=sql).repartition(80, 'asin').dropDuplicates(['asin', 'asin_price']).cache()
print("self.df_asin_profit_rate:") print("self.df_asin_profit_rate:")
self.df_asin_profit_rate.show(10, truncate=True) # self.df_asin_profit_rate.show(10, truncate=True)
if self.date_type == 'month': if self.date_type == 'month':
# 读累加表 dim_st_detail_history,用于判断 is_first_ever_text(全历史首次出现) # 读累加表 dim_st_detail_history,用于判断 is_first_ever_text(全历史首次出现)
...@@ -500,7 +589,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -500,7 +589,7 @@ class DwtAbaStAnalytics(Templates):
# 后续 join 上 → is_first_ever_text=0;join 不上 → is_first_ever_text=1(全历史首次) # 后续 join 上 → is_first_ever_text=0;join 不上 → is_first_ever_text=1(全历史首次)
self.df_history_st = self.spark.sql(sqlQuery=sql).repartition(80, 'search_term').cache() self.df_history_st = self.spark.sql(sqlQuery=sql).repartition(80, 'search_term').cache()
print("self.df_history_st:") print("self.df_history_st:")
self.df_history_st.show(10, truncate=True) # self.df_history_st.show(10, truncate=True)
# 读 dws_st_theme 计算搜索词属性标签 st_attribute_label # 读 dws_st_theme 计算搜索词属性标签 st_attribute_label
has_theme_data = self.spark.sql(f""" has_theme_data = self.spark.sql(f"""
...@@ -529,7 +618,7 @@ class DwtAbaStAnalytics(Templates): ...@@ -529,7 +618,7 @@ class DwtAbaStAnalytics(Templates):
""" """
self.df_st_attribute = self.spark.sql(sqlQuery=sql).repartition(80, 'search_term').cache() self.df_st_attribute = self.spark.sql(sqlQuery=sql).repartition(80, 'search_term').cache()
print("self.df_st_attribute:") print("self.df_st_attribute:")
self.df_st_attribute.show(10, truncate=True) # self.df_st_attribute.show(10, truncate=True)
# ============================================================ # ============================================================
# 读自身历史 11 个月分区的 rank,用于计算峰值月 peak_month / 常年可卖 all_year_text_flag # 读自身历史 11 个月分区的 rank,用于计算峰值月 peak_month / 常年可卖 all_year_text_flag
...@@ -563,7 +652,37 @@ class DwtAbaStAnalytics(Templates): ...@@ -563,7 +652,37 @@ class DwtAbaStAnalytics(Templates):
""" """
self.df_history_rank = self.spark.sql(sqlQuery=sql).repartition(80, 'search_term').cache() self.df_history_rank = self.spark.sql(sqlQuery=sql).repartition(80, 'search_term').cache()
print("self.df_history_rank:") print("self.df_history_rank:")
self.df_history_rank.show(10, truncate=True) # self.df_history_rank.show(10, truncate=True)
# ============================================================
# 读 ods_st_quantity_being_sold.filter_refinements,清洗出 st_filter_refinements
# 【业务背景】前端要在 ABA 搜索词页展示该词所在市场的"筛选维度"(材质/风格/功能/场景等商品属性),
# 数据源是爬虫抓的 filter_refinements(JSON:维度名->值列表),需过滤掉品牌/平台/物流/交易类维度
# 【为什么读 ODS】该字段在 DIM 层(dim_st_detail)已被丢弃,只能回 ODS 原始表取
# 【为什么要去重】ods_st_quantity_being_sold 未按 search_term 去重,同一词多次抓取会有多行、
# 各行 filter_refinements 值和抓取时间都不同 → 按 created_time 降序只取最新一行(对齐 dwd_st_brand_badge)
# 【数据边界】该字段 2026-06 才开始抓,之前的月份 filter_refinements 全为 null → 清洗后统一为 "{}",无害
# ============================================================
sql = f"""
select search_term, filter_refinements, created_time
from ods_st_quantity_being_sold
where site_name = '{self.site_name}'
and date_type = '{self.date_type}'
and date_info = '{self.date_info}'
and search_term is not null
"""
df_fr = self.spark.sql(sqlQuery=sql)
# 按 search_term 分区、created_time 降序,row_number=1 只保留最新一行
# 【只按时间取最新、不看内容】即便最新行 filter_refinements 为空也取它(业务确认口径);
# 放在清洗前:去重只看时间与内容无关,先去重能让每个词只跑一次清洗 UDF,避免无用计算
w_fr = Window.partitionBy('search_term').orderBy(F.col('created_time').desc_nulls_last())
self.df_st_filter = df_fr.withColumn('_rn', F.row_number().over(w_fr)) \
.filter(F.col('_rn') == 1) \
.withColumn('st_filter_refinements', self.u_clean_filter(F.col('filter_refinements'))) \
.select('search_term', 'st_filter_refinements') \
.repartition(80, 'search_term').cache()
print("self.df_st_filter:")
# self.df_st_filter.show(10, truncate=True)
def handle_data(self): def handle_data(self):
# 对基础计算表进行关联 # 对基础计算表进行关联
...@@ -581,26 +700,29 @@ class DwtAbaStAnalytics(Templates): ...@@ -581,26 +700,29 @@ class DwtAbaStAnalytics(Templates):
# 语种处理 # 语种处理
self.handle_calc_lang() self.handle_calc_lang()
# ============================================================
# month-only 字段统一处理:month 流程计算/关联,非 month 流程统一 lit 占位
# 涉及 5 个字段(都只在 month 有业务意义):
# - is_first_ever_text : 全历史首次出现标记(handle_first_ever_flag,基于累加表 dim_st_detail_history)
# - peak_month / all_year_text_flag : 峰值月 + 常年可卖(handle_peak_month,基于近 12 个月 rank)
# - st_attribute_label : 属性标签(left join read_data 读好的 dws_st_theme 聚合,占位 "-1")
# - st_filter_refinements : 市场筛选维度(left join read_data 清洗好的 ods filter_refinements,占位 "{}")
# join 不上的词在 handle_column 阶段各自 na.fill 占位;非 month 直接 lit 占位(占位值语义见各行注释)
# ============================================================
if self.date_type == 'month': if self.date_type == 'month':
self.handle_first_ever_flag() # 全历史首次出现标记(基于累加表 dim_st_detail_history) 不是月流程填充 -1 self.handle_first_ever_flag()
self.handle_peak_month() # 峰值月 + 常年可卖(基于自身历史 11 个月分区 + 当月 rank) self.handle_peak_month()
else:
self.df_save = self.df_save.withColumn('is_first_ever_text', F.lit(-1))
# 非 month 流程业务不需要峰值月/常年可卖,直接占位(peak_month 空串→PG 转空数组 {},all_year_text_flag 占位 -1)
self.df_save = self.df_save.withColumn(
'peak_month', F.lit('')
).withColumn(
'all_year_text_flag', F.lit(-1)
)
# 附加属性标签字段 st_attribute_label
# month 流程:left join read_data 阶段已读取的 dws_st_theme 聚合结果,join 不上的词在 handle_column 阶段 fillna("-1")
# 非 month 流程:业务不需要计算,直接 lit("-1") 占位(与 PG 端 '{-1}' 数组占位语义一致)
if self.date_type == 'month':
self.df_save = self.df_save.join(self.df_st_attribute, on='search_term', how='left') self.df_save = self.df_save.join(self.df_st_attribute, on='search_term', how='left')
self.df_st_attribute.unpersist() self.df_st_attribute.unpersist()
self.df_save = self.df_save.join(self.df_st_filter, on='search_term', how='left')
self.df_st_filter.unpersist()
else: else:
self.df_save = self.df_save.withColumn('st_attribute_label', F.lit('-1')) self.df_save = self.df_save \
.withColumn('is_first_ever_text', F.lit(-1)) \
.withColumn('peak_month', F.lit('')) \
.withColumn('all_year_text_flag', F.lit(-1)) \
.withColumn('st_attribute_label', F.lit('-1')) \
.withColumn('st_filter_refinements', F.lit('{}'))
# 处理输出字段 # 处理输出字段
self.handle_column() self.handle_column()
...@@ -1108,7 +1230,8 @@ class DwtAbaStAnalytics(Templates): ...@@ -1108,7 +1230,8 @@ class DwtAbaStAnalytics(Templates):
"seller_asin_proportion", # 前三页ASIN数最多卖家的ASIN数占比 "seller_asin_proportion", # 前三页ASIN数最多卖家的ASIN数占比
"st_attribute_label", # 搜索词属性标签(逗号分隔字符串,sqoop 导出 PG 后转 VARCHAR[]) "st_attribute_label", # 搜索词属性标签(逗号分隔字符串,sqoop 导出 PG 后转 VARCHAR[])
"peak_month", # 峰值月:近12个月 rank 最小的月份,并列保留,'YYYY-MM' 逗号分隔(PG 转 VARCHAR[]) "peak_month", # 峰值月:近12个月 rank 最小的月份,并列保留,'YYYY-MM' 逗号分隔(PG 转 VARCHAR[])
"all_year_text_flag" # 常年可卖:近12个月每月都出现在 ABA 月度搜索词中=1,否则=0 "all_year_text_flag", # 常年可卖:近12个月每月都出现在 ABA 月度搜索词中=1,否则=0
"st_filter_refinements" # 搜索词市场筛选维度(清洗后 JSON:维度名->值列表,已过滤黑名单维度;PG 端 jsonb,空为 {})
) )
# 空值处理 # 空值处理
...@@ -1129,7 +1252,8 @@ class DwtAbaStAnalytics(Templates): ...@@ -1129,7 +1252,8 @@ class DwtAbaStAnalytics(Templates):
"seller_asin_proportion": -1, # 分子为 null(搜索词全无账号) → 占位 -1 "seller_asin_proportion": -1, # 分子为 null(搜索词全无账号) → 占位 -1
"st_attribute_label": "-1", # 词典无匹配 → 占位 "-1"(Java 侧转 null 返前端) "st_attribute_label": "-1", # 词典无匹配 → 占位 "-1"(Java 侧转 null 返前端)
"peak_month": "", # 极端兜底(rank 全 null)→ 空串,PG string_to_array 转空数组 {} "peak_month": "", # 极端兜底(rank 全 null)→ 空串,PG string_to_array 转空数组 {}
"all_year_text_flag": 0 # 理论必有值,兜底 0(非 month 流程已提前填 -1,不受影响) "all_year_text_flag": 0, # 理论必有值,兜底 0(非 month 流程已提前填 -1,不受影响)
"st_filter_refinements": "{}" # 空/异常/无维度剩余/join不上 → 空对象 "{}"(PG jsonb)
}) })
# 日期字段补全 # 日期字段补全
......
...@@ -197,7 +197,11 @@ if __name__ == '__main__': ...@@ -197,7 +197,11 @@ if __name__ == '__main__':
# 同 st_attribute_label:copy 表先 ALTER 成 VARCHAR 让 sqoop 写字符串,交换前再 string_to_array 转回 VARCHAR[] # 同 st_attribute_label:copy 表先 ALTER 成 VARCHAR 让 sqoop 写字符串,交换前再 string_to_array 转回 VARCHAR[]
"peak_month", "peak_month",
# 常年可卖标记:标量 int,sqoop 直写,无需中转 # 常年可卖标记:标量 int,sqoop 直写,无需中转
"all_year_text_flag" "all_year_text_flag",
# 搜索词市场筛选维度:Hive 端是清洗后 JSON STRING(如 '{"Flavor":["Chocolate"]}'),PG 端是 jsonb
# Sqoop 不能直接写 jsonb:copy 表先把该列 ALTER 成 varchar 让 sqoop 写字符串,
# 交换前再 ALTER 回 jsonb(USING ...::jsonb),同 st_attribute_label 的中转思路
"st_filter_refinements"
] ]
# 处理导出表 # 处理导出表
export_master_tb = f"{export_base_tb}_{date_type}_{year_str}" export_master_tb = f"{export_base_tb}_{date_type}_{year_str}"
...@@ -239,9 +243,14 @@ if __name__ == '__main__': ...@@ -239,9 +243,14 @@ if __name__ == '__main__':
# copy 表继承自正式分区表(含 st_attribute_label VARCHAR[]), # copy 表继承自正式分区表(含 st_attribute_label VARCHAR[]),
# 但 Sqoop 不支持直接写入 PG 数组类型,必须先把 copy 表的该列临时改成 VARCHAR # 但 Sqoop 不支持直接写入 PG 数组类型,必须先把 copy 表的该列临时改成 VARCHAR
# 等 Sqoop 完成后、分区交换之前,再 ALTER 回 VARCHAR[](见下方 exchange_pg_part_tb 前的处理) # 等 Sqoop 完成后、分区交换之前,再 ALTER 回 VARCHAR[](见下方 exchange_pg_part_tb 前的处理)
# st_filter_refinements 在 master/copy 里是 jsonb 且带 DEFAULT '{}'::jsonb,
# Sqoop 不能写 jsonb → 先 DROP DEFAULT(避免改类型时默认值转换报错)再 ALTER 成 varchar
# USING st_filter_refinements::text 把 jsonb 显式转字符串;varchar 不限长,避免长 JSON 被截断
sql_alter_to_varchar = f""" sql_alter_to_varchar = f"""
ALTER TABLE {export_tb_copy} ALTER COLUMN st_attribute_label TYPE VARCHAR(200); ALTER TABLE {export_tb_copy} ALTER COLUMN st_attribute_label TYPE VARCHAR(200);
ALTER TABLE {export_tb_copy} ALTER COLUMN peak_month TYPE VARCHAR(200); ALTER TABLE {export_tb_copy} ALTER COLUMN peak_month TYPE VARCHAR(200);
ALTER TABLE {export_tb_copy} ALTER COLUMN st_filter_refinements DROP DEFAULT;
ALTER TABLE {export_tb_copy} ALTER COLUMN st_filter_refinements TYPE VARCHAR USING st_filter_refinements::text;
""" """
DBUtil.engine_exec_sql(engine, sql_alter_to_varchar) DBUtil.engine_exec_sql(engine, sql_alter_to_varchar)
...@@ -320,6 +329,10 @@ if __name__ == '__main__': ...@@ -320,6 +329,10 @@ if __name__ == '__main__':
ALTER TABLE {export_tb_copy} ALTER TABLE {export_tb_copy}
ALTER COLUMN peak_month TYPE VARCHAR[] ALTER COLUMN peak_month TYPE VARCHAR[]
USING string_to_array(coalesce(peak_month, ''), ',')::varchar[]; USING string_to_array(coalesce(peak_month, ''), ',')::varchar[];
ALTER TABLE {export_tb_copy}
ALTER COLUMN st_filter_refinements TYPE jsonb
USING st_filter_refinements::jsonb;
""" """
DBUtil.engine_exec_sql(engine, sql_alter_back) DBUtil.engine_exec_sql(engine, sql_alter_back)
......
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