Scala中使用newAPIHadoopRDD读取HBase表任务挂起问题求助
我尝试用newAPIHadoopRDD读取HBase表t1,代码如下:
import org.apache.hadoop.hbase.{HBaseConfiguration, HTableDescriptor} import org.apache.hadoop.hbase.client.{HBaseAdmin, Result} import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapreduce.TableInputFormat val tableName = "t1" val hconf = HBaseConfiguration.create() hconf.set(TableInputFormat.INPUT_TABLE, "t1") val hBaseRDD = sc.newAPIHadoopRDD(hconf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result]) println("records found : " + hBaseRDD.count())
但执行任何action操作(比如count()、collect())时,程序直接挂起,既不报错也不返回结果,就像进入了死循环。按Ctrl+C终止后得到以下报错:
scala> hBaseRDD.count 18/01/30 07:23:15 WARN repl.Signaling: Cancelling all active jobs, this can take a while. Press Ctrl+C again to exit now. 18/01/30 07:23:15 WARN spark.ExecutorAllocationManager: No stages are running, but numRunningTasks != 0 org.apache.spark.SparkException: Job 0 cancelled as part of cancellation of all jobs at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1430) at org.apache.spark.scheduler.DAGScheduler.handleJobCancellation(DAGScheduler.scala:1370) at org.apache.spark.scheduler.DAGScheduler$$anonfun$doCancelAllJobs$1.apply$mcVI$sp(DAGScheduler.scala:716) at org.apache.spark.scheduler.DAGScheduler$$anonfun$doCancelAllJobs$1.apply(DAGScheduler.scala:716) at org.apache.spark.scheduler.DAGScheduler$$anonfun$doCancelAllJobs$1.apply(DAGScheduler.scala:716) at scala.collection.mutable.HashSet.foreach(HashSet.scala:78) at org.apache.spark.scheduler.DAGScheduler.doCancelAllJobs(DAGScheduler.scala:716) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:1623) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1600) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1589) at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48) at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:623) at org.apache.spark.SparkContext.runJob(SparkContext.scala:1930) at org.apache.spark.SparkContext.runJob(SparkContext.scala:1943) at org.apache.spark.SparkContext.runJob(SparkContext.scala:1956) at org.apache.spark.SparkContext.runJob(SparkContext.scala:1970) at org.apache.spark.rdd.RDD.count(RDD.scala:1157) ... 48 elided
有没有人遇到过同样的问题?或者知道怎么解决?
我之前也碰到过类似的无响应情况,通常是几个常见原因导致的,你可以逐一排查:
HBase配置未正确加载:
默认HBaseConfiguration.create()只会加载基础默认配置,如果是集群环境或者HBase配置文件不在默认路径,必须手动把hbase-site.xml加入Spark的classpath。比如启动spark-shell时加--files /path/to/hbase-site.xml,或者在代码里通过hconf.addResource(new Path("/path/to/hbase-site.xml"))加载。这一步至关重要,因为Spark需要知道HBase的ZooKeeper地址等核心信息才能建立连接。依赖版本不兼容:
Spark和HBase的版本兼容性问题是这类无响应的高发原因。比如Spark用2.x但HBase是1.x,或者两者依赖的Hadoop客户端版本不匹配。你需要确保Spark依赖中包含对应HBase版本的客户端jar包,版本要严格对齐。比如启动spark-shell时可以加依赖参数:--packages org.apache.hbase:hbase-client:1.4.12,org.apache.hbase:hbase-common:1.4.12(根据你的实际HBase版本调整)。ZooKeeper连接异常:
如果Spark连不上HBase的ZooKeeper集群,就会无限等待连接。你可以先用zkCli.sh手动测试ZooKeeper的连通性,看看能不能找到HBase的节点(比如/hbase)。同时检查代码里的配置,确保hbase.zookeeper.quorum参数正确设置:hconf.set("hbase.zookeeper.quorum", "zk-node1,zk-node2,zk-node3")。集群资源不足:
如果Spark集群的executor内存、CPU核数不够,任务可能会卡在调度阶段。你可以尝试调整资源参数,比如启动spark-shell时加--executor-memory 4g --num-executors 2 --executor-cores 2,给任务分配足够的运行资源。表权限或存在性问题:
先确认HBase中t1表确实存在,且Spark运行的用户有读取权限。可以用HBase命令行工具hbase shell执行list查看表列表,scan t1测试表的可访问性。
建议先从HBase配置和版本兼容性入手排查,这两个是最常见的诱因。如果还是不行,去查看Spark的driver日志和executor日志,里面会有更详细的等待原因,比如连接超时、资源争抢等信息。
内容的提问来源于stack exchange,提问作者Vincenzo

