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

PySpark写入文件后出现字节缺失问题求助

问题:PySpark排序gensort生成文件后valsort验证失败

使用PySpark加载gensort生成的文件,排序后写入新文件,脚本可正常读取数据并按关键字排序,但用valsort验证输出时触发错误:

sump pump fatal error: pfunc_get_rec: partial record of 90 bytes found at end of input

手动对比正确输出与PySpark生成的输出,vimdiff显示无差异,但diff和cmp命令检测出文件存在差异。

原脚本如下:

#!/usr/bin/python3
import sys
from pyspark.sql import SparkSession

if len(sys.argv) != 5:
    print("Usage: ./sorter.py -i [input filename] -o [output filename]")
    sys.exit(1)

input_filename = sys.argv[2]

output_filename = sys.argv[4]

spark = SparkSession.builder \
                    .master("local[*]") \
                    .appName("sorter") \
                    .getOrCreate()

input_rdd = spark.sparkContext.textFile(input_filename)

print("# partitions: {}".format(input_rdd.getNumPartitions()))

sorted_list = input_rdd.map(lambda x: (x[:10], x[:])) \
                        .sortByKey() \
                        .collect()

with open(output_filename, "w") as ofile:
    for line in sorted_list:
        ofile.write(line[1] + '\n')

问题根源

gensort生成的是固定长度的二进制记录文件(每条记录100字节:10字节key + 90字节value),并非普通文本文件。Spark的textFile会按换行符分割内容,而gensort的记录不含换行符,导致读取时要么将整个文件当成一行,要么在错误位置拆分;后续写入时手动添加\n,彻底破坏了原文件的固定长度结构,最终导致valsort验证时识别到不完整的记录。

解决方案

改用二进制读取方式处理,严格保持记录的固定长度结构:

#!/usr/bin/python3
import sys
from pyspark.sql import SparkSession

if len(sys.argv) != 5:
    print("Usage: ./sorter.py -i [input filename] -o [output filename]")
    sys.exit(1)

input_filename = sys.argv[2]
output_filename = sys.argv[4]

spark = SparkSession.builder \
                    .master("local[*]") \
                    .appName("sorter") \
                    .getOrCreate()

# 读取二进制文件,每个文件对应(文件名, 二进制内容)的元组
binary_rdd = spark.sparkContext.binaryFiles(input_filename)

def split_records(content):
    # gensort每条记录固定100字节
    record_size = 100
    records = []
    # 按固定长度拆分二进制内容
    for i in range(0, len(content), record_size):
        record = content[i:i+record_size]
        if len(record) == record_size:
            # 提取前10字节二进制作为排序key
            key = record[:10]
            records.append((key, record))
    return records

# 拆分记录、排序后收集结果
sorted_records = binary_rdd.flatMap(lambda x: split_records(x[1])) \
                           .sortByKey() \
                           .map(lambda x: x[1]) \
                           .collect()

# 以二进制模式写入,不添加任何额外字符
with open(output_filename, "wb") as ofile:
    for record in sorted_records:
        ofile.write(record)

修改说明

  • 用binaryFiles读取原始二进制内容,避免文本换行分割的干扰
  • 按100字节固定长度拆分记录,确保每条记录完整
  • 使用二进制的前10字节作为排序key,匹配gensort的排序规则
  • 写入时采用二进制模式(wb),直接写入原始记录,不添加换行符等额外内容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 20:15:34