You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

关键逻辑说明

  1. 主键去重:以date_key作为唯一标识,通过dropDuplicates(["date_key"])确保每条日期记录唯一,避免重复数据
  2. 增量判断:利用spark.catalog.tableExists检查表是否存在,自动区分首次全量生成和后续增量追加场景
  3. 灵活扩展:后续运行时只需修改start_date和end_date为需要补充的日期区间,即可自动生成增量并完成去重合并
  4. IO优化:通过一次合并去重后直接覆盖写入,替代原有的"追加→读取→去重→覆盖"流程,减少IO操作次数

内容的提问来源于stack exchange,提问作者Abinandhana G

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 15:17:22