如何使用Scala Spark解压LZ4压缩格式的CSV文件?
Spark解压LZ4压缩CSV文件的实现方案
Spark可以处理LZ4压缩格式的CSV文件,但原生读取方式存在格式兼容性限制,以下是针对相关问题的分析与可行解决方法:
失败方法分析
原生CSV读取指定compression参数
使用sparkSession.read.csv并指定compression="lz4"时,Spark依赖Hadoop的LZ4Codec,仅支持Hadoop定制的LZ4压缩格式(带有特定头部的块压缩)。如果文件是标准LZ4块或帧格式,该方法无法正确识别,会返回空结果。
代码示例:sparkSession.read .option("delimiter", ",") .option("compression", "lz4") .csv("data.csv.lz4")Text格式读取后解压
用text格式读取LZ4文件时,Spark会将二进制数据按行分割,但LZ4是整体压缩格式,拆分后的字节片段无法被正确解压,因此返回空结果。
代码示例:val rez = sparkSession.read .format("text") .load("/data.csv.lz4") rez .map{ row => val bytes = row.getAs[Array[Byte]]("value") unzipLZ4Content(bytes) } .toDF("value") .show(truncate = false) def unzipLZ4Content(bytes: Array[Byte]) = Using .Manager { use => val bais = use(new ByteArrayInputStream(bytes)) val lz4is = use(new LZ4BlockInputStream(bais)) Source .fromInputStream(lz4is) .getLines() .toSeq } .getOrElse(Seq(""))
可行解决方案:二进制读取+兼容多格式解压
通过Spark的binaryFiles读取完整的压缩文件二进制数据,再使用支持标准LZ4块/帧格式的解压工具处理,即可正确解析内容。
完整代码实现
import java.io.ByteArrayInputStream import org.apache.spark.SparkContext import org.apache.spark.sql.SparkSession import scala.util.Using import scala.io.Source import net.jpountz.lz4.{FramedLZ4CompressorInputStream, BlockLZ4CompressorInputStream} object LZ4CsvReader { def main(args: Array[String]): Unit = { val sparkSession = SparkSession.builder() .appName("LZ4CsvReader") .master("local[*]") // 本地模式,生产环境可移除 .getOrCreate() val rez = sparkSession.sparkContext .binaryFiles("/data.csv.lz4") rez .map{ case (_, pds: SparkContext.PortableDataStream) => val bytes = pds.toArray() unzipBinaryContent(bytes) } .toDF("value") .show(truncate = false) sparkSession.stop() } def unzipBinaryContent(bytes: Array[Byte]) = Using .Manager { use => val bais = use(new ByteArrayInputStream(bytes)) val lz4is = try { use(new FramedLZ4CompressorInputStream(bais)) } catch { case _: Exception => bais.reset() use(new BlockLZ4CompressorInputStream(bais)) } Source.fromInputStream(lz4is).getLines().toSeq } .getOrElse(Seq("")) }
依赖说明
需要在项目中添加LZ4 Java库依赖(以Maven为例):
<dependency> <groupId>net.jpountz.lz4</groupId> <artifactId>lz4</artifactId> <version>1.3.0</version> </dependency>
该方案通过尝试两种LZ4输入流(帧级和块级),兼容不同标准的LZ4压缩格式,确保能正确解压并读取CSV内容。
内容的提问来源于stack exchange,提问作者Jelly
相关产品推荐
相关产品推荐

