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
相关产品推荐
相关产品推荐

