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
相关产品推荐
相关产品推荐

