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

