如何在PySpark中读取指定数据类型的大二进制文件?
用PySpark正确读取固定格式二进制文件
问题背景
有一个固定格式的二进制文件,可通过numpy和pandas读取,代码如下:
import numpy as np import pandas as pd dt = np.dtype([('col1', np.int64), ('col2', np.float32), ('col3', np.int32)]) df = pd.DataFrame(np.fromfile(file, dtype=dt), columns=dt.names)
由于文件可能超出内存,希望直接用PySpark读取,无需先创建pandas DataFrame。尝试过以下方法但遇到问题:
- 使用
spark.read.format("binaryFile"):无法指定列类型,且测试文件也出现java.lang.OutOfMemoryError - 使用
spark.sparkContext.binaryFiles读取后映射:返回的是包含整列数据的单个条目,而非每行记录
解决方案
核心思路是按固定记录长度拆分二进制数据(每条记录占16字节:int64(8) + float32(4) + int32(4)),具体步骤如下:
1. 定义数据类型与记录长度
import numpy as np # 定义与原文件匹配的numpy数据类型 dt = np.dtype([('col1', np.int64), ('col2', np.float32), ('col3', np.int32)]) # 计算单条记录的字节长度 record_length = dt.itemsize # 结果为16
2. 读取二进制文件并拆分记录
使用binaryFiles读取文件(每个元素是(文件名, 二进制内容)),然后将二进制内容按固定长度拆分,过滤掉最后可能不完整的记录:
rdd = spark.sparkContext.binaryFiles(file_path) def split_records(file_content): content = file_content[1] # 获取二进制内容 total_length = len(content) # 计算完整记录数,过滤不完整的尾部数据 num_records = total_length // record_length # 按固定长度拆分二进制数据 for i in range(num_records): start = i * record_length end = start + record_length yield content[start:end] # 拆分后得到每条记录的二进制RDD record_rdd = rdd.flatMap(split_records)
3. 解析二进制记录为元组
用numpy解析每个二进制块,转换成Spark能识别的元组:
def parse_record(binary_data): # 从二进制数据解析出numpy结构化数组 arr = np.frombuffer(binary_data, dtype=dt) # 返回单个记录的元组 return (arr['col1'][0], arr['col2'][0], arr['col3'][0]) parsed_rdd = record_rdd.map(parse_record)
4. 转换为Spark DataFrame
指定Schema将RDD转换为DataFrame:
from pyspark.sql.types import StructType, StructField, LongType, FloatType, IntegerType # 定义DataFrame Schema schema = StructType([ StructField("col1", LongType(), nullable=False), StructField("col2", FloatType(), nullable=False), StructField("col3", IntegerType(), nullable=False) ]) # 转换为DataFrame df = spark.createDataFrame(parsed_rdd, schema=schema) # 查看前5行验证 df.show(5)
执行后得到期望结果:
+-------------+--------+-------+ | col1| col2| col3| +-------------+--------+-------+ |2317613314222|1.551823| 556| |2317614400012|1.206112| 131609| |2317615429391|1.022747| 131888| |2317621608598|2.082569| 131643| |2317622053589|1.260681| 271| +-------------+--------+-------+
关键说明
- 避免使用
wholeTextFiles:该方法会将整个文件加载到内存,大文件容易OOM;binaryFiles按文件加载后拆分处理,更灵活 - 固定长度拆分:必须确保每条记录的字节数准确,否则解析会出错
- 内存优化:拆分后每条记录独立处理,不会一次性加载整个文件到内存,适合大文件场景
内容的提问来源于stack exchange,提问作者BBG
相关产品推荐
相关产品推荐

