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

如何读取数据湖中增量加载的新分区文件?(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 05:33:16