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

PySpark批量小文件转Delta表并增量追加新文件的最佳实现方法

基于Delta Lake三层架构的CSV增量入湖方案

所有按天生成的原始CSV统一写入Bronze层做原始数据留存,全程保留溯源字段保证可回溯,后续清洗去重、维度关联加工后流入Silver层,聚合指标、业务宽表在Gold层产出,全程无需回读源CSV文件。

一、PySpark追加小CSV到Delta表的最优实现路径

核心思路是开启Delta原生小文件优化、固定Schema、加审计字段,避免手动处理小文件、Schema漂移等问题:

  • 初始化SparkSession时提前开启Delta自动优化配置,无需事后单独跑小文件合并作业:
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("CSVToDeltaIncremental") \
    .config("spark.databricks.delta.autoCompact.enabled", "true") \
    .config("spark.databricks.delta.autoOptimize.optimizeWrite", "true") \
    .config("spark.sql.files.ignoreCorruptFiles", "true") \
    .getOrCreate()

配置开启后,写入时Delta会自动将零散的小数据合并为128MB左右的最优文件大小,适配对象存储、HDFS的读性能要求,不要手动调用coalesce(1)这类强行合并为单文件的操作,容易导致数据倾斜、写入超时。

  • 读取CSV时提前固定Schema,禁止Spark自动推断Schema,既提升读取性能,也避免类型漂移导致的写入错误:
from pyspark.sql.types import StructType, StructField, StringType, TimestampType, LongType
# 按实际CSV字段定义即可
csv_schema = StructType([
    StructField("user_id", LongType(), False),
    StructField("biz_time", TimestampType(), False),
    StructField("biz_content", StringType(), True)
])

DELTA_BRONZE_PATH = "/lakehouse/bronze/raw_biz_daily"
CSV_ROOT_PATH = "/source/csv/daily/"
  • 写入时追加溯源字段,锁Schema保证数据质量,用追加模式写入:
from pyspark.sql.functions import input_file_name, current_timestamp

# new_file_list为后续筛选出的新增CSV文件路径列表
new_csv_df = spark.read.schema(csv_schema) \
    .option("header", "true") # CSV带表头则开启,无表头则关闭
    .csv(new_file_list) \
    .withColumn("source_file", input_file_name()) \
    .withColumn("ingest_ts", current_timestamp())

new_csv_df.write \
    .format("delta") \
    .mode("append") \
    .option("mergeSchema", "false") # 字段不匹配直接报错,阻断脏数据写入
    .save(DELTA_BRONZE_PATH)

二、增量筛选未处理新增CSV文件的实现方案

不要靠文件修改时间筛选新文件,容易漏处理补传的历史文件,最稳定的幂等方案是和Delta表已记录的处理文件做比对:

  • 第一步:从已写入的Bronze层Delta表中,取出所有已经处理过的源文件路径。因为写入时已经保留了source_file字段,直接查询即可,文件路径为字符串,数据量极小可以直接拉取到内存:
from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, DELTA_BRONZE_PATH)
processed_files = set(
    row.source_file for row in delta_table.toDF()
    .select("source_file")
    .distinct()
    .collect()
)
  • 第二步:遍历CSV根目录下所有CSV文件,过滤掉已处理的文件,得到本次待写入的新增文件列表:
# 调用Hadoop API遍历文件,兼容本地盘、HDFS、S3、OSS等所有存储介质
hadoop_conf = spark._jsc.hadoopConfiguration()
hadoop_fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf)
csv_root_jpath = spark._jvm.org.apache.hadoop.fs.Path(CSV_ROOT_PATH)

all_csv_files = [
    str(file_status.getPath()) 
    for file_status in hadoop_fs.listFiles(csv_root_jpath, True)
    if file_status.getPath().getName().endswith(".csv")
]

new_file_list = [f for f in all_csv_files if f not in processed_files]
  • 无新增文件时直接跳过本次作业,存在新增文件时走前面的读取、写入逻辑即可。如果CSV按天存放在分区子目录下,可以先按目录修改时间过滤最近N天的目录再做文件比对,进一步减少遍历开销。

生产环境不建议用Structured Streaming直读CSV目录做自动增量,一旦目录下出现写入未完成的临时文件、空文件、损坏文件会直接导致作业失败。上述批处理方案按天/按小时调度即可,稳定性更高,作业重跑不会重复写入,历史补传的文件也会被自动识别处理,完全适配天级小文件的入湖场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 16:24:34