Spark高效实现追加数据并去重后保存为Parquet格式的方案
PySpark时间维度表增量生成与去重优化方案
核心实现思路
- 首次执行:生成连续10年的全量时间维度数据,直接写入目标存储
- 后续执行:根据指定的增量日期范围生成新数据,与已有数据合并后按主键(如
date_key)去重,最终覆盖写入目标存储,省去先追加再读取去重的冗余操作
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import date_format, year, month, dayofmonth, dayofweek, weekofyear, quarter import datetime def create_spark_session(): return SparkSession.builder \ .appName("TimeDimensionGenerator") \ .getOrCreate() def generate_time_dim(start_date, end_date): # 生成连续日期序列 date_range = spark.sql(f"SELECT sequence(to_date('{start_date}'), to_date('{end_date}'), interval 1 day) AS dates") \ .selectExpr("explode(dates) AS date") # 构建时间维度全量字段 time_dim = date_range \ .withColumn("date_key", date_format("date", "yyyyMMdd").cast("int")) \ .withColumn("year", year("date")) \ .withColumn("month", month("date")) \ .withColumn("day", dayofmonth("date")) \ .withColumn("day_of_week", dayofweek("date")) \ .withColumn("week_of_year", weekofyear("date")) \ .withColumn("quarter", quarter("date")) \ .withColumn("is_weekend", (dayofweek("date").isin(1,7)).cast("boolean")) return time_dim if __name__ == "__main__": spark = create_spark_session() target_table_path = "/path/to/time_dimension_parquet" target_table_name = "time_dimension" # 首次运行配置:生成当前日期往前推10年到当日的全量数据 # 后续运行可直接修改start_date/end_date为需要追加的增量日期范围 today = datetime.date.today() first_run_start = today - datetime.timedelta(days=365*10) start_date = first_run_start.strftime("%Y-%m-%d") end_date = today.strftime("%Y-%m-%d") # 生成目标日期范围的数据 incremental_df = generate_time_dim(start_date, end_date) # 处理已有数据与增量数据的合并去重 if spark.catalog.tableExists(target_table_name): existing_df = spark.read.table(target_table_name) # 合并后按date_key去重,保留最新生成的记录(或任意唯一记录) final_df = incremental_df.unionByName(existing_df).dropDuplicates(["date_key"]) else: # 首次运行直接使用全量生成的数据 final_df = incremental_df # 写入去重后的完整数据集 final_df.write \ .mode("overwrite") \ .format("parquet") \ .saveAsTable(target_table_name) spark.stop()
关键逻辑说明
- 主键去重:以
date_key作为唯一标识,通过dropDuplicates(["date_key"])确保每条日期记录唯一,避免重复数据 - 增量判断:利用
spark.catalog.tableExists检查表是否存在,自动区分首次全量生成和后续增量追加场景 - 灵活扩展:后续运行时只需修改
start_date和end_date为需要补充的日期区间,即可自动生成增量并完成去重合并 - IO优化:通过一次合并去重后直接覆盖写入,替代原有的"追加→读取→去重→覆盖"流程,减少IO操作次数
内容的提问来源于stack exchange,提问作者Abinandhana G
相关产品推荐
相关产品推荐

