ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Spark 读写 HBase 的两种方式(RDD、DataFrame)配 TaoToken:settings.json 骨架与连通性验证

Spark 读写 HBase 的两种方式(RDD、DataFrame)配 TaoToken:settings.json 骨架与连通性验证 1. 为什么 Spark 读写 HBase 总在环境这关卡住Spark 和 HBase 放在一起用是很多数据管道里绕不开的组合HBase 负责海量行的随机读写Spark 负责批量计算和 SQL 分析。但真正动手时卡住大多数人的往往不是业务逻辑而是环境本身——ZooKeeper 地址写错、依赖包版本对不上、TableInputFormat和TableOutputFormat分属两个包、Phoenix 的 JDBC URL 拼错一个字符就连不上。这些问题的共同点是报错信息模糊排查方向发散最后变成反复改配置、反复重启。这篇聚焦两件事。第一先把 Spark 侧访问外部服务时统一走的那条通道配置好用一份可复制的settings.json骨架把 Key 和 API 地址固定下来避免每个任务里散落硬编码。第二分别用 RDD 和 DataFrame 两条链路各跑一次真实的读写验证确认 HBase 环境可用并对比两种方式各自适合什么场景。RDD 链路走的是 HBase 原生mapred/mapreduce接口DataFrame 链路走 Phoenix 的 JDBC 与phoenix-spark数据源。两条链路都验证通过你才算真正把环境跑通了。适合已经装好 Spark 和 HBase、准备写第一个读写任务的人也适合之前跑通过但换机器后配置丢失、需要快速重建的人。下面所有配置和代码都可以直接复制改掉地址和表名就能用。2. TaoToken 前置settings.json 骨架与 Key 获取在写 Spark 任务之前先把访问通道的配置统一收口。很多团队的做法是在每个spark-submit里用--conf传一堆参数或者把 Key 写死在代码里时间一长没人记得哪个 Key 对应哪个环境。更稳的做法是准备一份settings.json把 API 地址和 Key 放在一起代码里只读这个文件。先到控制台创建 Key入口在 API Keys 页面https://taotoken.net/console/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi_keysutm_campaignrewrite创建后拿到形如sk-开头的字符串只显示一次复制保存。API 的基础地址是https://taotoken.net/api注意这个地址不带任何查询参数直接作为 base URL 使用。下面是一份settings.json骨架放在项目src/main/resources下打包时会进 classpath{ taotoken: { baseUrl: https://taotoken.net/api, apiKey: sk-替换成你自己的Key, timeoutMs: 30000, retry: 2 }, hbase: { zkQuorum: localhost, zkPort: 2181, tableName: SparkHBase, columnFamily: cf1 }, phoenix: { jdbcUrl: jdbc:phoenix:localhost:2181, table: PHOENIXTEST } }读取这份配置的 Scala 工具类可以这样写用 Jackson 解析避免手写字符串拼接import com.fasterxml.jackson.databind.ObjectMapper import scala.io.Source object Settings { private val mapper new ObjectMapper() lazy val root: Map[String, Any] { val stream getClass.getResourceAsStream(/settings.json) val json Source.fromInputStream(stream, UTF-8).mkString mapper.readValue(json, classOf[Map[String, Any]]) } def taotoken(key: String): String root(taotoken).asInstanceOf[Map[String, String]](key) def hbase(key: String): String root(hbase).asInstanceOf[Map[String, String]](key) }这样代码里取 ZooKeeper 地址就是Settings.hbase(zkQuorum)换环境只改 JSON不动代码。Key 的获取和模型调用验证可以顺手在模型对话页做一次确认通道本身是通的https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_contentmodelsutm_campaignrewrite注意settings.json里如果放了真实 Key不要提交到公开仓库。生产环境建议用环境变量覆盖代码里先读环境变量、读不到再回退到 JSON。3. RDD 链路saveAsHadoopDataset 写入与 newAPIHadoopRDD 读取RDD 链路是 HBase 最原生的对接方式不依赖 Phoenix直接操作Put和Result对象。它的优点是控制粒度细能精确指定行键、列族、列名缺点是代码偏底层字段多了会写得很啰嗦。3.1 依赖与包路径的坑先确认pom.xml里 HBase 相关依赖的版本和集群一致。最容易踩的坑是TableOutputFormat和TableInputFormat分属两个包用途类所在包写入TableOutputFormatorg.apache.hadoop.hbase.mapred读取TableInputFormatorg.apache.hadoop.hbase.mapreduce写入用的是mapred老接口读取用的是mapreduce新接口。如果 import 写反了编译能过但运行时报ClassNotFoundException或格式不匹配。依赖大致如下dependency groupIdorg.apache.hbase/groupId artifactIdhbase-client/artifactId version2.4.11/version /dependency dependency groupIdorg.apache.hbase/groupId artifactIdhbase-mapreduce/artifactId version2.4.11/version /dependency3.2 写入saveAsHadoopDataset写入的核心是把每行数据构造成(ImmutableBytesWritable, Put)的二元组再交给saveAsHadoopDataset。行键和所有列值都必须用Bytes.toBytes转换直接传 String 会报类型错误。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.sql.SparkSession object SparkHBaseRDDWrite { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(SparkHBaseRDDWrite) .getOrCreate() val sc spark.sparkContext val tableName Settings.hbase(tableName) val cf Settings.hbase(columnFamily) val hbaseConf HBaseConfiguration.create() hbaseConf.set(hbase.zookeeper.quorum, Settings.hbase(zkQuorum)) hbaseConf.set(hbase.zookeeper.property.clientPort, Settings.hbase(zkPort)) hbaseConf.set(TableOutputFormat.OUTPUT_TABLE, tableName) val jobConf new JobConf(hbaseConf) jobConf.setOutputFormat(classOf[TableOutputFormat]) val raw sc.makeRDD(Array(2,jack,16, 1,Lucy,15, 5,mike,17, 3,Lily,14)) val rdd raw.map(_.split(,)).map { arr val put new Put(Bytes.toBytes(arr(0))) put.addColumn(Bytes.toBytes(cf), Bytes.toBytes(name), Bytes.toBytes(arr(1))) put.addColumn(Bytes.toBytes(cf), Bytes.toBytes(age), Bytes.toBytes(arr(2))) (new ImmutableBytesWritable, put) } rdd.saveAsHadoopDataset(jobConf) spark.stop() } }跑完后到 HBase shell 里scan SparkHBase能看到 4 行记录每行两个列说明写入链路通了。3.3 读取newAPIHadoopRDD读取用sc.newAPIHadoopRDD传入TableInputFormat和Result类型返回的 RDD 每个元素是(ImmutableBytesWritable, Result)。取值时通过列族和列名定位同样用Bytes.toString还原。import org.apache.hadoop.hbase.mapreduce.TableInputFormat import org.apache.hadoop.hbase.client.Result val hbaseConf HBaseConfiguration.create() hbaseConf.set(hbase.zookeeper.quorum, Settings.hbase(zkQuorum)) hbaseConf.set(hbase.zookeeper.property.clientPort, Settings.hbase(zkPort)) hbaseConf.set(TableInputFormat.INPUT_TABLE, Settings.hbase(tableName)) val hBaseRDD sc.newAPIHadoopRDD( hbaseConf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) hBaseRDD.foreach { case (_, result) val key Bytes.toString(result.getRow) val name Bytes.toString(result.getValue(Bytes.toBytes(cf1), Bytes.toBytes(name))) val age Bytes.toString(result.getValue(Bytes.toBytes(cf1), Bytes.toBytes(age))) println(sRow key:$key\tname:$name\tage:$age) }输出 4 行行键和列值都对得上读取链路就验证完了。RDD 方式适合字段少、需要精确控制行键、或者要复用已有 HBase 原生逻辑的场景。4. DataFrame 链路通过 Phoenix 读写 HBaseDataFrame 链路走 Phoenix好处是能用 SQL 和 DataFrame API字段多的时候不用一个个Put。代价是多一层 Phoenix 依赖且表结构要先在 Phoenix 里建好。4.1 依赖与 JDBC 读取pom.xml里加 Phoenix 的两个包版本要和 Phoenix 服务端一致dependency groupIdorg.apache.phoenix/groupId artifactIdphoenix-core/artifactId version5.1.2/version /dependency dependency groupIdorg.apache.phoenix/groupId artifactIdphoenix-spark/artifactId version5.1.2/version /dependency读取有两种写法。第一种走标准 JDBC指定 Phoenix 驱动val url Settings.root(phoenix).asInstanceOf[Map[String, String]](jdbcUrl) val dbtable Settings.root(phoenix).asInstanceOf[Map[String, String]](table) val rdf spark.read .format(jdbc) .option(driver, org.apache.phoenix.jdbc.PhoenixDriver) .option(url, url) .option(dbtable, dbtable) .load() rdf.printSchema() rdf.show()第二种走phoenix-spark数据源用table和zkUrl两个 optionval df spark.read .format(org.apache.phoenix.spark) .options(Map(table - dbtable, zkUrl - url)) .load() df.printSchema()两种方式返回的 DataFrame 结构一致printSchema能看到PK和COL1两列。4.2 写入与结果核对写入用phoenix-spark数据源注意目标表必须先在 Phoenix 里建好否则会报表不存在import org.apache.spark.sql.SaveMode df.write .format(org.apache.phoenix.spark) .mode(SaveMode.Overwrite) .options(Map(table - PHOENIXTESTCOPY, zkUrl - url)) .save()到 Phoenix 命令行里核对SELECT * FROM PHOENIXTESTCOPY;能看到和源表一样的 4 行数据说明 DataFrame 写入也通了。DataFrame 方式适合字段多、需要做聚合和 join、或者团队更习惯 SQL 的场景。它和 RDD 方式不是替代关系而是按任务特点选随机点查和精细控制选 RDD批量分析和宽表选 DataFrame。5. 本篇常见错排查环境跑不通时九成问题集中在下面几类按顺序排查能省很多时间。ZooKeeper 连不上。报错通常是Connection refused或超时。先确认hbase.zookeeper.quorum和端口和集群实际一致再确认 Spark 运行节点能访问该地址。如果 HBase 是伪分布式localhost只在同一台机器有效跨节点要换成真实主机名或 IP。TableInputFormat 找不到。编译通过但运行报ClassNotFoundException多半是hbase-mapreduce依赖没打进包或者spark-submit时没带--jars。用mvn package打 fat jar或者显式指定依赖路径。写入报列族不存在。Put.addColumn里的列族必须已经在 HBase 表里创建。建表命令是create SparkHBase, cf1列族名要和代码里Bytes.toBytes(cf1)完全一致大小写敏感。Phoenix JDBC 连不上。检查jdbc:phoenix:后面的地址和端口Phoenix 默认走 ZooKeeper 的 2181。如果 Phoenix 服务没起JDBC 会一直重试直到超时。另外phoenix-core版本和服务端不一致时会报协议不匹配。DataFrame 写入报表不存在。phoenix-spark的save不会自动建表必须先在 Phoenix 里CREATE TABLE。表名大小写在 Phoenix 里默认转大写代码里写小写可能匹配不上。Key 读取为 null。settings.json没进 classpath或者字段名拼错。用getClass.getResourceAsStream(/settings.json)确认能读到读不到就检查src/main/resources目录和打包配置。提示排查时先把日志级别调到 WARN 以上HBase 客户端默认会打大量 INFO 日志容易淹没真正的报错。Logger.getLogger(org.apache.hadoop.hbase).setLevel(Level.WARN)能压掉大部分噪音。6. 通道与接入文档环境验证通过后把 Key 和接入方式固定下来后续换机器或加任务时直接复用。API Keys 管理页可以创建和吊销 Keyhttps://taotoken.net/console/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentapi_keysutm_campaignrewrite接入参数和调用示例在文档里包含 base URL、鉴权头和常见错误码说明https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_contentdocutm_campaignrewrite如果后续要把这套通道用在长期运行的编码任务或 Agent 流程里Coding Plan 页面有对应的额度方案https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentcoding_planutm_campaignrewrite我自己的习惯是settings.json只放非敏感配置Key 走环境变量注入代码里做一次回退读取。这样本地调试和集群提交用同一份代码只换环境变量不用改任何一行逻辑。RDD 和 DataFrame 两条链路都验证过之后再往上叠业务逻辑心里就有底了。
返回列表