Cloudera环境Spark2连接HBase不被支持,求数据插入解决办法
当然有可行方案!我之前在CDH集群上用Spark2操作HBase的时候也碰到过这个问题,下面是几个经过实践验证的方法,你可以根据自己的场景来选:
方案1:直接使用HBase原生Java API
这种方法绕过CDH自带的Spark On HBase组件,直接在Spark任务里调用HBase的客户端API来实现数据写入,兼容性最好,几乎所有HBase版本都支持。
- 步骤:
- 先在你的Spark项目里引入HBase客户端依赖,一定要注意版本和集群的HBase版本完全匹配,比如Maven依赖可以这么加:
<dependency> <groupId>org.apache.hbase</groupId> <artifactId>hbase-client</artifactId> <version>你的HBase版本号</version> </dependency> <dependency> <groupId>org.apache.hbase</groupId> <artifactId>hbase-common</artifactId> <version>你的HBase版本号</version> </dependency> - 编写Spark代码时,建议用
mapPartitions来初始化HBase连接(避免每条数据都创建连接,导致资源浪费),示例代码如下:import org.apache.hadoop.hbase.{HBaseConfiguration, TableName} import org.apache.hadoop.hbase.client.{Connection, ConnectionFactory, Put} import org.apache.hadoop.hbase.util.Bytes import org.apache.spark.sql.SparkSession object Spark2HBaseWriter { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("Spark2HBaseInsert").getOrCreate() import spark.implicits._ // 模拟需要插入的业务数据 val data = Seq(("user_001", "info", "name", "Alice"), ("user_002", "info", "age", "25")) val df = data.toDF("rowkey", "cf", "col", "value") // 配置HBase的ZK地址等信息 val hbaseConf = HBaseConfiguration.create() hbaseConf.set("hbase.zookeeper.quorum", "zk-node1,zk-node2,zk-node3") hbaseConf.set("hbase.zookeeper.property.clientPort", "2181") df.rdd.foreachPartition { partitionData => var connection: Connection = null try { // 每个分区创建一次HBase连接 connection = ConnectionFactory.createConnection(hbaseConf) val table = connection.getTable(TableName.valueOf("user_info")) partitionData.foreach { row => val put = new Put(Bytes.toBytes(row.getAs[String]("rowkey"))) put.addColumn( Bytes.toBytes(row.getAs[String]("cf")), Bytes.toBytes(row.getAs[String]("col")), Bytes.toBytes(row.getAs[String]("value")) ) table.put(put) } table.close() } catch { case e: Exception => e.printStackTrace() } finally { // 确保连接关闭 if (connection != null) connection.close() } } spark.stop() } } - 打包提交任务时,如果集群环境没有自带HBase依赖,可以用
spark-submit的--jars参数指定HBase的客户端jar包,或者把依赖打包到你的应用jar里。
- 先在你的Spark项目里引入HBase客户端依赖,一定要注意版本和集群的HBase版本完全匹配,比如Maven依赖可以这么加:
方案2:通过Apache Phoenix作为中间层
Phoenix是HBase的SQL层,提供了JDBC接口,Spark2可以通过Phoenix的Spark插件轻松读写HBase数据,这种方法的好处是可以用SQL来操作HBase,学习成本较低。
- 步骤:
- 先确保你的CDH集群已经部署了Phoenix,并且版本和HBase兼容(CDH一般会有对应的Phoenix版本包)。
- 在Spark项目中引入Phoenix的Spark依赖:
<dependency> <groupId>org.apache.phoenix</groupId> <artifactId>phoenix-spark</artifactId> <version>你的Phoenix版本号</version> </dependency> - 编写代码时直接用Spark的DataFrame API写入,示例如下:
import org.apache.spark.sql.SparkSession object Spark2PhoenixWriter { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("Spark2PhoenixInsert").getOrCreate() import spark.implicits._ // 模拟数据,注意字段名要和Phoenix表的列名对应 val data = Seq(("user_001", "Alice"), ("user_002", "Bob")) val df = data.toDF("ROWKEY", "NAME") // 写入Phoenix对应的HBase表 df.write .format("org.apache.phoenix.spark") .mode("append") .option("table", "USER_INFO") // Phoenix中创建的表名 .option("zkUrl", "zk-node1,zk-node2,zk-node3:2181") .save() spark.stop() } }
CREATE TABLE USER_INFO (ROWKEY VARCHAR PRIMARY KEY, NAME VARCHAR) COLUMN_ENCODED_BYTES=0;来创建表。
方案3:使用HBase官方的Spark Connector
HBase官方后来推出了hbase-spark模块,支持Spark2,和CDH自带的Spark On HBase不是同一个组件,你可以单独引入这个依赖来使用。
- 步骤:
- 引入对应的依赖,版本要匹配你的Spark2和HBase版本:
<dependency> <groupId>org.apache.hbase</groupId> <artifactId>hbase-spark</artifactId> <version>你的HBase版本号</version> </dependency> - 使用
HBaseTableCatalog来定义表结构,然后用DataFrame API写入,示例代码:import org.apache.hadoop.hbase.spark.HBaseTableCatalog import org.apache.spark.sql.SparkSession object Spark2HBaseConnectorWriter { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("Spark2HBaseConnectorInsert").getOrCreate() import spark.implicits._ // 模拟数据 val data = Seq(("user_001", "25"), ("user_002", "30")) val df = data.toDF("rowkey", "age") // 定义HBase表的Catalog,描述表名、rowkey、列族和列的映射 val catalog = s"""{ | "table":{"namespace":"default", "name":"user_info"}, | "rowkey":"rowkey", | "columns":{ | "rowkey":{"cf":"rowkey", "col":"rowkey", "type":"string"}, | "age":{"cf":"info", "col":"age", "type":"string"} | } |}""".stripMargin df.write .options(Map( HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "5" // 如果表不存在则创建,指定预分区数 )) .format("org.apache.hadoop.hbase.spark") .save() spark.stop() } }
- 引入对应的依赖,版本要匹配你的Spark2和HBase版本:
内容的提问来源于stack exchange,提问作者yAsH
相关产品推荐
相关产品推荐

