如何使用Databricks Auto Loader将不同类型CSV写入独立Parquet表
解决方案
你当前的代码会把sourcePath路径下所有CSV文件读取到同一个流中,最终写入同一个Delta表,所以会得到合并后的大表。可以通过以下两种方案实现按CSV类型/Schema拆分生成独立表:
方案一:按子目录拆分(最推荐,性能最优)
如果不同Schema的CSV可以存放在sourcePath下的独立子目录(比如user/放用户类CSV、order/放订单类CSV),直接为每个子目录启动独立的流任务即可,每个流单独维护自身的Schema和checkpoint,不会互相干扰:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("split_csv_loader").getOrCreate() # 按实际CSV分类调整配置 csv_configs = [ {"table_name": "user", "sub_dir": "user/", "checkpoint": f"{pathCheckpoint}/user/", "save_path": f"{pathResult}/user/"}, {"table_name": "order", "sub_dir": "order/", "checkpoint": f"{pathCheckpoint}/order/", "save_path": f"{pathResult}/order/"}, {"table_name": "product", "sub_dir": "product/", "checkpoint": f"{pathCheckpoint}/product/", "save_path": f"{pathResult}/product/"}, ] # 为每类CSV启动独立流任务 for conf in csv_configs: spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("delimiter", "~|~") \ .option("cloudFiles.inferColumnTypes","true") \ .option("cloudFiles.schemaLocation", conf["checkpoint"]) \ .load(f"{sourcePath}/{conf['sub_dir']}") \ .writeStream \ .format("delta") \ .option("mergeSchema", "true") \ .option("checkpointLocation", conf["checkpoint"]) \ .start(conf["save_path"])
方案二:同目录混合存放动态路由
如果所有CSV都混存在同一个sourcePath下无法拆分目录,可以用foreachBatch处理每一批次数据,根据文件名、文件路径或者数据Schema自动路由写入对应的目标表:
from pyspark.sql import functions as F def batch_processor(batch_df, batch_id): # 提取源文件路径,从文件名中提取归属表名(可根据自己的文件命名规则调整逻辑) batch_df = batch_df.withColumn("file_path", F.col("_metadata.file_path")) batch_df = batch_df.withColumn("table_name", F.split(F.element_at(F.split(F.col("file_path"), "/"), -1), "_")[0]) # 按表名分批次写入 distinct_tables = batch_df.select("table_name").distinct().collect() for row in distinct_tables: table_name = row["table_name"] table_data = batch_df.filter(F.col("table_name") == table_name).drop("file_path", "table_name") table_data.write.format("delta") \ .mode("append") \ .option("mergeSchema", "true") \ .save(f"{pathResult}/{table_name}/") # 启动流任务 spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("delimiter", "~|~") \ .option("cloudFiles.inferColumnTypes","true") \ .option("cloudFiles.schemaLocation", pathCheckpoint) \ .load(sourcePath) \ .writeStream \ .foreachBatch(batch_processor) \ .option("checkpointLocation", pathCheckpoint) \ .start()
注意事项
- 如果不同CSV的Schema差异极大,不推荐使用方案二,全局统一的Schema推断可能会出现字段类型兼容报错,优先使用方案一独立维护每个表的Schema。
- 方案二中如果需要按Schema而非文件名判断归属,可以在
batch_processor中校验batch_df.columns的字段列表,匹配到对应目标表后再写入。
内容的提问来源于stack exchange,提问作者Leonardo Lima
相关产品推荐
相关产品推荐

