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

使用newAPIHadoopRDD读取HBase未刷写Memstore数据无显示问题排查

问题:Spark通过newAPIHadoopRDD无法读取HBase Memstore中未刷写的记录

现象说明

  • 读取HBase中数百条未刷写到HDFS的记录时,Spark无法获取这些数据;通过HBase Shell强制刷写或系统触发刷写(Memstore大小超过阈值)后,这些记录能在Spark中正常显示。
  • 首次刷写后,新写入的内容可以立即在Spark中读取到。
  • 未刷写的Memstore记录,HBase原生客户端的Scan API可以正常读取。
  • 因提前未知Schema,无法使用spark-hbase connector,只能采用newAPIHadoopRDD API实现读取。

现有Spark读取代码

@transient val conf : Configuration = HBaseConfiguration.create()
conf.set("hbase.zookeeper.quorum", ZOOKEEPER_QUORUM);
conf.set("hbase.zookeeper.property.clientPort", ZOOKEEPER_PORT);
conf.set(TableInputFormat.INPUT_TABLE, "test")
val hbaseValueRDD = spark.sparkContext.newAPIHadoopRDD(conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result]).map(x => x._2)
val hbaseValueFlattenedRDD = hbaseValueRDD.map(x => inputSplitter(x)).flatMap(x => x)
var rowDF = hbaseValueFlattenedRDD.toDF(COL_KEY, COL_VALUE)

尝试过的无效方案

曾尝试传递自定义Scan对象并设置参数(比如强制从replica-0读取),但无效果:

// Function to convert Scan object to a base64 string
def convertScanToString(scan: Scan): String = {
  val proto = ProtobufUtil.toScan(scan)
  Base64.encodeBytes(proto.toByteArray)
}
@transient val scan = new Scan()
scan.setReplicaId(0)
@transient val scanStr = convertScanToString(scan)
@transient val conf : Configuration = HBaseConfiguration.create()
conf.set("hbase.zookeeper.quorum", ZOOKEEPER_QUORUM);
conf.set("hbase.zookeeper.property.clientPort", ZOOKEEPER_PORT);
conf.set(TableInputFormat.INPUT_TABLE, "test")
spark.sparkContext.newAPIHadoopRDD(conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result])

可正常读取的HBase原生客户端代码

val table = connection.getTable(TableName.valueOf(Bytes.toBytes("test")))
val scan = table.getScanner(new Scan())
scan.asScala.foreach(result => {
    println(result)
})

原因分析

Spark使用的TableInputFormat默认基于HBase RegionServer的HFile扫描逻辑,只会读取已经持久化到HDFS的磁盘数据,不会主动加载Memstore中的未刷写内存数据。而HBase原生Scan API会直接与RegionServer的Memstore和HFile交互,合并返回内存与磁盘中的全量数据,所以能读取到未刷写的记录。

首次刷写后新写入的内容能被Spark读取,是因为此时TableInputFormat的扫描逻辑会感知到Region的最新状态,后续的内存数据可能在扫描时被临时纳入读取范围(但这个逻辑不保证稳定)。


解决方案

1. 配置参数强制读取Memstore

在HBase Configuration中添加hbase.mapreduce.scan.read.memstore参数并设为true,强制TableInputFormat读取Memstore中的未刷写数据:

@transient val conf : Configuration = HBaseConfiguration.create()
conf.set("hbase.zookeeper.quorum", ZOOKEEPER_QUORUM);
conf.set("hbase.zookeeper.property.clientPort", ZOOKEEPER_PORT);
conf.set(TableInputFormat.INPUT_TABLE, "test")
// 开启读取Memstore数据
conf.set("hbase.mapreduce.scan.read.memstore", "true")
val hbaseValueRDD = spark.sparkContext.newAPIHadoopRDD(conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result]).map(x => x._2)
// 后续处理代码不变

该参数默认值为false,开启后会让TableInputFormat在扫描时合并Memstore和HFile的数据。

2. 自定义InputFormat(进阶方案)

如果上述参数无效,可以自定义一个基于HBase原生Scan逻辑的InputFormat,直接复用原生Scan读取全量数据的逻辑,再在Spark中使用这个自定义InputFormat。这种方式需要额外的代码开发,但能完全对齐原生客户端的读取行为。

3. 手动触发刷写(临时应急方案)

如果前两种方案无法快速落地,可以在Spark读取前通过HBase API手动触发刷写,但频繁刷写会增加IO开销,仅适合临时场景:

val admin = connection.getAdmin()
admin.flush(TableName.valueOf("test"))

内容的提问来源于stack exchange,提问作者Manjunath NS

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 11:24:51