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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 02:52:01