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

Spark如何将被Bz2整体压缩的Parquet文件解压为原生Parquet文件

问题原因分析
  • 第一种textFile操作报错原因:textFile会按行解析文件内容为字符串格式,而你的文件是二进制的bzip2压缩Parquet,转成字符串后会损坏Parquet的二进制结构,且String类型不符合Parquet输出格式要求的输入类型,所以抛出类型不匹配错误。
  • 第二种直接读Parquet报错原因:Parquet文件的元数据存储在文件末尾的Footer区域,Spark读取Parquet时需要先随机定位到文件末尾读取Footer,但外层的bzip2是不可切分的流式压缩格式,Spark无法直接在压缩包内做随机寻址,因此无法读取到Parquet的Footer信息,抛出IO异常。
解决方法

方法1:Spark批量处理(适合多文件场景)

核心逻辑是先读取压缩包的二进制内容,手动解压后直接写入为原生二进制Parquet文件,示例代码如下:

import org.apache.commons.compress.compressors.bzip2.BZip2CompressorInputStream
import java.io.ByteArrayOutputStream
import org.apache.hadoop.io.{BytesWritable, NullWritable}

// 读取bz2文件的二进制内容,支持批量匹配多个bz2文件
val bz2FileRDD = spark.sparkContext.binaryFiles("/mnt/shahgau/test/*.bz2")

val rawParquetRDD = bz2FileRDD.flatMap { case (filePath, streamData) =>
  val inputStream = streamData.open()
  val bzip2In = new BZip2CompressorInputStream(inputStream)
  val outputStream = new ByteArrayOutputStream()
  val buffer = new Array[Byte](8192)
  var readLen = 0
  
  try {
    while ({readLen = bzip2In.read(buffer); readLen > 0}) {
      outputStream.write(buffer, 0, readLen)
    }
    Some(outputStream.toByteArray)
  } finally {
    bzip2In.close()
    outputStream.close()
    inputStream.close()
  }
}

// 写出为原始Parquet文件
rawParquetRDD
  .map(bytes => (NullWritable.get(), new BytesWritable(bytes)))
  .saveAsHadoopFile(
    "/mnt/shahgau/test/raw_parquet_output",
    classOf[NullWritable],
    classOf[BytesWritable],
    classOf[org.apache.hadoop.mapreduce.lib.output.FileOutputFormat]
  )

写出完成后,到输出目录下取出生成的文件,重命名为.parquet后缀即可正常读取。

方法2:命令行直接解压(适合单文件/小文件场景)

直接用Hadoop命令结合bzip2工具解压,效率更高:

hadoop fs -cat /mnt/shahgau/test/000000_0.bz2 | bzip2 -d | hadoop fs -put - /mnt/shahgau/test/000000_0.raw.parquet

解压完成后即可直接用Spark读取:

spark.read.parquet("/mnt/shahgau/test/000000_0.raw.parquet").show()

内容的提问来源于stack exchange,提问作者Gaurang Shah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 13:30:01