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

