Commit 90ede8f8 by chenyuanjie

流量选品每日刷新任务-调整为部分列更新

parent 30f81277
"""
@Author : HuangJian
@SourceTable :
①ods_seller_account_syn
②ods_seller_asin_account
③ods_seller_account_feedback
@SinkTable :
①dim_fd_asin_info
@CreateTime : 2022/12/19 9:56
@UpdateTime : 2022/12/19 9:56
"""
import os
import sys
sys.path.append(os.path.dirname(sys.path[0])) # 上级目录
from pyspark.sql.window import Window
from pyspark.sql import functions as F
from utils.spark_util import SparkUtil
from utils.hdfs_utils import HdfsUtils
from utils.common_util import CommonUtil
class DwtFdAsinInfo(object):
def __init__(self, site_name='us', date_type='month', date_info='2026-06'):
super().__init__()
self.hive_tb = "dim_fd_asin_info"
self.site_name = site_name
self.date_type = date_type
self.date_info = date_info
self.partition_dict = {
"site_name": site_name,
}
# 落表路径校验
self.hdfs_path = CommonUtil.build_hdfs_path(self.hive_tb, partition_dict=self.partition_dict)
app_name = f"{self.hive_tb}:{self.site_name}_{self.date_type}_{self.date_info}"
self.spark = SparkUtil.get_spark_session(app_name)
self.partitions_num = CommonUtil.reset_partitions(self.site_name, 80)
# 初始化全局变量df--ods获取数据的原始df
self.df_seller_account_syn = self.spark.sql("select 1+1;")
self.df_seller_account_feedback = self.spark.sql("select 1+1;")
self.df_fd_asin = self.spark.sql("select 1+1;")
# 初始化全局变量df--dwd层转换输出的df
self.df_save = self.spark.sql(f"select 1+1;")
def read_data(self):
# 获取爬虫店铺记录表
print("获取 ods_seller_account_syn")
sql = f"""
select id as fd_account_id, seller_id as unique_id, account_name as fd_account_name,
lower(account_name) as fd_account_name_lower, url as fd_url
from ods_seller_account_syn where site_name = '{self.site_name}'
"""
# seller_id 本身在 ods_seller_account_syn 里是唯一的,不需要再去重
self.df_seller_account_syn = self.spark.sql(sqlQuery=sql)
print(sql)
# 获取店铺详情表:只读传参指定的这一个分区(date_type+date_info 即调度传入的最新分区)
print("获取 ods_seller_account_feedback")
sql = f"""
select seller_id as unique_id, country_name as fd_country_name, created_at
from ods_seller_account_feedback
where site_name = '{self.site_name}' and date_type = '{self.date_type}' and date_info = '{self.date_info}'
"""
self.df_seller_account_feedback = self.spark.sql(sqlQuery=sql)
print(sql)
window = Window.partitionBy('unique_id').orderBy(F.col('created_at').desc())
self.df_seller_account_feedback = self.df_seller_account_feedback.withColumn(
'rank', F.row_number().over(window)
).filter('rank = 1').withColumn(
'fb_crawl_date', F.date_format(F.col('created_at'), 'yyyy-MM-dd HH:mm:ss')
).drop('rank', 'created_at')
# 获取店铺与asin的对应关系库(店铺与asin所有历史对应关系表)
print("获取 ods_seller_asin_account")
sql = f"""
select seller_id as unique_id, asin from ods_seller_asin_account where site_name='{self.site_name}'
"""
self.df_fd_asin = self.spark.sql(sqlQuery=sql)
self.df_fd_asin = self.df_fd_asin.drop_duplicates(['unique_id', 'asin'])
print(sql)
def save_data(self):
df_save = self.df_seller_account_syn.join(
self.df_seller_account_feedback, on='unique_id', how='left'
).join(
self.df_fd_asin, on='unique_id', how='left'
)
df_save = df_save.select(
F.col('fd_account_id'),
F.col('unique_id').alias('fd_unique'),
F.col('fd_account_name'),
F.col('fd_account_name_lower'),
F.col('fd_country_name'),
F.col('fd_url'),
F.col('asin'),
F.date_format(F.current_timestamp(), 'yyyy-MM-dd HH:mm:ss').alias('created_at'),
F.date_format(F.current_timestamp(), 'yyyy-MM-dd HH:mm:ss').alias('updated_at'),
F.col('fb_crawl_date'),
F.lit(self.site_name).alias('site_name'),
)
print(f"清除hdfs目录中:{self.hdfs_path}")
HdfsUtils.delete_file_in_folder(self.hdfs_path)
df_save = df_save.repartition(self.partitions_num)
partition_by = ["site_name"]
print(f"当前存储的表名为:{self.hive_tb},分区为{partition_by}", )
df_save.write.saveAsTable(name=self.hive_tb, format='hive', mode='append', partitionBy=partition_by)
print("success")
def run(self):
self.read_data()
self.save_data()
if __name__ == '__main__':
site_name = sys.argv[1] # 参数1:站点
date_type = sys.argv[2] # 参数2:类型:week/4_week/month/quarter
date_info = sys.argv[3] # 参数3:年-周/年-月/年-季, 比如: 2022-1
handle_obj = DwtFdAsinInfo(site_name=site_name, date_type=date_type, date_info=date_info)
handle_obj.run()
"""
@Author : CT
@Description : 流量选品每日刷新任务合并脚本——把原来分散的几个"每日刷新"任务合并到一起,
统一只处理最新一个月,且全部改成"发生变化才更新"的部分列更新,减少 Doris 压力
【背景】
selection.{site}_flow_asin_month_{yyyy_mm}(月表) / selection.{site}_flow_asin_365day(年表)
这两张物化表里,只有以下 9 个字段依赖 dwd_asin_profit_rate_latest(利润率) /
dwd_keepa_asin_detail(Keepa) 这两张每日刷新的源表,其余几十个字段都是月度/年度批跑时
算好、不需要每天动的:
ocean_profit, air_profit,
launch_time, launch_time_type,
tracking_since, tracking_since_type,
package_length, package_width, package_height, item_weight
月表/年表都是 UNIQUE KEY(asin) + enable_unique_key_merge_on_write=true,支持部分列更新
(INSERT INTO table(asin, 这9个字段) SELECT ...),不需要整表/整分区重算。
【4 个子任务,main() 里按顺序执行】
1) 刷新月表:selection.{site}_flow_asin_month_{最新月},只更新这9个字段
2) 刷新年表:selection.{site}_flow_asin_365day,限定 latest_date_info=最新月 这批 asin
3) 刷新利润率趋势中间表:dwt.dwt_asin_profit_rate_history,只处理最新1个月(原来是近3月)
4) Spark 聚合中间表全量历史 → 趋势物理表 selection.{site}_asin_profit_rate_trend
1/2/3 都只在"新值和已存值不一致"时才真正写入(Doris 支持 NULL-safe 的 <=>,已在生产
环境验证过 有值<=>NULL 会判定为不相等),避免每天对没有变化的行也触发一次 MOW 合并写;
4 本身就是全量重聚合一遍很小的物理表,不受"只刷新最新月"影响,逻辑不变。
【站点范围】
目前只支持 us 站点(这几个每日变化字段目前都是 us 站点专属逻辑)
【无入参】
不接收 sys.argv,所有逻辑动态推断;海豚定时调度直接执行 spark-submit 即可
执行示例:
spark-submit dws_flow_asin_refresh.py
"""
import os
import sys
sys.path.append(os.path.dirname(sys.path[0]))
from pyspark.sql import functions as F, Window
from utils.spark_util import SparkUtil
from utils.db_util import DBUtil
from utils.DorisHelper import DorisHelper
SITE_NAME = 'us' # 目前只支持 us 站点
# ===== Doris 连接 / 执行 =====
def _doris_connect(use_type='selection'):
"""统一 pymysql 连接,database='selection' 让 UDF / selection 库可见"""
import pymysql
conn_info = DorisHelper.get_connection_info(use_type)
return pymysql.connect(
host=conn_info['ip'],
port=conn_info['jdbc_port'],
user=conn_info['user'],
password=conn_info['pwd'],
database='selection',
charset='utf8mb4',
autocommit=True,
)
def _exec_doris_sql(sql_list, use_type='selection'):
"""通过 pymysql 走 Doris jdbc_port 执行 DDL / DML,返回最后一条语句的受影响行数
(pymysql 的 execute() 本身就会返回受影响行数,这里顺手往外传,避免调用方为了拿到
"写了几行"还要额外发一次 COUNT 查询)"""
conn = _doris_connect(use_type)
try:
cur = conn.cursor()
affected_rows = 0
for sql in sql_list:
print(f"[Doris SQL] {sql[:250]}{'...' if len(sql) > 250 else ''}")
affected_rows = cur.execute(sql)
cur.close()
return affected_rows
finally:
conn.close()
def get_latest_month(spark):
"""通过 MySQL workflow_everyday.MAX(report_date) 探测最新的流量选品月份(只取1个月,不再往前推3月)"""
sql_max_month = (
"select MAX(report_date) as date_info from workflow_everyday "
"where site_name = 'us' and date_type = 'month' and page = '流量选品'"
)
print(f"sql_max_month = {sql_max_month}")
mysql_con = DBUtil.get_connection_info('mysql', 'us')
max_date_info = SparkUtil.read_jdbc_query(
session=spark, url=mysql_con['url'],
pwd=mysql_con['pwd'], username=mysql_con['username'], query=sql_max_month,
).collect()[0]['date_info']
assert max_date_info is not None, "workflow_everyday 流量选品月度记录为空, 无法推断最新月份"
print(f"最新月份: {max_date_info}")
return str(max_date_info)
# ===== 任务1/2 共用:月表 + 年表的"9个每日变化字段"部分列更新 =====
def _daily_field_exprs():
"""9个每日变化字段各自的"新值"表达式(与 dwt_flow_asin_month.py 里的口径保持一致),
SELECT 取新值 和 WHERE 判断是否变化 两处都要用到,这里只定义一次"""
launch_time_expr = "COALESCE(f.launch_time, kp.keepa_launch_time)"
launch_time_type_expr = f"""CASE
WHEN {launch_time_expr} IS NULL THEN 0
WHEN DATEDIFF(f.asin_crawl_date, {launch_time_expr}) <= 30 THEN 1
WHEN DATEDIFF(f.asin_crawl_date, {launch_time_expr}) <= 90 THEN 2
WHEN DATEDIFF(f.asin_crawl_date, {launch_time_expr}) <= 180 THEN 3
WHEN DATEDIFF(f.asin_crawl_date, {launch_time_expr}) <= 360 THEN 4
WHEN DATEDIFF(f.asin_crawl_date, {launch_time_expr}) <= 720 THEN 5
WHEN DATEDIFF(f.asin_crawl_date, {launch_time_expr}) <= 1080 THEN 6
ELSE 7
END"""
tracking_since_expr = "FROM_UNIXTIME((CAST(kp.tracking_since AS BIGINT) + 21564000) * 60)"
tracking_since_type_expr = f"""CASE
WHEN kp.tracking_since IS NULL OR kp.tracking_since <= 0 THEN 0
WHEN DATEDIFF(f.asin_crawl_date, {tracking_since_expr}) <= 30 THEN 1
WHEN DATEDIFF(f.asin_crawl_date, {tracking_since_expr}) <= 90 THEN 2
WHEN DATEDIFF(f.asin_crawl_date, {tracking_since_expr}) <= 180 THEN 3
WHEN DATEDIFF(f.asin_crawl_date, {tracking_since_expr}) <= 360 THEN 4
WHEN DATEDIFF(f.asin_crawl_date, {tracking_since_expr}) <= 720 THEN 5
WHEN DATEDIFF(f.asin_crawl_date, {tracking_since_expr}) <= 1080 THEN 6
ELSE 7
END"""
item_weight_expr = "CASE WHEN kp.item_weight > 0 THEN kp.item_weight ELSE kp.package_weight END"
# (字段名, 新值表达式) 顺序即最终 INSERT 的列顺序
return [
("ocean_profit", "pr.ocean_profit"),
("air_profit", "pr.air_profit"),
("launch_time", launch_time_expr),
("launch_time_type", launch_time_type_expr),
("tracking_since", tracking_since_expr),
("tracking_since_type", tracking_since_type_expr),
("package_length", "kp.package_length"),
("package_width", "kp.package_width"),
("package_height", "kp.package_height"),
("item_weight", item_weight_expr),
]
def build_daily_fields_refresh_sql(site_name, table_name, extra_where_sql=None):
"""构建月表/年表共用的部分列更新 SQL:
只更新9个每日变化字段,且只对"新值跟表里已存的值不一致"的行才写入
(<=> 是 NULL-safe 比较,已在生产 Doris 上验证过 有值<=>NULL 会判定为不相等)
extra_where_sql: 年表需要额外限定 latest_date_info=最新月,月表不需要(整表就是那一个月)
"""
field_exprs = _daily_field_exprs()
select_list = ",\n ".join(f"{expr} AS {name}" for name, expr in field_exprs)
insert_cols = ", ".join(name for name, _ in field_exprs)
unchanged_conditions = " AND ".join(f"({expr}) <=> f.{name}" for name, expr in field_exprs)
where_parts = [extra_where_sql] if extra_where_sql else []
where_parts.append(f"NOT (\n {unchanged_conditions}\n )")
where_clause = " AND ".join(where_parts)
return f"""
INSERT INTO `selection`.`{table_name}` (asin, {insert_cols})
SELECT
f.asin,
{select_list}
FROM `selection`.`{table_name}` f
LEFT JOIN `dwd`.`dwd_asin_profit_rate_latest` pr
ON f.asin = pr.asin AND f.price = pr.price AND pr.site_name = '{site_name}'
LEFT JOIN `dwd`.`dwd_keepa_asin_detail` kp
ON f.asin = kp.asin AND kp.site_name = '{site_name}'
WHERE {where_clause}
"""
# ===== 任务1:刷新月表 =====
def refresh_month_table(site_name, latest_month):
"""任务1: 刷新 selection.{site}_flow_asin_month_{最新月} 的9个每日变化字段(表已由月度主流程建好)"""
table_name = f"{site_name}_flow_asin_month_{latest_month.replace('-', '_')}"
print(f"\n========== [任务1/月表] selection.{table_name} ==========")
sql = build_daily_fields_refresh_sql(site_name, table_name)
affected_rows = _exec_doris_sql([sql])
print(f"[完成] selection.{table_name},本次更新 {affected_rows} 行")
# ===== 任务2:刷新年表 =====
def refresh_year_table(site_name, latest_month):
"""任务2: 刷新 selection.{site}_flow_asin_365day,限定 latest_date_info=最新月 这批 asin 的9个每日变化字段"""
table_name = f"{site_name}_flow_asin_365day"
print(f"\n========== [任务2/年表] selection.{table_name} (latest_date_info={latest_month}) ==========")
extra_where_sql = f"f.latest_date_info = '{latest_month}'"
sql = build_daily_fields_refresh_sql(site_name, table_name, extra_where_sql=extra_where_sql)
affected_rows = _exec_doris_sql([sql])
print(f"[完成] selection.{table_name},本次更新 {affected_rows} 行")
# ===== 任务3:刷新利润率趋势中间表(只处理最新1个月)=====
def refresh_profit_trend_staging(site_name, latest_month):
"""任务3: 刷新 dwt.dwt_asin_profit_rate_history,只处理最新1个月(原来是近3月)
Doris MOW + sequence_col=update_time 本身已保证"最新 update_time 生效"的正确性,
这里额外加 pr.update_time <=> h.update_time 判断纯粹是为了减少不必要的 MOW 写入,
不判断具体利润率数值——利润率有没有变,跟 update_time 有没有推进是同一件事
返回本次实际写入的行数,供 main() 判断中间表有没有变化、要不要跳过任务4的全量重聚合
"""
print(f"\n========== [任务3/趋势中间表] dwt.dwt_asin_profit_rate_history (date_info={latest_month}) ==========")
sql = f"""
INSERT INTO `dwt`.`dwt_asin_profit_rate_history`
(site_name, asin, date_info, price, ocean_profit, air_profit, update_time)
SELECT
'{site_name}' AS site_name,
m.asin,
m.date_info,
m.price,
pr.ocean_profit,
pr.air_profit,
pr.update_time
FROM `dwt`.`{site_name}_flow_asin_month` m
INNER JOIN `dwd`.`dwd_asin_profit_rate_latest` pr
ON m.asin = pr.asin AND m.price = pr.price AND pr.site_name = '{site_name}'
LEFT JOIN `dwt`.`dwt_asin_profit_rate_history` h
ON h.site_name = '{site_name}' AND h.asin = m.asin AND h.date_info = m.date_info
WHERE m.date_info = '{latest_month}'
AND m.price > 0
AND NOT (pr.update_time <=> h.update_time)
"""
affected_rows = _exec_doris_sql([sql])
print(f"[完成] dwt.dwt_asin_profit_rate_history,本次写入 {affected_rows} 行")
return affected_rows
# ===== 任务4:Spark 聚合写入趋势物理表(全量重聚合,不受"只刷新最新月"影响)=====
def spark_aggregate_and_export_trend(spark, site_name):
"""任务4: Spark 读 dwt.dwt_asin_profit_rate_history 全量历史 → 聚合成每个asin的时间序列数组
→ INSERT OVERWRITE 写 selection.{site}_asin_profit_rate_trend
用 Spark 做 collect_list 聚合是为了避开 Doris 端 GROUP BY + COLLECT_LIST 的内存爆炸问题;
这一步跟"只刷新最新月"无关,本来就是把中间表全量重聚合一遍,物理表本身很小
"""
trend_table = f"{site_name}_asin_profit_rate_trend"
print(f"\n========== [任务4/趋势物理表] dwt.dwt_asin_profit_rate_history({site_name}) -> selection.{trend_table} ==========")
table_identifier = "dwt.dwt_asin_profit_rate_history"
read_fields = "site_name,asin,date_info,price,ocean_profit,air_profit"
df = DorisHelper.spark_import_with_connector(spark, table_identifier, read_fields) \
.filter(F.col('site_name') == site_name) \
.select('asin', 'date_info', 'price', 'ocean_profit', 'air_profit') \
.repartition(40, 'asin')
# groupBy 内按 date_info 排序,保证 collect_list 出来的几个数组是按时间顺序一一对应的
win = Window.partitionBy('asin').orderBy(F.col('date_info').asc())
df_sorted = df.withColumn('rn', F.row_number().over(win))
df_agg = df_sorted.sortWithinPartitions('asin', 'rn').groupBy('asin').agg(
F.collect_list('date_info').alias('date_info_arr'),
F.collect_list('price').alias('price_arr'),
F.collect_list('ocean_profit').alias('ocean_profit_arr'),
F.collect_list('air_profit').alias('air_profit_arr'),
).cache()
cnt = df_agg.count()
print(f"Spark 聚合后 asin 数: {cnt:,}")
df_agg.show(5, truncate=False)
# Doris connector 需要 ARRAY 字段转 JSON 字符串(StreamLoad ARRAY 接收 JSON 数组字符串格式)
df_save = df_agg.select(
F.col('asin'),
F.to_json(F.col('date_info_arr')).alias('date_info_arr'),
F.to_json(F.col('price_arr')).alias('price_arr'),
F.to_json(F.col('ocean_profit_arr')).alias('ocean_profit_arr'),
F.to_json(F.col('air_profit_arr')).alias('air_profit_arr'),
)
table_columns = "asin, date_info_arr, price_arr, ocean_profit_arr, air_profit_arr"
DorisHelper.spark_export_with_columns(
df_save=df_save,
db_name="selection",
table_name=trend_table,
table_columns=table_columns,
)
df_agg.unpersist()
print(f"[完成] selection.{trend_table} 已更新, 总 asin 数 {cnt:,}")
def main():
spark = SparkUtil.get_spark_session("DwsFlowAsinRefresh")
latest_month = get_latest_month(spark)
# 任务1: 月表 selection.{site}_flow_asin_month_{最新月}
refresh_month_table(SITE_NAME, latest_month)
# 任务2: 年表 selection.{site}_flow_asin_365day
refresh_year_table(SITE_NAME, latest_month)
# 任务3: 利润率趋势中间表 dwt.dwt_asin_profit_rate_history
affected_rows = refresh_profit_trend_staging(SITE_NAME, latest_month)
# 任务4: Spark 聚合写入趋势物理表 selection.{site}_asin_profit_rate_trend
if affected_rows > 0:
spark_aggregate_and_export_trend(spark, SITE_NAME)
else:
print(f"\n===== [任务4] dwt_asin_profit_rate_history 本次无变化,跳过趋势物理表重聚合 =====")
print("\nsuccess!")
if __name__ == "__main__":
main()
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