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

如何使用Scala Spark解压LZ4压缩格式的CSV文件?

Spark解压LZ4压缩CSV文件的实现方案

Spark可以处理LZ4压缩格式的CSV文件,但原生读取方式存在格式兼容性限制,以下是针对相关问题的分析与可行解决方法:

失败方法分析

  1. 原生CSV读取指定compression参数
    使用sparkSession.read.csv并指定compression="lz4"时,Spark依赖Hadoop的LZ4Codec,仅支持Hadoop定制的LZ4压缩格式(带有特定头部的块压缩)。如果文件是标准LZ4块或帧格式,该方法无法正确识别,会返回空结果。
    代码示例:

    sparkSession.read
      .option("delimiter", ",")
      .option("compression", "lz4")
      .csv("data.csv.lz4")
    
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 13:43:18