如何读取数据湖中增量加载的新分区文件?(PySpark场景)
解决方案:PySpark增量读取Landing层新增分区到Bronze Delta表
核心思路是追踪上次处理的最大分区日期,对比Landing区的所有分区,筛选出新增的分区进行处理,同时用可靠的元数据存储记录处理状态。
步骤1:创建元数据存储表
需要一个元数据表来记录每个数据源的最后处理分区日期,确保每次执行都能明确增量起点。用Delta表存储元数据可保证原子性和并发安全:
from pyspark.sql.types import StructType, StructField, StringType, DateType, TimestampType import datetime # 定义元数据结构 metadata_schema = StructType([ StructField("data_source", StringType(), nullable=False), StructField("last_processed_date", DateType(), nullable=False), StructField("updated_at", TimestampType(), nullable=False) ]) # 初始化元数据表(仅第一次执行) initial_metadata = spark.createDataFrame([ ("your_target_source", datetime.date(1900, 1, 1), datetime.datetime.now()) ], schema=metadata_schema) initial_metadata.write.format("delta") \ .mode("overwrite") \ .save("/path/to/bronze/metadata")
步骤2:获取Landing区的所有分区日期
通过Hadoop FileSystem API列出Landing区的分区文件夹,提取DATE=后的日期值:
from pyspark.sql.functions import col, to_date from pyspark.sql.types import StringType landing_base_path = "/path/to/landing/your_source" fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) path_statuses = fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(landing_base_path)) # 提取所有有效分区日期 landing_partition_dates = [] for status in path_statuses: full_path = status.getPath().toString() if "DATE=" in full_path: date_str = full_path.split("DATE=")[-1].strip() landing_partition_dates.append(date_str) # 转为Spark可处理的日期DataFrame landing_dates_df = spark.createDataFrame(landing_partition_dates, StringType()).toDF("date_str") \ .withColumn("partition_date", to_date(col("date_str"), "yyyy-MM-dd")) \ .filter(col("partition_date").isNotNull())
步骤3:筛选新增分区
读取元数据表,获取上次处理的最大日期,筛选出Landing区中未处理的分区:
from delta.tables import DeltaTable # 读取元数据 metadata_table = DeltaTable.forPath(spark, "/path/to/bronze/metadata") last_processed_df = metadata_table.toDF() \ .filter(col("data_source") == "your_target_source") \ .select("last_processed_date") # 获取上次处理日期(首次执行默认1900-01-01) last_processed_date = last_processed_df.first()["last_processed_date"] # 筛选新增分区 new_partitions_df = landing_dates_df.filter(col("partition_date") > last_processed_date) new_partition_dates = new_partitions_df.select("date_str").rdd.flatMap(lambda x: x).collect()
步骤4:读取新增分区并写入Bronze层
如果存在新增分区,读取对应路径的数据并写入Delta表,最后更新元数据:
if new_partition_dates: # 构造新增分区的完整路径 partition_paths = [f"{landing_base_path}/DATE={date}" for date in new_partition_dates] # 读取Landing层数据(替换为你的数据源格式,比如csv、parquet等) raw_data = spark.read.format("parquet") \ .load(partition_paths) # 写入Bronze层Delta表(按DATE分区,append模式) raw_data.write.format("delta") \ .mode("append") \ .partitionBy("DATE") \ .save("/path/to/bronze/your_source") # 更新元数据:记录本次处理的最大日期 max_processed_date_str = max(new_partition_dates) update_data = spark.createDataFrame([ ("your_target_source", max_processed_date_str, datetime.datetime.now()) ], schema=metadata_schema) # 使用Merge操作保证并发安全 metadata_table.alias("target").merge( update_data.alias("source"), "target.data_source = source.data_source" ).whenMatchedUpdate(set={ "last_processed_date": "source.last_processed_date", "updated_at": "source.updated_at" }).whenNotMatchedInsert(values={ "data_source": "source.data_source", "last_processed_date": "source.last_processed_date", "updated_at": "source.updated_at" }).execute() else: print("No new partitions to process.")
关键注意事项
- 分区日期格式校验:确保Landing区的
DATE=格式是yyyy-MM-dd,否则需要调整to_date的格式参数 - 并发安全:使用Delta表的Merge操作更新元数据,避免多任务同时执行时的状态覆盖问题
- 容错处理:可添加异常捕获逻辑,确保元数据仅在数据写入成功后更新
内容的提问来源于stack exchange,提问作者guigahbr
相关产品推荐
相关产品推荐

