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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:00:00