如何通过Spark原生读取LZMA压缩的xz格式二进制文件
解决Spark通过hadoop-xz原生读取XZ压缩二进制文件的问题
问题根源
Spark的binaryFile数据源设计为读取文件原始字节流,不会自动应用配置的压缩解码器;而text数据源基于TextInputFormat,会自动识别文件扩展名并调用对应的Codec解压,因此text文件能正常处理,但二进制文件用binaryFile读时无法触发hadoop-xz的解压逻辑。
hadoop-xz的XZCodec本身支持二进制流解压,只是需要在读取二进制文件时主动触发它的使用。
解决方案:自定义InputFormat实现原生读取
通过自定义支持XZ解压的FileInputFormat,结合Spark的hadoopFile接口读取,即可实现基于hadoop-xz的原生二进制文件读取,最终转换为DataFrame使用。
1. 实现自定义XZBinaryInputFormat
import org.apache.hadoop.fs.Path import org.apache.hadoop.io.{BytesWritable, NullWritable} import org.apache.hadoop.mapreduce.{InputSplit, JobContext, RecordReader, TaskAttemptContext} import org.apache.hadoop.mapreduce.lib.input.FileInputFormat import org.apache.hadoop.mapreduce.lib.input.FileSplit import io.sensesecure.hadoop.xz.XZCodec import org.apache.commons.io.IOUtils import java.io.ByteArrayOutputStream class XZBinaryInputFormat extends FileInputFormat[NullWritable, BytesWritable] { // XZ压缩文件默认不可拆分,这里强制设置为不可拆分 override def isSplitable(context: JobContext, filename: Path): Boolean = false override def createRecordReader(split: InputSplit, context: TaskAttemptContext): RecordReader[NullWritable, BytesWritable] = { new RecordReader[NullWritable, BytesWritable] { private var inputStream: java.io.InputStream = _ private var resultBytes: BytesWritable = _ private var hasProcessed = false override def initialize(split: InputSplit, context: TaskAttemptContext): Unit = { val fileSplit = split.asInstanceOf[FileSplit] val fs = fileSplit.getPath.getFileSystem(context.getConfiguration) val rawStream = fs.open(fileSplit.getPath) // 使用hadoop-xz的XZCodec创建解压流 val xzCodec = new XZCodec() xzCodec.setConf(context.getConfiguration) inputStream = xzCodec.createInputStream(rawStream) } override def nextKeyValue(): Boolean = { if (!hasProcessed) { val outputStream = new ByteArrayOutputStream() IOUtils.copy(inputStream, outputStream) resultBytes = new BytesWritable(outputStream.toByteArray) hasProcessed = true true } else { false } } override def getCurrentKey: NullWritable = NullWritable.get() override def getCurrentValue: BytesWritable = resultBytes override def getProgress: Float = if (hasProcessed) 1.0f else 0.0f override def close(): Unit = inputStream.close() } } }
2. 在Spark中使用自定义InputFormat读取并转换为DataFrame
import org.apache.hadoop.io.{BytesWritable, NullWritable} import org.apache.spark.sql.SparkSession val spark = SparkSession.builder() .master("local[*]") .appName("XZBinaryReader") .getOrCreate() // 读取XZ压缩的二进制文件 val xzBinaryRDD = spark.sparkContext.hadoopFile[NullWritable, BytesWritable, XZBinaryInputFormat]( "path/to/binaryRawFile.xz" ) // 转换为DataFrame,content列存储解压后的二进制字节数组 val binaryDataDF = xzBinaryRDD.map(_._2.copyBytes()).toDF("content") // 后续处理示例 binaryDataDF.select("content").show(false)
说明
- 自定义
XZBinaryInputFormat通过hadoop-xz的XZCodec对输入流进行解压,确保二进制内容被正确还原。 - 该方式属于Spark原生的Hadoop接口调用,完全基于hadoop-xz实现解压,无需手动处理流操作。
内容的提问来源于stack exchange,提问作者Mehdi
相关产品推荐
相关产品推荐

