测试HBase Spark分布式扫描示例时出现DoNotRetryIOException异常
我正在尝试运行HBase Spark分布式扫描的Java示例,代码如下:
public class DistributedHBaseScanToRddDemo { public static void main(String[] args) { JavaSparkContext jsc = getJavaSparkContext("hbasetable1"); Configuration hbaseConf = getHbaseConf(0, "", ""); JavaHBaseContext javaHbaseContext = new JavaHBaseContext(jsc, hbaseConf); Scan scan = new Scan(); scan.setCaching(100); JavaRDD<Tuple2<ImmutableBytesWritable, Result>> javaRdd = javaHbaseContext.hbaseRDD(TableName.valueOf("hbasetable1"), scan); List<String> results = javaRdd.map(new ScanConvertFunction()).collect(); System.out.println("Result Size: " + results.size()); } public static Configuration getHbaseConf(int pRimeout, String pQuorumIP, String pClientPort) { Configuration hbaseConf = HBaseConfiguration.create(); hbaseConf.setInt("timeout", 120000); hbaseConf.set("hbase.zookeeper.quorum", "10.56.36.14"); hbaseConf.set("hbase.zookeeper.property.clientPort", "2181"); return hbaseConf; } public static JavaSparkContext getJavaSparkContext(String pTableName) { SparkConf sparkConf = new SparkConf().setAppName("JavaHBaseBulkPut" + pTableName); sparkConf.setMaster("local"); sparkConf.set("spark.testing.memory", "471859200"); JavaSparkContext jsc = new JavaSparkContext(sparkConf); return jsc; } private static class ScanConvertFunction implements Function<Tuple2<ImmutableBytesWritable, Result>, String> { public String call(Tuple2<ImmutableBytesWritable, Result> v1) throws Exception { return Bytes.toString(v1._1().copyBytes()); } } }
执行时抛出如下异常:
Exception in thread "main" org.apache.hadoop.hbase.DoNotRetryIOException: /10.56.48.219:16020 is unable to read call parameter from client 10.56.49.148; java.lang.UnsupportedOperationException: GetRegionLoad at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method) at sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62) at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45) at java.lang.reflect.Constructor.newInstance(Constructor.java:422) at org.apache.hadoop.hbase.ipc.RemoteWithExtrasException.instantiateException(RemoteWithExtrasException.java:93) at org.apache.hadoop.hbase.ipc.RemoteWithExtrasException.unwrapRemoteException(RemoteWithExtrasException.java:83) at org.apache.hadoop.hbase.shaded.protobuf.ProtobufUtil.makeIOExceptionOfException(ProtobufUtil.java:368) at org.apache.hadoop.hbase.shaded.protobuf.ProtobufUtil.getRemoteException(ProtobufUtil.java:345) at org.apache.hadoop.hbase.shaded.protobuf.ProtobufUtil.getRegionLoad(ProtobufUtil.java:1746) at org.apache.hadoop.hbase.client.HBaseAdmin.getRegionLoad(HBaseAdmin.java:2089) at org.apache.hadoop.hbase.mapreduce.RegionSizeCalculator.init(RegionSizeCalculator.java:82) at org.apache.hadoop.hbase.mapreduce.RegionSizeCalculator.<init>(RegionSizeCalculator.java:60) at org.apache.hadoop.hbase.mapreduce.TableInputFormatBase.oneInputSplitPerRegion(TableInputFormatBase.java:293) at org.apache.hadoop.hbase.mapreduce.TableInputFormatBase.getSplits(TableInputFormatBase.java:257) at org.apache.hadoop.hbase.mapreduce.TableInputFormat.getSplits(TableInputFormat.java:254) at org.apache.spark.rdd.NewHadoopRDD.getPartitions(NewHadoopRDD.scala:121) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:248) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:246) at scala.Option.getOrElse(Option.scala:121) at org.apache.spark.rdd.RDD.partitions(RDD.scala:246) at org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:248) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:246) at scala.Option.getOrElse(Option.scala:121) at org.apache.spark.rdd.RDD.partitions(RDD.scala:246) at org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:248) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:246) at scala.Option.getOrElse(Option.scala:121) at org.apache.spark.rdd.RDD.partitions(RDD.scala:246) at org.apache.spark.SparkContext.runJob(SparkContext.scala:1911) at org.apache.spark.rdd.RDD$$anonfun$collect$1.apply(RDD.scala:893) at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112) at org.apache.spark.rdd.RDD.withScope(RDD.scala:358) at org.apache.spark.rdd.RDD.collect(RDD.scala:892) at org.apache.spark.api.java.JavaRDDLike$class.collect(JavaRDDLike.scala:360) at org.apache.spark.api.java.AbstractJavaRDDLike.collect(JavaR...
请问该如何解决此异常?
从异常栈里的核心错误UnsupportedOperationException: GetRegionLoad可以看出,问题出在HBase客户端调用了一个服务端不支持的RPC方法GetRegionLoad。这种情况90%以上是因为HBase客户端与集群服务端的版本不匹配——你的项目依赖的HBase客户端版本比集群上的HBase版本高,导致客户端尝试调用一个服务端还未实现的方法,反之亦然。
另外,异常发生在RegionSizeCalculator.init方法中,这个组件是用来计算HBase表的Region大小,进而生成合理的Spark分区的,它在初始化时会调用getRegionLoad获取Region的负载信息,而服务端不支持这个方法就会抛出异常。
1. 对齐HBase客户端与服务端版本(最优先解决)
这是最根本的解决办法:
- 确认你的HBase集群的版本(比如通过
hbase version命令在集群节点上查看)。 - 修改项目的依赖配置(Maven/Gradle),将所有HBase相关的依赖(
hbase-client、hbase-spark、hbase-common、hbase-server等)的版本调整为和集群完全一致。 - 注意排除掉依赖中可能引入的冲突版本,确保最终打包的jar包中只有一套HBase的类文件。
2. 跳过Region元数据计算(临时 workaround)
如果暂时无法调整版本,可以通过配置跳过RegionSizeCalculator的初始化,避免触发getRegionLoad调用:
在你的getHbaseConf方法中添加如下配置:
hbaseConf.setBoolean("hbase.mapreduce.inputtable.skiptablemetadata", true);
这个配置会让TableInputFormatBase跳过Region大小的计算,直接按每个Region生成一个Spark分区,虽然可能会影响分区的合理性,但可以先让程序跑起来。
3. 确认Spark与HBase的兼容性
除了HBase版本,还要确保你使用的Spark版本和HBase版本是兼容的。比如:
- Spark 2.x通常兼容HBase 1.x和2.x的部分版本;
- Spark 3.x则更适配HBase 2.x及以上版本。
可以参考HBase官方文档中的兼容性矩阵来确认组合是否合适。
4. 验证ZooKeeper连接正确性
虽然不是直接原因,但确保你的hbase.zookeeper.quorum和hbase.zookeeper.property.clientPort配置完全正确,能正常连接到集群的ZooKeeper服务,避免因连接问题间接引发的异常。
内容的提问来源于stack exchange,提问作者Mahesha999

