Zeppelin Spark解释器写入HBase失败求技术协助
我明白你遇到的困扰——本地Spark环境能正常跑的写HBase代码,放到Zeppelin里就失效了。结合你的代码片段和常见的Zeppelin+Spark+HBase集成问题,咱们一步步排查解决:
1. 先揪出代码里的关键错误
你的getJob方法里有个致命的配置错误:job.setOutputValueClass(classOf[Result]),这完全不对!TableOutputFormat需要的输出值类型是Put,不是Result(Result是读HBase时用的)。这大概率是导致代码在Zeppelin里跑不通的核心原因之一。
修正后的getJob方法应该是这样:
def getJob(tableName: String): Job = { sc.hadoopConfiguration.set("hbase.zookeeper.quorum", "dn1.hadoop,dn2.hadoop,dn3.hadoop,dn4.hadoop,dn5.hadoop") sc.hadoopConfiguration.set("hbase.zookeeper.property.clientPort", "2181") sc.hadoopConfiguration.set(TableOutputFormat.OUTPUT_TABLE, tableName) val job = Job.getInstance(sc.hadoopConfiguration) job.setOutputKeyClass(classOf[ImmutableBytesWritable]) job.setOutputValueClass(classOf[Put]) // 这里改成Put! job.setOutputFormatClass(classOf[TableOutputFormat[ImmutableBytesWritable]]) job }
另外,你的mapPartitions逻辑里,最终要返回(ImmutableBytesWritable, Put)类型的元组,而且addColumn的参数要完整(你代码里it.getAs[String](3...是截断的),修正后的处理逻辑:
val hbaseRDD = list.mapPartitions( iter => { iter.map( it => { val rowkey = it.getAs[String](8).getBytes val put = new Put(rowkey) // 完整的列添加逻辑,注意转成字节数组 put.addColumn(Bytes.toBytes("A"), Bytes.toBytes("CI"), Bytes.toBytes(it.getAs[String](3))) // 如果有其他列,继续添加 // put.addColumn(Bytes.toBytes("A"), Bytes.toBytes("COL2"), Bytes.toBytes(it.getAs[String](4))) (new ImmutableBytesWritable, put) // 返回TableOutputFormat需要的元组 } ) } ) // 最后执行保存 hbaseRDD.saveAsNewAPIHadoopDataset(getJob("你的HBase表名").getConfiguration)
2. 检查Zeppelin的HBase依赖是否缺失
本地Spark环境你可能手动引入了HBase的依赖包,但Zeppelin的Spark解释器默认不一定包含这些。解决方法有两种:
- 方法一:用Zeppelin的%dep魔法命令在线加载依赖
在你的notebook开头添加这段代码(替换成你实际的HBase版本):%dep z.load("org.apache.hbase:hbase-client:2.4.11") z.load("org.apache.hbase:hbase-common:2.4.11") z.load("org.apache.hbase:hbase-server:2.4.11") z.load("org.apache.hbase:hbase-mapreduce:2.4.11") - 方法二:手动添加本地依赖
把HBase集群的hbase-site.xml和相关依赖jar包(比如hbase-client-*.jar、hbase-mapreduce-*.jar)复制到Zeppelin的Spark解释器目录(比如$ZEPPELIN_HOME/interpreter/spark/),然后重启Zeppelin服务。
3. 验证HBase配置是否生效
Zeppelin的Spark上下文可能没有正确加载HBase的配置,你可以在代码里加一行打印验证:
println("ZooKeeper地址:" + sc.hadoopConfiguration.get("hbase.zookeeper.quorum"))
如果输出为空或者不对,说明配置没生效,建议把HBase集群的hbase-site.xml放到Zeppelin的Spark配置目录,确保配置能被加载。
4. 排查权限问题
Zeppelin的运行用户可能没有写入目标HBase表的权限。你可以用HBase Shell给用户赋权:
grant 'zeppelin运行用户', 'RW', '你的HBase表名'
5. 用日志定位剩余问题
如果以上步骤都试过还是不行,打开Zeppelin的调试日志:
sc.setLogLevel("DEBUG")
然后运行代码,查看Zeppelin日志目录($ZEPPELIN_HOME/logs/)里的Spark解释器日志,里面会有具体的错误信息(比如连接超时、列族不存在、权限拒绝等),根据日志再针对性解决。
内容的提问来源于stack exchange,提问作者vanben

