新闻详情

新闻详情

首页 / 资讯中心 / 详情

Spark 读写 HBase 实战:saveAsHadoopDataset 与 saveAsNewAPIHadoopDataset 配置骨架

发布时间:2026/9/29 3:48:50来源:尧图网络
Spark 读写 HBase 实战:saveAsHadoopDataset 与 saveAsNewAPIHadoopDataset 配置骨架
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 里需要补齐这些 jarHBase 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 的差异对照维度saveAsHadoopDatasetsaveAsNewAPIHadoopDatasetOutputFormat 包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.newAPIHadoopRDDInputFormat 是TableInputFormatKey 是ImmutableBytesWritableValue 是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_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 可以留着备用。最后留一个实用习惯每次改完配置先跑 count 校验再跑业务逻辑。count 对了说明连接和表映射没问题剩下的才是数据本身的事。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

MTK平台充电调试全解析:从硬件通路到快充协议 2026/9/29 4:40:58

MTK平台充电调试全解析:从硬件通路到快充协议

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
Windows蓝屏0x000007E深度解析:从SYSTEM_THREAD_EXCEPTION_NOT_HANDLED到根因定位 2026/9/29 4:40:58

Windows蓝屏0x000007E深度解析:从SYSTEM_THREAD_EXCEPTION_NOT_HANDLED到根因定位

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
TCRT5000+STM32CubeMX循迹小车实战指南 2026/9/29 4:40:58

TCRT5000+STM32CubeMX循迹小车实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
Computer-Use 与 Browser-Agent 实战:用 Playwright + Accessibility Tree 搭一套可复现的 WebArena 评测骨架 2026/9/29 4:40:51

Computer-Use 与 Browser-Agent 实战:用 Playwright + Accessibility Tree 搭一套可复现的 WebArena 评测骨架

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
【电脑智能操控神器】小龙虾 OpenClaw 适配 Win11 完整教学:TaoToken 统一 Key 接入与 config.toml 配置骨架 2026/9/29 4:40:51

【电脑智能操控神器】小龙虾 OpenClaw 适配 Win11 完整教学:TaoToken 统一 Key 接入与 config.toml 配置骨架

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
OpenClaw限流有救了!免费Nvidia API+阿里云百炼接入指南(TaoToken统一Key配置版) 2026/9/29 4:40:51

OpenClaw限流有救了!免费Nvidia API+阿里云百炼接入指南(TaoToken统一Key配置版)

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉