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

Spark如何读取orc文件为RDD格式而非DataFrame格式

读取ORC为RDD的实现方案

完全可以直接读取ORC文件为RDD格式,无需经过DataFrame转换环节,避免中间过程的内存开销,不同API的实现方式如下:

Scala API

直接调用Hadoop原生ORC输入格式读取:

import org.apache.hadoop.hive.ql.io.orc.OrcNewInputFormat
import org.apache.hadoop.io.NullWritable
import org.apache.hadoop.hive.ql.io.orc.OrcStruct

// 直接读取ORC为RDD[OrcStruct]
val orcRDD = sc.newAPIHadoopFile(
  "hdfs://your-orc-file-path/*",
  classOf[OrcNewInputFormat],
  classOf[NullWritable],
  classOf[OrcStruct]
).map(_._2)

// 提取字段示例,无需转Row
val parsedRDD = orcRDD.map { struct =>
  val userId = struct.getFieldValue("user_id").toString
  val payAmount = struct.getFieldValue("pay_amount").asInstanceOf[Long]
  (userId, payAmount)
}

PySpark API

调用Hadoop InputFormat接口实现读取:

from pyspark import SparkContext

sc = SparkContext.getOrCreate()
# 读取ORC为RDD
orc_rdd = sc.newAPIHadoopFile(
    "hdfs://your-orc-file-path/*",
    inputFormatClass="org.apache.hadoop.hive.ql.io.orc.OrcNewInputFormat",
    keyClass="org.apache.hadoop.io.NullWritable",
    valueClass="org.apache.hadoop.hive.ql.io.orc.OrcStruct"
).map(lambda x: x[1])
更推荐的优化方案:优化DataFrame读取参数避免OOM

实际上DataFrame内置了Tungsten内存管理和Catalyst优化,内存效率远高于手动处理RDD,大部分场景下读取ORC触发OOM是参数配置不合理导致,你可以先尝试调整以下参数,无需切换RDD:

  • 缩小单分区最大数据量,避免单个分片过大占满Executor内存:
    --conf spark.sql.files.maxPartitionBytes=33554432 # 单分区最大32M,默认是128M,可根据实际情况调低
    --conf spark.sql.files.openCostInBytes=4194304
    
  • 列裁剪+谓词下推,利用ORC列存特性只加载需要的字段,跳过不需要的数据读取:
    // 只加载需要的3列,同时过滤不需要的数据,ORC会直接跳过无关列和无关块的读取,内存开销大幅降低
    val df = spark.read.orc("path")
      .select("user_id", "pay_amount", "create_time")
      .where("create_time >= '2024-01-01'")
    
  • 补充Executor内存配置,你当前的配置未设置Executor内存,建议根据集群资源调整:
    --conf spark.executor.memory=8g
    --conf spark.executor.memoryOverhead=2g
    
现有配置优化建议

你当前的配置可以补充以下参数,不管用RDD还是DataFrame都能降低OOM概率:

  • 优化Kryo序列化配置,降低序列化内存开销:
    --conf spark.kryo.classesToRegister=org.apache.hadoop.hive.ql.io.orc.OrcStruct
    
  • 调整内存分配比例,给计算环节预留更多内存:
    --conf spark.memory.fraction=0.7
    --conf spark.memory.storageFraction=0.2
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 04:18:05