如何使用Apache Spark读取zip压缩包内文件生成Spark DataFrame
Spark读取ZIP包内6G文本文件并入库实现方式
首先明确两个前提:
- 原有读取普通分隔文本、写入SQL库的逻辑不需要大改,核心只需要调整ZIP包的读取部分
- ZIP格式本身是不可切分的,直接单线程读取6G文件效率极低还容易内存溢出,下面给两种可直接落地的方案
方案1:最小改动快速实现(适合测试、小集群临时跑数场景)
Spark内置支持ZIP格式的压缩解码,你只需要在原有read.csv的参数里加一行压缩格式指定即可,不需要改其他逻辑:
# 路径直接填zip包路径即可,其他参数和你之前读普通txt完全一致 df = spark.read.csv( path="MyFilePath/MyCompressedFile.zip", sep="|", header="true", inferSchema="true", compression="zip" # 仅新增这一行指定压缩格式 ) # 校验数据和之前完全一样 df.show() # 写入SQL库的逻辑无任何改动,以JDBC写入为例 df.write.format("jdbc") \ .option("url", "你的数据库JDBC连接串") \ .option("dbtable", "目标表名") \ .option("user", "数据库账号") \ .option("password", "数据库密码") \ .mode("overwrite") # 按业务需求选append/overwrite等模式 \ .save()
注意:这个方案的缺点是ZIP不可切分,Spark只会启动1个Task读取整个6G文件,并行度为1,读取速度慢,单Task内存压力大,数据量大时很容易报OOM,只适合临时跑数用
方案2:分布式并行读取(生产环境推荐,性能稳定)
6G单文件不算小,生产环境建议用分布式解压读取的方式,把解压和解析逻辑分发到多个Executor并行执行,避免单点性能瓶颈,代码示例如下:
import zipfile import io # 以二进制文件方式读取ZIP,设置合理的最小分区数,单分区对应处理500M-1G数据即可,6G文件建议设8-16个分区 zip_rdd = spark.sparkContext.binaryFiles("MyFilePath/MyCompressedFile.zip", minPartitions=12) # 定义逐行解析逻辑,避免一次性加载全量文件到内存 def parse_zip_line(zip_item): _, zip_content = zip_item with zipfile.ZipFile(io.BytesIO(zip_content)) as zf: # 已知包内只有1个文本文件,直接取第一个文件流 with zf.open(zf.namelist()[0]) as f: for line in f: # 按指定分隔符拆分,编码按实际文件编码调整 yield line.decode("utf-8").rstrip("\n").split("|") # 执行解析得到结构化RDD parsed_rdd = zip_rdd.flatMap(parse_zip_line) # 取第一行作为表头,和原有header=true逻辑对齐 header_row = parsed_rdd.first() # 过滤掉表头行后转成DataFrame,和原有inferSchema逻辑一致会自动推断字段类型 df = parsed_rdd.filter(lambda x: x != header_row).toDF(header_row) # 校验数据 df.show() # 写入SQL库的逻辑和原有流程完全一致 df.write.format("jdbc") \ .option("url", "你的数据库JDBC连接串") \ .option("dbtable", "目标表名") \ .option("user", "数据库账号") \ .option("password", "数据库密码") \ .mode("overwrite") \ .save()
实操提示:如果提前定义好明确的表Schema(不用inferSchema自动推断),读取性能还能再提升30%以上,也能避免自动类型推断带来的字段类型错误
避坑提醒
不要用wholeTextFiles接口读取这个ZIP包,这个接口会把整个文件全量加载到单个Executor内存,6G文件100%会触发OOM。
内容的提问来源于stack exchange,提问作者nam
相关产品推荐
相关产品推荐

