遍历Spark DataFrame耗时过长且触发OutOfMemoryError: GC开销超限错误
问题分析与解决方案
兄弟,你的问题我太熟了——用Spark处理大文件写HBase踩过的坑基本全中了,咱们一步步来拆解问题,然后给你落地的优化方案:
首先看你现有代码的几个致命问题:
collect()直接把数据拉到Driver端:哪怕你加了limit(100000),collect()也会把这10万条数据全塞到Driver的内存里,500万条的话直接爆内存是板上钉钉的事。而且collect()之后的foreach是在Driver单线程跑,完全浪费了Spark的分布式集群能力,速度能快才怪。- 单HBase连接串行写入:你在Driver创建了一个
hTable,然后单线程一条条写HBase,这种串行操作根本发挥不出HBase的写入性能,相当于用跑车拉货却只让一个轮子转。 - RowKey用自增数字:这会导致HBase的热点问题——所有写入请求都挤到同一个RegionServer,直接把服务器压垮,写入速度慢到离谱。
接下来给你三个从基础到进阶的解决方案,你可以根据自己的环境选:
方案一:紧急修复——用foreachPartition替代collect().foreach
这是最基础的修复,不用改太多代码,就能把写入改成分布式的:
import org.apache.hadoop.hbase.client.{HTable, Put} import org.apache.hadoop.hbase.util.Bytes dfXML.foreachPartition { partition => // 每个Partition单独创建HBase连接,避免跨节点共享连接的问题 val hBaseConf = HBaseConfiguration.create() hBaseConf.set("hadoop.security.authentication", "kerberos") hBaseConf.set("hbase.zookeeper.quorum", cluster) hBaseConf.set("hbase.zookeeper.property.client.port", "2181") UserGroupInformation.setConfiguration(hBaseConf) UserGroupInformation.loginUserFromKeytab("user", "file:///tmp/keytab.keytab") val hTable = new HTable(hBaseConf, "ns:table_name") hTable.setAutoFlush(false, true) // 关闭自动刷写,攒够一批再提交 hTable.setWriteBufferSize(1024 * 1024 * 5) // 设置5MB的写入缓冲区,可根据集群情况调大到10-20MB partition.foreach { elem => // 用Option处理空值,比你原来的if-else简洁多了 val A = Option(elem.getString(0)).getOrElse("") val B = Option(elem.getString(1)).getOrElse("") val C = Option(elem.getString(2)).getOrElse("") val D = Option(elem.getString(3)).getOrElse("") val as_of_date = elem.getDate(4).toString val last_updated_date = elem.getTimestamp(5).toString // 重点!别用自增数字当RowKey!这里给你个示例,用业务字段+时间戳,让写入均匀分布 val rowKey = s"${elem.getString(0)}_${System.currentTimeMillis()}" val put = new Put(Bytes.toBytes(rowKey)) // 统一用Bytes工具类处理字节转换,避免乱码 put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("A"), Bytes.toBytes(A)) put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("B"), Bytes.toBytes(B)) put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("C"), Bytes.toBytes(C)) put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("D"), Bytes.toBytes(D)) put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("as_of_date"), Bytes.toBytes(as_of_date)) put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("last_updated_date"), Bytes.toBytes(last_updated_date)) hTable.put(put) } hTable.flushCommits() // 每个Partition处理完手动提交剩余数据 hTable.close() // 用完关闭连接,避免资源泄漏 }
核心改进点:
foreachPartition:让每个Executor的Partition独立处理数据,分布式写入,把集群的算力用起来- 每个Partition创建一次连接:避免频繁创建销毁连接的开销,也不会在Driver端堆积数据
- 批量写入:关闭自动刷写+设置缓冲区,把多条Put攒成一批提交,大幅提升写入吞吐量
- 优化RowKey:彻底解决HBase热点问题,让写入请求均匀分散到各个RegionServer
方案二:优雅进阶——用Spark官方HBase Connector
Spark官方提供了HBase集成的Connector,不用手动管理连接和Put操作,代码简洁,性能更稳定,推荐优先用这个:
第一步:添加依赖(如果用Maven)
<dependency> <groupId>org.apache.hbase</groupId> <artifactId>hbase-spark</artifactId> <version>你的HBase版本,比如2.4.9,要和集群版本匹配</version> </dependency>
第二步:改写代码
import org.apache.hadoop.hbase.spark.HBaseTableCatalog import org.apache.spark.sql.functions._ // 定义HBase表的元数据Catalog,相当于告诉Spark怎么映射DataFrame字段到HBase的列 val catalog = s"""{ |"table":{"namespace":"ns", "name":"table_name"}, |"rowkey":"rowkey", |"columns":{ | "rowkey":{"cf":"rowkey", "col":"rowkey", "type":"string"}, | "A":{"cf":"cf", "col":"A", "type":"string"}, | "B":{"cf":"cf", "col":"B", "type":"string"}, | "C":{"cf":"cf", "col":"C", "type":"string"}, | "D":{"cf":"cf", "col":"D", "type":"string"}, | "as_of_date":{"cf":"cf", "col":"as_of_date", "type":"string"}, | "last_updated_date":{"cf":"cf", "col":"last_updated_date", "type":"string"} |} |}""".stripMargin // 生成合理的RowKey,这里用业务字段A+时间戳,你可以根据自己的业务调整 val dfWithRowKey = dfXML.withColumn("rowkey", concat(col("A"), lit("_"), current_timestamp().cast("string"))) // 写入HBase,newTable参数可以指定预分区数,比如5,根据你的数据量调整 dfWithRowKey.write.options( Map(HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "5") ).format("org.apache.hadoop.hbase.spark").save()
优势:
- 无需手动管理连接:Spark自动帮你处理连接池、批量提交等底层逻辑,减少出错概率
- 支持预分区:创建表时直接指定预分区数,提前优化HBase的Region分布,从根源避免热点
- 代码极简:几行代码搞定写入,不用写繁琐的Put和连接逻辑
方案三:极致优化——针对500万条数据的调优
如果要处理500万条数据,还要进一步优化集群和配置:
XML读取优化:
- 调整Spark分区数:把
spark.sql.shuffle.partitions设为集群节点数×核数×2(比如10节点×4核=40,就设为80),让数据均匀分布到各个Executor - 如果你是多个XML文件,开启
option("recursiveFileLookup", "true")并行读取
- 调整Spark分区数:把
Spark资源配置优化:
- Driver内存:
--driver-memory 16G(根据Driver节点的内存调整,至少8G以上) - Executor配置:
--executor-memory 8G --executor-cores 4(每个Executor分配足够的内存和核数,避免OOM)
- Driver内存:
HBase集群优化:
- 调整RegionServer并发数:把
hbase.regionserver.handler.count设为100(默认是30),提升HBase的写入并发能力 - 开启WAL异步刷写:设置
hbase.regionserver.wal.async为true,提升写入速度(注意:如果集群断电可能丢失少量数据,根据业务容忍度调整)
- 调整RegionServer并发数:把
总结
优先选方案二的官方Connector,代码简洁还稳定;如果因为环境限制不能用Connector,就用方案一的foreachPartition紧急修复。另外一定要注意RowKey的设计,这是HBase写入性能的关键,千万不能用自增数字!
内容的提问来源于stack exchange,提问作者CeeJay
相关产品推荐
相关产品推荐

