如何在Databricks中批量读取数据湖内所有.csv格式文件
在Databricks中批量读取数据湖内同Schema的CSV文件
要一次性读取数据湖内所有后缀为.csv的文件(无论文件名和所在子文件夹),可以利用Spark的递归路径匹配功能,结合统一的Schema配置来实现,具体操作如下:
1. 配置文件路径匹配规则
由于CSV文件分布在不同层级的子文件夹中,需要使用递归通配符**匹配所有子目录,再用*.csv筛选目标文件。路径格式根据你的数据湖类型调整(以ADLS Gen2为例):
file_path = "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/**/*.csv"
若是S3、DBFS等其他存储类型,替换对应的路径前缀即可。
2. 批量读取CSV文件
方式一:指定Schema(推荐,性能更优)
因为所有文件Schema一致,提前定义Schema可以避免Spark多次推断Schema,大幅提升读取效率。假设你的数据结构包含id(整数)、name(字符串)、value(浮点数)三个字段:
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, FloatType # 定义统一Schema custom_schema = StructType([ StructField("id", IntegerType(), nullable=True), StructField("name", StringType(), nullable=True), StructField("value", FloatType(), nullable=True) ]) # 读取所有符合条件的CSV文件 df = spark.read \ .option("header", "true") # CSV含表头则设为true,无表头设为false .schema(custom_schema) \ .csv(file_path)
方式二:自动推断Schema
如果不确定Schema结构,也可以让Spark自动推断,适合小数据集场景:
df = spark.read \ .option("header", "true") \ .option("inferSchema", "true") \ .csv(file_path)
3. 可选优化配置
- 处理坏文件:若部分CSV格式异常,可添加参数跳过错误行:
df = spark.read \ .option("header", "true") \ .schema(custom_schema) \ .option("mode", "DROPMALFORMED") # 丢弃格式错误的行 .csv(file_path) - 自定义分隔符:如果CSV使用非逗号分隔符,添加
option("delimiter", ";")指定分隔符。
验证读取结果
读取完成后,可通过以下命令确认数据:
df.show() # 显示前20行数据 df.printSchema() # 打印Schema确认结构 df.count() # 统计总记录数
内容的提问来源于stack exchange,提问作者Ankit Sawa
相关产品推荐
相关产品推荐

