在Databricks中直接读取ADLS上ZIP内的大CSV文件为Spark DataFrame
解决Databricks中读取ADLS挂载ZIP内大UTF-16-LE TXT文件的方案
搞定这个问题其实不难,你之前用Pandas踩的坑都是因为单进程处理大文件的局限性——不管是解压不全、字段数不匹配报错,还是连接拒绝,本质上都是Pandas没办法高效处理分布式存储上的超大文件。换成Spark的分布式方案就完美解决了,下面给你PySpark和Scala两种实现方式,都能直接读取DBFS上ZIP里的大TXT文件,保证全量数据、正常展示,还不会搞挂集群。
PySpark 实现方案
Spark可以直接识别ZIP压缩包并读取内部的文本文件,不需要手动解压(避免占用DBFS空间和集群资源),同时分布式处理特性不会把整个文件加载到Driver内存。
基础读取代码(处理字段数不一致+UTF-16编码)
# Databricks环境下无需手动初始化SparkSession,直接用内置的spark变量即可 # 如果是本地环境才需要手动创建,这里省略 # 直接读取ZIP内的TXT文件 df = spark.read \ .option("encoding", "UTF-16LE") # 注意Spark里编码名是UTF-16LE,不是UTF-16-LE .option("header", "false") # 如果你的TXT有表头,改成"true" .option("mode", "PERMISSIVE") # 允许字段数不匹配的行,错误行会被放到_corrupt_record列 .option("delimiter", ",") \ .csv("dbfs:/path/to/your/1.3GB-file.zip") # DBFS路径用dbfs://开头或者/dbfs/...都可以 # 给列命名(如果没有表头的话) df = df.toDF("字段1", "字段2", "额外字段") # 预留额外字段位置,对应那些有3个字段的行 # 查看前10行数据,用show()而不是head(),避免把大量数据拉到Driver df.show(10)
进阶:手动定义Schema(提高性能+精准控制)
大文件不要用自动推断Schema(inferSchema=true),会很慢,手动定义Schema更高效:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 根据你的实际数据类型定义Schema,比如字段1是字符串,字段2是整数 custom_schema = StructType([ StructField("字段1", StringType(), nullable=True), StructField("字段2", IntegerType(), nullable=True) ]) df = spark.read \ .option("encoding", "UTF-16LE") \ .option("header", "false") \ .option("mode", "PERMISSIVE") \ .option("delimiter", ",") \ .schema(custom_schema) \ .csv("dbfs:/path/to/your/1.3GB-file.zip") # 过滤掉坏行(如果需要) filtered_df = df.filter(df["_corrupt_record"].isNull()) filtered_df.show(10)
Scala 实现方案
如果你的团队更倾向于Scala,这个方案逻辑和PySpark一致,只是语法不同:
// Databricks环境下直接用内置的spark变量 val df = spark.read .option("encoding", "UTF-16LE") .option("header", "false") .option("mode", "PERMISSIVE") .option("delimiter", ",") .csv("dbfs:/path/to/your/1.3GB-file.zip") // 重命名列 val namedDF = df.toDF("字段1", "字段2", "额外字段") // 查看数据 namedDF.show(10) // 过滤坏行(可选) val filteredDF = namedDF.filter($"_corrupt_record".isNull) filteredDF.show(10)
关键注意事项
- 避免手动解压ZIP:Spark会自动处理压缩包,手动解压6GB文件不仅占DBFS空间,还可能导致集群节点内存不足。
- 编码匹配:一定要用
UTF-16LE,和你之前修复Pandas错误时的编码对应,Spark不识别UTF-16-LE这种带连字符的写法。 - 用show()替代head():
head()会把数据拉到Driver节点,大文件下容易OOM;show()只展示指定行数,分布式安全。 - 字段数处理:
PERMISSIVE模式是最灵活的,既不会直接报错中断,还能把错误行单独存起来方便后续排查;如果不需要坏行,可以换成DROPMALFORMED模式直接丢弃。
内容的提问来源于stack exchange,提问作者Yogesh Kulkarni
相关产品推荐
相关产品推荐

