☰
Spark 读写 HBase 实战:saveAsHadoopDataset 与 saveAsNewAPIHadoopDataset 配置骨架
2026/9/29 3:48:49 网站建设 项目流程

1. Spark 读写 HBase 的真实痛点:两个 API 到底该用哪个

如果你正在用 Spark 把 RDD 批量灌进 HBase,或者从 HBase 拉数据转成 RDD 做分析,大概率会在saveAsHadoopDataset和saveAsNewAPIHadoopDataset这两个方法上卡一下。它们名字只差一个NewAPI,但背后走的是 Hadoop 两套完全不同的 OutputFormat 体系:一套是老的org.apache.hadoop.mapred,一套是新的org.apache.hadoop.mapreduce。选错了,轻则编译报错,重则任务跑起来一直重连 ZooKeeper 却不报明确错误。

这篇就聚焦 Spark RDD 与 HBase 之间的批量读写链路,把两个写入 API 的适用差异讲清楚,给出可直接复制的 SparkConf / HBaseConfiguration 配置骨架、表名与列族映射示例,以及写入后 scan 校验、读取后 count 校验的验证动作。适合已经能跑 Spark 本地任务、想打通 HBase 读写但被配置项绕晕的同学。核心检索词就三个:spark、hbase、rdd,围绕它们把链路走通。

我试过在本地 IDE 里连远程 HBase 集群,最容易出问题的不是业务逻辑,而是依赖包和 ZooKeeper 地址。下面按“先配环境、再写数据、再读数据、最后排障”的顺序展开,每一步都给可复制的骨架。

2. 前置准备:依赖、ZooKeeper 连接与建表

2.1 依赖包别导错包名

老 API 和新 API 的类名高度相似,但包路径不同,这是最常见的坑。写入时:

  • 老 API 用org.apache.hadoop.hbase.mapred.TableOutputFormat
  • 新 API 用org.apache.hadoop.hbase.mapreduce.TableOutputFormat

读取时统一用org.apache.hadoop.hbase.mapreduce.TableInputFormat。如果你把mapred的 OutputFormat 传给saveAsNewAPIHadoopDataset,运行时会直接抛类型不匹配。

classpath 里需要补齐这些 jar:HBase lib 目录下的hbase开头包、hadoop开头包,另外几个容易被漏掉但缺了会出问题的:zookeeper-3.4.6.jar、metrics-core-2.2.0.jar、htrace-core-3.1.0-incubating.jar、guava-12.0.1.jar。少了 metrics 或 htrace,典型现象是日志里RpcRetryingCaller: Call exception反复刷,任务不报错但一直重连。Spark 侧还需要spark-assembly对应的程序集 jar。

2.2 ZooKeeper 连接两种方式

Spark 应用要连 HBase,本质是先连 ZooKeeper 再借它找到 HBase。两种配法:

第一种是把hbase-site.xml放进 classpath,让 HBaseConfiguration 自动加载。第二种是在代码里显式 set。不配的话默认连localhost:2181,直接connection refused。本文用第二种,方便在 IDE 里改。

val conf = HBaseConfiguration.create() conf.set("hbase.zookeeper.quorum", "slave1,slave2,slave3") conf.set("hbase.zookeeper.property.clientPort", "2181")

2.3 表先在 HBase Shell 建好

虽然代码里能用 HBaseAdmin 建表,但不建议在 Spark 任务里做。表结构变更和写入混在一起,出问题不好定位。先在 shell 里建:

create 'account', 'cf'

表名account,列族cf。后面所有读写都围绕它。

3. 可复制配置:两个写入 API 的骨架对比

3.1 saveAsHadoopDataset(老 API)

老 API 的 RDD 元素类型必须是(ImmutableBytesWritable, Put),且 OutputFormat 来自mapred包。JobConf 直接由 HBaseConfiguration 构造。

import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.client.Put import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapred.TableOutputFormat import org.apache.hadoop.hbase.util.Bytes import org.apache.hadoop.mapred.JobConf import org.apache.spark.{SparkConf, SparkContext} object WriteOldAPI { def main(args: Array[String]): Unit = { val sparkConf = new SparkConf().setAppName("HBaseWriteOld").setMaster("local") val sc = new SparkContext(sparkConf) val conf = HBaseConfiguration.create() conf.set("hbase.zookeeper.quorum", "slave1,slave2,slave3") conf.set("hbase.zookeeper.property.clientPort", "2181") val tableName = "account" val jobConf = new JobConf(conf) jobConf.setOutputFormat(classOf[TableOutputFormat]) jobConf.set(TableOutputFormat.OUTPUT_TABLE, tableName) val indataRDD = sc.makeRDD(Array("1,jack,15", "2,Lily,16", "3,mike,16")) val rdd = indataRDD.map(_.split(',')).map { arr => val put = new Put(Bytes.toBytes(arr(0))) put.add(Bytes.toBytes("cf"), Bytes.toBytes("name"), Bytes.toBytes(arr(1))) put.add(Bytes.toBytes("cf"), Bytes.toBytes("age"), Bytes.toBytes(arr(2).toInt)) (new ImmutableBytesWritable, put) } rdd.saveAsHadoopDataset(jobConf) sc.stop() } }

注意put.add三个参数依次是列族、列名、值,值必须用Bytes.toBytes转。行键这里用字符串,和读取时保持一致。

3.2 saveAsNewAPIHadoopDataset(新 API)

新 API 走mapreduce包,配置挂在sc.hadoopConfiguration上,Job 从它构造。RDD 元素类型同样是(ImmutableBytesWritable, Put),但 OutputFormat 换成mapreduce版本。

import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.client.Put import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapreduce.TableOutputFormat import org.apache.hadoop.hbase.util.Bytes import org.apache.hadoop.mapreduce.Job import org.apache.spark.{SparkConf, SparkContext} object WriteNewAPI { def main(args: Array[String]): Unit = { val sparkConf = new SparkConf().setAppName("HBaseWriteNew").setMaster("local") val sc = new SparkContext(sparkConf) val tableName = "account" sc.hadoopConfiguration.set("hbase.zookeeper.quorum", "slave1,slave2,slave3") sc.hadoopConfiguration.set("hbase.zookeeper.property.clientPort", "2181") sc.hadoopConfiguration.set(TableOutputFormat.OUTPUT_TABLE, tableName) val job = new Job(sc.hadoopConfiguration) job.setOutputKeyClass(classOf[ImmutableBytesWritable]) job.setOutputValueClass(classOf[Put]) job.setOutputFormatClass(classOf[TableOutputFormat[ImmutableBytesWritable]]) val indataRDD = sc.makeRDD(Array("4,Tom,20", "5,Jerry,21")) val rdd = indataRDD.map(_.split(',')).map { arr => val put = new Put(Bytes.toBytes(arr(0))) put.add(Bytes.toBytes("cf"), Bytes.toBytes("name"), Bytes.toBytes(arr(1))) put.add(Bytes.toBytes("cf"), Bytes.toBytes("age"), Bytes.toBytes(arr(2).toInt)) (new ImmutableBytesWritable, put) } rdd.saveAsNewAPIHadoopDataset(job.getConfiguration) sc.stop() } }

3.3 两个 API 的差异对照

维度saveAsHadoopDatasetsaveAsNewAPIHadoopDataset
OutputFormat 包org.apache.hadoop.hbase.mapredorg.apache.hadoop.hbase.mapreduce
配置载体JobConfJob + sc.hadoopConfiguration
输出 Value 类型PutPut
适用场景老代码迁移、依赖旧 mapred新项目、与 mapreduce 生态一致

提示:新项目优先用saveAsNewAPIHadoopDataset,它和 HBase 官方 mapreduce 示例一致,后续接 TableInputFormat 读取时配置风格统一,少一次心智切换。

4. 读取 HBase 转 RDD 与验证动作

4.1 newAPIHadoopRDD 读取骨架

读取用sc.newAPIHadoopRDD,InputFormat 是TableInputFormat,Key 是ImmutableBytesWritable,Value 是Result。

import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.client.Result import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.mapreduce.TableInputFormat import org.apache.hadoop.hbase.util.Bytes import org.apache.spark.{SparkConf, SparkContext} object ReadHBase { def main(args: Array[String]): Unit = { val sparkConf = new SparkConf().setAppName("HBaseRead").setMaster("local") val sc = new SparkContext(sparkConf) val tableName = "account" val conf = HBaseConfiguration.create() conf.set("hbase.zookeeper.quorum", "slave1,slave2,slave3") conf.set("hbase.zookeeper.property.clientPort", "2181") conf.set(TableInputFormat.INPUT_TABLE, tableName) val hBaseRDD = sc.newAPIHadoopRDD( conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) val count = hBaseRDD.count() println("total rows: " + count) hBaseRDD.foreach { case (_, result) => val key = Bytes.toString(result.getRow) val name = Bytes.toString(result.getValue(Bytes.toBytes("cf"), Bytes.toBytes("name"))) val age = Bytes.toInt(result.getValue(Bytes.toBytes("cf"), Bytes.toBytes("age"))) println("Rowkey:" + key + " Name:" + name + " Age:" + age) } sc.stop() } }

4.2 写入后 scan 校验

写完不要只看任务成功,去 HBase Shell 里 scan 一下:

scan 'account', {LIMIT => 5}

确认行键、列族cf、列名name/age都在,值类型对得上。如果 scan 出来是空,先查 ZooKeeper 地址是否写对,再查表名大小写。

4.3 读取后 count 校验

读取侧用hBaseRDD.count()和预期行数比对。写入 5 行就读出 5,说明链路通。如果 count 为 0 但 scan 有数据,多半是TableInputFormat.INPUT_TABLE没设或设错。

注意:getValue返回 null 时Bytes.toString会抛空指针,生产代码里建议先判空再转。

5. 本篇常见错排查

5.1 RpcRetryingCaller 反复重连

日志刷RpcRetryingCaller: Call exception但不报错,通常是 classpath 缺metrics-core或htrace-core。补齐这两个 jar 后重启任务。

5.2 connection refused

默认连localhost:2181失败。检查hbase.zookeeper.quorum是否设了真实集群地址,端口是否 2181。

5.3 类找不到或类型不匹配

老 API 的TableOutputFormat传给新 API 方法,或反过来。核对 import 是mapred还是mapreduce。

5.4 写入成功但读不到

行键类型不一致。写入用Bytes.toBytes(arr(0))字符串,读取时Bytes.toString(result.getRow)才对得上。如果写入用了toInt,读取也要按 int 解析。

5.5 依赖冲突

不同 package 下有同名类,IDE 自动导入容易导错。手动检查 import 语句,确保 OutputFormat、InputFormat 来自正确包。

6. 把链路固化下来:配置骨架与后续接入

到这里,Spark RDD 读写 HBase 的完整链路就通了:依赖补齐、ZooKeeper 显式配置、老新两个写入 API 按包路径区分、读取用newAPIHadoopRDD、写入后 scan、读取后 count。把上面三段骨架存成模板,换表名和列族就能复用。

如果你在排障或接入阶段需要统一管理模型调用和密钥,可以走 API Keys 加接入文档这条线:API Keys 在 https://taotoken.net/api-keys ,接入文档在 https://taotoken.net/doc ,两个页面配合看能少走弯路。验证模型行为时用模型对话 https://taotoken.net/models 直接试;如果是长期编码或 Agent 场景,Coding Plan https://taotoken.net/coding-plan 更合适。官网入口 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 可以留着备用。

最后留一个实用习惯:每次改完配置,先跑 count 校验再跑业务逻辑。count 对了,说明连接和表映射没问题,剩下的才是数据本身的事。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询