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

如何从xz压缩二进制文件中解析并使用定长记录?

问题描述

我需要读取包含定长记录的xz压缩二进制文件并解析内容,使用了hadoop-xz的XZCodec编解码器,但运行代码时抛出异常,提示FixedLengthRecordReader不支持读取压缩文件。

尝试的代码

from pyspark.sql import SparkSession
from struct import unpack
spark = SparkSession.builder.master('local[*]').getOrCreate()
sc=spark.sparkContext
sc._jsc.hadoopConfiguration().set("io.compression.codecs","io.sensesecure.hadoop.xz.XZCodec")
trace_format='<Q2B2B4B2Q4Q'
file=sc.binaryRecords('test.xz',64)
file.flatMap(lambda line: unpack(trace_format,line)[0]).collect()

异常信息

23/03/21 14:19:09 ERROR Executor: Exception in task 0.0 in stage 2.0 (TID 2)
java.io.IOException: FixedLengthRecordReader does not support reading compressed files
    at org.apache.spark.input.FixedLengthBinaryRecordReader.initialize(FixedLengthBinaryRecordReader.scala:89)
    at org.apache.spark.rdd.NewHadoopRDD$$anon$1.liftedTree1$1(NewHadoopRDD.scala:221)
    at org.apache.spark.rdd.NewHadoopRDD$$anon$1.<init>(NewHadoopRDD.scala:218)
    at org.apache.spark.rdd.NewHadoopRDD.compute(NewHadoopRDD.scala:173)
    at org.apache.spark.rdd.NewHadoopRDD.compute(NewHadoopRDD.scala:72)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
    at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:136)
    at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)
23/03/21 14:19:09 WARN TaskSetManager: Lost task 0.0 in stage 2.0 (TID 2) (10.110.30.68 executor driver): java.io.IOException: FixedLengthRecordReader does not support reading compressed files
    at org.apache.spark.input.FixedLengthBinaryRecordReader.initialize(FixedLengthBinaryRecordReader.scala:89)
    at org.apache.spark.rdd.NewHadoopRDD$$anon$1.liftedTree1$1(NewHadoopRDD.scala:221)
    at org.apache.spark.rdd.NewHadoopRDD$$anon$1.<init>(NewHadoopRDD.scala:218)
    at org.apache.spark.rdd.NewHadoopRDD.compute(NewHadoopRDD.scala:173)
    at org.apache.spark.rdd.NewHadoopRDD.compute(NewHadoopRDD.scala:72)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
    at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:136)
    at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)

23/03/21 14:19:09 ERROR TaskSetManager: Task 0 in stage 2.0 failed 1 times; aborting job
解决方案

Spark的binaryRecords底层依赖FixedLengthBinaryRecordReader,这个Reader确实不支持直接读取压缩文件——压缩文件无法随机定位,而固定长度读取需要按偏移量精准定位。可以通过以下两种方法解决:

方法一:先解压文件再读取

如果文件规模不大,可先将xz文件解压为原始二进制文件,再用binaryRecords读取:

# 解压xz文件
xz -d test.xz

修改代码读取解压后的文件:

from pyspark.sql import SparkSession
from struct import unpack
spark = SparkSession.builder.master('local[*]').getOrCreate()
sc=spark.sparkContext
trace_format='<Q2B2B4B2Q4Q'
# 读取解压后的原始文件
file=sc.binaryRecords('test',64)
file.flatMap(lambda line: unpack(trace_format,line)[0]).collect()

方法二:读取压缩文件流后手动拆分定长记录

对于大规模文件,先解压不现实,可直接读取压缩文件的字节流,在每个分区内手动拆分固定长度的记录:

from pyspark.sql import SparkSession
from struct import unpack
import lzma

spark = SparkSession.builder.master('local[*]').getOrCreate()
sc=spark.sparkContext

# 读取压缩文件为字节流RDD
compressed_rdd = sc.binaryFiles('test.xz').values()

def split_into_fixed_length(chunk, record_size=64):
    buffer = bytearray()
    # 解压当前分区的字节块
    with lzma.LZMAFile(chunk) as f:
        buffer.extend(f.read())
    # 拆分固定长度记录
    records = []
    for i in range(0, len(buffer), record_size):
        if i + record_size <= len(buffer):
            records.append(buffer[i:i+record_size])
    return records

# 拆分记录并解析
trace_format='<Q2B2B4B2Q4Q'
result = compressed_rdd.flatMap(split_into_fixed_length)\
                      .map(lambda line: unpack(trace_format, line)[0])\
                      .collect()

这种方法的核心是利用binaryFiles读取压缩文件的原始字节,在分区内用lzma库解压,再手动按固定长度拆分记录,避免依赖Spark原生的固定长度Reader。

注意事项

  • 若处理HDFS上的xz文件,需确保每个executor节点都安装了lzma库(可通过pip install lzma安装)
  • 文件过大时,调整Spark的分区数,避免单个分区处理的数据量过大导致内存溢出

内容的提问来源于stack exchange,提问作者Sai Aravind

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 03:07:16