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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 12:40:28