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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 09:39:02