☰
Spark 读写 HBase 与 SparkStreaming 操作:TaoToken 统一 Key 接入配置骨架
2026/10/3 6:45:47 网站建设 项目流程

1. Spark 读写 HBase 与 SparkStreaming 操作:TaoToken 统一 Key 接入配置骨架

Spark 读写 HBase 这件事,说难不难,说简单也容易踩坑。批处理侧用newAPIHadoopRDD读、saveAsNewAPIHadoopDataset写,流处理侧用 SparkStreaming 消费 Kafka 再落 HBase,这套组合在实时数仓里非常常见。真正让人头疼的不是 API 本身,而是连接配置散落在各处、ZooKeeper 参数写错一个字母就报local proxy failed、Kafka offset 管理不当导致重复消费。这篇内容聚焦 Spark 批处理读写 HBase 与 SparkStreaming 流式写入场景,结合 Kafka 数据源与 Phoenix 查询视角,交付一套可复制的连接配置骨架,同时把 TaoToken 统一 Key/API 通道的 settings.json 与 config.toml 片段给出来,让你一次性跑通批流两套 HBase 操作。

适合谁看:正在做 Spark + HBase 集成的数据开发、需要把 Kafka 数据实时写入 HBase 的工程师、以及想用 Phoenix 做 HBase 查询但不想手写大量 JDBC 代码的同学。核心检索词就是 Spark 读写 HBase、SparkStreaming 操作、Kafka 数据源、Phoenix 查询,这几个词会贯穿全文。

我试过在同一个项目里同时维护批处理和流处理两套 HBase 连接配置,结果因为 ZooKeeper 地址不一致排查了大半天。后来把配置抽成统一骨架,问题就少了很多。下面从原问题与场景开始,一步步把配置和验证动作讲清楚。

2. TaoToken 前置:统一 Key 与 API 通道准备

在动手写 Spark 代码之前,先把 TaoToken 的统一 Key 和 API 通道准备好。TaoToken 在这里的角色是给模型对话、coding-plan、console 等场景提供统一的 API 入口,方便你在 Spark 项目里做配置管理时有一个集中的 Key 来源。官网地址是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 入口是 https://taotoken.net/api ,注意 API 地址不加 UTM 参数。

你需要先拿到一个可用的 API Key。进入 console 页面创建 Key,地址是 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console_key&utm_campaign=rewrite 。创建完成后,Key 只在创建时显示一次,记得复制保存。如果你后续要用 Claude Code 做代码润色或辅助生成 Spark 配置,可以走 ClaudeCodeAnthropic 通道,地址是 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_content=claudecode&utm_campaign=rewrite 。

为什么 Spark 项目里要提 TaoToken?因为很多同学在写 Spark 作业时,会顺手让模型帮忙生成配置片段或排查报错。这时候如果 API Key 散落在各个工具里,管理起来很乱。统一用 TaoToken 的 Key,配合 settings.json 和 config.toml 做集中配置,后面换 Key 或加通道都方便。API Keys 管理页在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=apikeys&utm_campaign=rewrite ,接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 。

这里要强调一点:TaoToken 是合法的 API 通道服务,不是所谓的中转。你在配置里填的 Base URL 就是 https://taotoken.net/api ,Key 就是 console 里创建的那串字符。模型对话入口在 https://taotoken.net/chat?utm_source=taotoken_aicg_blog_end&utm_content=chat&utm_campaign=rewrite ,可以用来验证 Key 是否可用。Coding Plan 入口在 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=codingplan&utm_campaign=rewrite ,适合长期做 Spark 作业开发的同学。

准备好 Key 之后,我们进入可复制配置环节。下面会给出 settings.json 和 config.toml 的完整片段,以及 Spark 侧 HBase 连接配置骨架。这些配置你可以直接复制到项目里,改掉 ZooKeeper 地址和表名就能用。

3. 可复制配置:settings.json、config.toml 与 Spark HBase 连接骨架

这一节是全文的核心,直接给可复制的配置片段。先看 TaoToken 侧的 settings.json,路径放在项目根目录的.taotoken/settings.json下,内容如下:

{ "api_base": "https://taotoken.net/api", "api_key": "sk-你的TaoTokenKey", "default_model": "claude-sonnet-4-20250514", "timeout_seconds": 60, "retry": { "max_attempts": 3, "backoff_ms": 800 }, "channels": { "chat": "https://taotoken.net/chat", "coding_plan": "https://taotoken.net/coding-plan", "console": "https://taotoken.net/console", "api_keys": "https://taotoken.net/api-keys", "doc": "https://taotoken.net/doc" } }

注意 api_key 字段填你在 console 创建的那串 Key,不要带多余空格。default_model 按你实际使用的模型 ID 填,这里只是示例。timeout_seconds 和 retry 按网络情况调整,内网环境可以适当调大。

再看 config.toml,路径放在~/.taotoken/config.toml,适合命令行工具读取:

[api] base_url = "https://taotoken.net/api" api_key = "sk-你的TaoTokenKey" timeout = 60 [model] default = "claude-sonnet-4-20250514" max_tokens = 4096 [retry] max_attempts = 3 backoff_ms = 800 [links] chat = "https://taotoken.net/chat" coding_plan = "https://taotoken.net/coding-plan" console = "https://taotoken.net/console" api_keys = "https://taotoken.net/api-keys" doc = "https://taotoken.net/doc"

这两个文件的作用是给模型辅助工具提供统一入口。接下来是 Spark 侧 HBase 连接配置骨架。我把它抽成一个 Scala 对象,批处理和流处理共用:

import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.mapreduce.TableInputFormat import org.apache.hadoop.hbase.mapreduce.TableOutputFormat import org.apache.hadoop.conf.Configuration object HBaseConfBuilder { def build(zkQuorum: String, zkPort: String, tableName: String): Configuration = { val conf = HBaseConfiguration.create() conf.set("hbase.zookeeper.quorum", zkQuorum) conf.set("hbase.zookeeper.property.clientPort", zkPort) conf.set("hbase.defaults.for.version.skip", "true") conf.set("zookeeper.znode.parent", "/hbase") conf.set(TableInputFormat.INPUT_TABLE, tableName) conf.set(TableOutputFormat.OUTPUT_TABLE, tableName) conf } }

这个骨架里,zkQuorum 填你的 ZooKeeper 集群地址,多个节点用逗号分隔。zkPort 默认 2181。tableName 按实际表名填。hbase.defaults.for.version.skip设为 true 可以跳过版本检查,避免一些兼容性报错。zookeeper.znode.parent默认是/hbase,如果你的 HBase 集群改过这个路径,要对应调整。

批处理读取用newAPIHadoopRDD,写入用saveAsNewAPIHadoopDataset。流处理侧在 foreachRDD 里拿到 offsetRanges 后,用foreachPartition建 HTable 连接,关闭自动提交并设置写缓冲区,最后 flushCommits。这些代码在后面的验证环节会给出完整可运行版本。

Phoenix 侧如果用 DataFrame 读写,依赖要加 phoenix-core 和 phoenix-spark,JDBC URL 格式是jdbc:phoenix:zkHost:2181。读的时候可以用spark.read.format("org.apache.phoenix.spark"),写的时候用df.write.format("org.apache.phoenix.spark").mode(SaveMode.Overwrite)。注意写入前表要提前建好,Phoenix 不会自动建表。

配置骨架给完了,下面进入验证请求环节。我会给出批处理和流处理两套验证动作,确保你跑通。

4. 验证请求与成功结果:批流两套 HBase 操作跑通

先验证批处理读取。用newAPIHadoopRDD读 HBase,代码骨架如下:

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.sql.SparkSession object ReadHBaseBatch { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("ReadHBaseBatch") .master("local[4]") .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .getOrCreate() val sc = spark.sparkContext val conf = HBaseConfBuilder.build("zk1,zk2,zk3", "2181", "student") val hbaseRDD = sc.newAPIHadoopRDD( conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) hbaseRDD.foreach { case (_, result) => val rowKey = Bytes.toString(result.getRow) val name = Bytes.toString(result.getValue("info".getBytes, "name".getBytes)) val age = Bytes.toString(result.getValue("info".getBytes, "age".getBytes)) println(s"RowKey=$rowKey Name=$name Age=$age") } println("总记录数: " + hbaseRDD.count()) spark.stop() } }

跑通后控制台会打印每行的 rowKey、name、age,最后输出总记录数。如果 count 为 0,先检查表里有没有数据,再检查TableInputFormat.INPUT_TABLE是否设对。

再验证批处理写入。用saveAsNewAPIHadoopDataset写 HBase:

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.sql.SparkSession object WriteHBaseBatch { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("WriteHBaseBatch").getOrCreate() val sc = spark.sparkContext val conf = HBaseConfBuilder.build("zk1,zk2,zk3", "2181", "student") val job = Job.getInstance(conf) job.setOutputKeyClass(classOf[ImmutableBytesWritable]) job.setOutputValueClass(classOf[Put]) job.setOutputFormatClass(classOf[TableOutputFormat[ImmutableBytesWritable]]) val data = sc.makeRDD(Array("1001,Alice,F,22", "1002,Bob,M,25")) val rdd = data.map(_.split(',')).map { arr => val put = new Put(Bytes.toBytes(arr(0))) put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("name"), Bytes.toBytes(arr(1))) put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("gender"), Bytes.toBytes(arr(2))) put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("age"), Bytes.toBytes(arr(3))) (new ImmutableBytesWritable, put) } rdd.saveAsNewAPIHadoopDataset(job.getConfiguration()) println("写入完成") spark.stop() } }

写入完成后,用 HBase shell 执行scan 'student'应该能看到 1001 和 1002 两行。注意job.setOutputValueClass这里设的是classOf[Put],有些老版本示例写的是classOf[Result],那是错的,写入必须用 Put。

流处理侧验证 SparkStreaming 消费 Kafka 写 HBase。核心是手动维护 offset,关闭自动提交:

val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "kafka1:9092,kafka2:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "spark_hbase_group", "auto.offset.reset" -> "latest", "enable.auto.commit" -> (false: java.lang.Boolean) ) val stream = KafkaUtils.createDirectStream[String, String]( scc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](Set("topic_hbase"), kafkaParams) ) stream.foreachRDD { rdd => if (!rdd.isEmpty()) { val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges rdd.foreachPartition { partition => val conf = HBaseConfBuilder.build("zk1,zk2,zk3", "2181", "test") val table = new HTable(conf, TableName.valueOf("test")) table.setAutoFlush(false, false) table.setWriteBufferSize(3 * 1024 * 1024) partition.foreach { record => val json = JSONObject.fromObject(record.value()) val put = new Put(Bytes.toBytes(json.get("rowkey").toString)) put.addColumn(Bytes.toBytes("f1"), Bytes.toBytes("data"), Bytes.toBytes(json.get("data").toString)) table.put(put) } table.flushCommits() table.close() } offsetRanges.foreach { o => println(s"topic=${o.topic} partition=${o.partition} from=${o.fromOffset} until=${o.untilOffset}") } } }

跑通后控制台会打印每个 partition 的 offset 范围,HBase 里能看到新写入的数据。offset 可以存到 Redis 或 ZooKeeper,下次启动时从上次位置继续消费,实现 exactly-once 语义。

Phoenix 侧验证用 DataFrame 读:

val df = spark.read .format("org.apache.phoenix.spark") .options(Map("table" -> "PHOENIXTEST", "zkUrl" -> "jdbc:phoenix:zk1:2181")) .load() df.printSchema() df.show(10)

成功的话会打印 schema 和前 10 行数据。写入用df.write.format("org.apache.phoenix.spark").mode(SaveMode.Overwrite).options(Map("table" -> "PHOENIXTESTCOPY", "zkUrl" -> "jdbc:phoenix:zk1:2181")).save(),注意目标表要提前建好。

验证动作给完了,下面进入常见错排查。这些报错都是我实际遇到过的,对照着看能省不少时间。

5. 本篇常见错排查:401、local proxy failed、reading choices、OAuth

第一个高频报错是 401 Unauthorized。如果你在调用 TaoToken API 时看到 401,先检查 settings.json 或 config.toml 里的 api_key 是否填对。常见问题是 Key 前后有空格、Key 被截断、或者用了过期的 Key。去 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=apikeys_fix&utm_campaign=rewrite 重新生成一个,替换后重启工具。另外确认 Base URL 是 https://taotoken.net/api ,不要多加斜杠或路径。

第二个报错是local proxy failed。这个在 Spark 连 HBase 时经常出现,原因是 ZooKeeper 地址或端口配错。检查hbase.zookeeper.quorum是否填了正确的 ZK 节点,hbase.zookeeper.property.clientPort是否是 2181。如果 ZK 集群开了 ACL 或 Kerberos,还要额外配置认证。另一个可能是zookeeper.znode.parent路径不对,默认是/hbase,有些集群会改成/hbase-secure之类。用echo stat | nc zkHost 2181可以确认 ZK 是否可达。

第三个报错是reading choices相关,通常出现在 Kafka 消费端。报错信息里会有Error reading field 'choices'或类似字样,这多半是反序列化器配错。检查key.deserializer和value.deserializer是否都设成了classOf[StringDeserializer]。如果 Kafka 里的消息不是字符串格式,要用对应的反序列化器。另外auto.offset.reset设成latest时,如果 topic 没有新数据,流会一直空转,可以临时改成earliest验证。

第四个是 OAuth 相关报错。如果你用 Claude Code 或类似工具接入 TaoToken,看到 OAuth token 失效或回调失败,先确认 ClaudeCodeAnthropic 通道的配置是否正确。地址是 https://taotoken.net/claude-code-anthropic?utm_source=taotoken_aicg_blog_end&utm_content=claudecode_fix&utm_campaign=rewrite 。OAuth 流程需要回调地址可达,本地开发时注意端口不要被占用。如果反复失败,改用 API Key 方式接入更稳定。

还有一个容易忽略的报错是 HBase 写入时的TableNotFoundException。这通常是因为表名大小写不一致,或者 namespace 没带。比如cbd:prod_base这种带 namespace 的表,配置里要写全。另外 Phoenix 写入时如果目标表不存在,也会报错,Phoenix 不会自动建表,要先用CREATE TABLE建好。

排查完这些,基本就能跑通批流两套 HBase 操作了。最后给一个语义一致的 CTA,方便你继续深入。

6. 继续深入:模型对话、Coding Plan 与接入文档

批流两套 HBase 操作跑通之后,如果你想继续优化配置或排查更复杂的报错,可以用 TaoToken 的模型对话入口快速验证想法。地址是 https://taotoken.net/chat?utm_source=taotoken_aicg_blog_end&utm_content=chat_cta&utm_campaign=rewrite ,把报错信息贴进去,让模型帮你分析。长期做 Spark 作业开发的同学,可以看 Coding Plan,地址是 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=codingplan_cta&utm_campaign=rewrite ,适合需要频繁生成和调试代码的场景。

接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc_cta&utm_campaign=rewrite ,里面有完整的 API 说明和配置示例。API Keys 管理在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=apikeys_cta&utm_campaign=rewrite ,需要换 Key 或加通道时来这里。官网首页 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= 可以看整体介绍。

最后分享一个实用技巧:把 HBaseConfBuilder 抽成独立对象后,批处理和流处理共用同一份 ZooKeeper 配置,改一处就全生效。另外 SparkStreaming 写 HBase 时,setAutoFlush(false, false)和setWriteBufferSize能显著提升写入吞吐,但记得在 partition 结束时调flushCommits(),否则数据可能丢。Phoenix 查询适合做聚合分析,但写入性能不如原生 HBase API,按场景选。

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

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

立即咨询