使用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

