ARTICLE DETAIL

资讯详情

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

Spark 从 HBase 读取写入数据:TaoToken 统一 Key 通道下的配置与验证

Spark 从 HBase 读取写入数据:TaoToken 统一 Key 通道下的配置与验证 1. Spark 读写 HBase 的真实痛点为什么你的任务总是卡在连接阶段Spark 从 HBase 读取写入数据是离线数仓和实时链路里非常高频的一环。简单说它解决的是「把 HBase 这张宽表当成 Spark 的数据源或数据汇」的问题Spark 负责分布式计算HBase 负责随机读写和行级存储。适合谁做用户画像、订单明细同步、埋点数据落库、维表关联的工程师基本都会碰到。但真正动手时问题往往不在业务逻辑而在连接配置。我见过太多任务卡在RpcRetryingCaller: Call exception上日志刷屏却不报错或者直接Connection refused因为默认去连了localhost:2181。核心原因就两个一是 ZooKeeper 地址没配对二是 classpath 里缺了 HBase 依赖的 jar 包。这篇内容聚焦两条链路读取HBase → RDD和写入RDD → HBase。我会给出可复制的连接参数、表与列族配置以及读写结果的校验动作。同时在统一 Key/API 通道下把模型调用和工程配置串起来让你端到端跑通。热词里的 spark、hbase、读取、写入数据都会落到具体代码和参数上。先说结论HBase 表最好提前在 hbase shell 里建好Spark 只负责读写不要指望应用启动时自动建表权限和 region 分配容易出幺蛾子。下面按步骤拆。2. TaoToken 前置准备统一 Key 通道与依赖清单在写 Spark 代码之前先把「通道」和「依赖」两件事理清楚。所谓统一 Key 通道是指你用一套 API Key 去访问模型能力同时工程侧的配置保持一致的 Base URL 和 Model ID 规范。这样在调试 Spark 任务时如果需要用模型辅助生成配置或排查日志不用来回切换账号。第一步拿到 Key。访问 https://taotoken.net/api-keys 创建一个 API Key复制保存。注意这个 Key 只在创建时完整显示一次丢了就重新生成。第二步确认 Base URL。API 入口是 https://taotoken.net/api 所有请求都走这个地址。如果你用的是兼容 OpenAI 协议的客户端Base URL 填这个即可。第三步选模型。做代码生成和配置排查可以用模型对话页面 https://taotoken.net/models 先试一下确认通道可用。长期跑编码任务或 Agent可以看 Coding Planhttps://taotoken.net/coding-plan 。依赖清单这块是 Spark 读写 HBase 最容易翻车的地方。你需要把以下 jar 包加入 classpathlib 目录下所有hadoop开头的 jarlib 目录下所有hbase开头的 jarzookeeper-3.4.6.jarmetrics-core-2.2.0.jar缺了会一直重连不报错htrace-core-3.1.0-incubating.jarguava-12.0.1.jar$SPARK_HOME/lib下的spark-assembly-1.6.1-hadoop2.4.0.jar不同 package 里可能有同名类导入时别导错。比如TableOutputFormat在org.apache.hadoop.hbase.mapred和org.apache.hadoop.hbase.mapreduce下都有用错方法就写不进去。连接 ZooKeeper 有两种方式一是把hbase-site.xml放进 classpath二是在HBaseConfiguration实例里手动 set。我建议手动设置因为集群地址变了不用重新打包。不设置的话默认连localhost:2181直接Connection refused。3. 可复制配置Spark 连接 HBase 的参数与代码片段这一节给可直接复制的配置。先看连接参数这是所有读写的基础。val conf HBaseConfiguration.create() conf.set(hbase.zookeeper.quorum, slave1,slave2,slave3) conf.set(hbase.zookeeper.property.clientPort, 2181)如果你用sc.hadoopConfiguration写法是sc.hadoopConfiguration.set(hbase.zookeeper.quorum, slave1,slave2,slave3) sc.hadoopConfiguration.set(hbase.zookeeper.property.clientPort, 2181)表名和列族要提前在 hbase shell 建好create account, cf写入用saveAsHadoopDataset旧 APIorg.apache.hadoop.hbase.mapred.TableOutputFormatimport 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} import org.apache.spark.rdd.RDD.rddToPairRDDFunctions object TestHBase { def main(args: Array[String]): Unit { val sparkConf new SparkConf().setAppName(HBaseTest).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).toInt)) 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() } }写入用saveAsNewAPIHadoopDataset新 APIorg.apache.hadoop.hbase.mapreduce.TableOutputFormatimport org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.mapreduce.TableOutputFormat import org.apache.spark._ import org.apache.hadoop.mapreduce.Job import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.client.{Result, Put} import org.apache.hadoop.hbase.util.Bytes object TestHBase3 { def main(args: Array[String]): Unit { val sparkConf new SparkConf().setAppName(HBaseTest).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[Result]) job.setOutputFormatClass(classOf[TableOutputFormat[ImmutableBytesWritable]]) 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.saveAsNewAPIHadoopDataset(job.getConfiguration()) } }读取用newAPIHadoopRDDimport org.apache.hadoop.hbase.{HBaseConfiguration, HTableDescriptor, TableName} import org.apache.hadoop.hbase.client.HBaseAdmin import org.apache.hadoop.hbase.mapreduce.TableInputFormat import org.apache.spark._ import org.apache.hadoop.hbase.util.Bytes object TestHBase2 { def main(args: Array[String]): Unit { val sparkConf new SparkConf().setAppName(HBaseTest).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 admin new HBaseAdmin(conf) if (!admin.isTableAvailable(tablename)) { val tableDesc new HTableDescriptor(TableName.valueOf(tablename)) admin.createTable(tableDesc) } val hBaseRDD sc.newAPIHadoopRDD(conf, classOf[TableInputFormat], classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable], classOf[org.apache.hadoop.hbase.client.Result]) val count hBaseRDD.count() println(count) hBaseRDD.foreach { case (_, result) val key Bytes.toString(result.getRow) val name Bytes.toString(result.getValue(cf.getBytes, name.getBytes)) val age Bytes.toInt(result.getValue(cf.getBytes, age.getBytes)) println(Row key: key Name: name Age: age) } sc.stop() admin.close() } }如果你用 Cline MCP 或 Codex 的auth.json来管理模型通道配置三件套要写全Base URL 填https://taotoken.net/apiKey 填你创建的 API KeyModel ID 按模型对话页面里显示的填。CC Switch 同理切换配置时别只改 Key 忘了 Base URL。4. 验证请求与成功结果读写结果校验动作代码跑起来只是第一步关键是验证数据真的进去了、真的读出来了。写入之后去 hbase shell 里查scan account你应该看到三行行键 1、2、3列族 cf 下有 name 和 age。如果 scan 为空说明写入没成功先看日志有没有TableOutputFormat相关的报错。读取验证看控制台输出。hBaseRDD.count()应该打印 3然后 foreach 打印Row key:1 Name:jack Age:15 Row key:2 Name:Lily Age:16 Row key:3 Name:mike Age:16如果 count 是 0检查TableInputFormat.INPUT_TABLE是否设对表名大小写敏感。如果 count 对但字段是 null检查列族名和列名是否和建表时一致。模型通道的验证可以用模型对话页面发一条测试请求确认返回正常。这一步是为了排除 Key 或 Base URL 的问题别把模型通道的报错和 HBase 的报错混在一起排查。实测下来读写都通之后建议把setMaster(local)换成集群模式再跑一遍因为 local 模式下 classpath 和集群模式可能不一致本地能跑不代表集群能跑。5. 本篇常见错排查401、local proxy failed、reading choices、OAuth这一节对照真实报错逐个拆。401 Unauthorized模型通道的 Key 不对或过期。检查https://taotoken.net/api-keys里的 Key 是否复制完整Base URL 是否是https://taotoken.net/api。如果用的是auth.json确认字段名和格式没写错。local proxy failed本地代理配置问题。如果你在客户端里配了代理地址但代理没启动就会报这个。检查客户端设置或者临时关掉代理直连。注意这里说的是客户端自身的网络设置不是让你去搞什么特殊网络工具。reading choices 报错通常是模型返回格式和客户端预期不一致。检查 Model ID 是否填对有些客户端要求模型名和通道支持的名称完全匹配。去模型对话页面确认当前可用的模型名。OAuth 相关报错如果你用 Claude Code 或 Anthropic 风格的客户端OAuth 流程没走完或 token 过期。重新走一遍授权或者改用 API Key 方式。Claude Code 的接入文档在 https://taotoken.net/doc 里面有详细步骤。HBase 侧报错RpcRetryingCaller: Call exception不断重连但不报错缺metrics-core-2.2.0.jar补上。Connection refusedZooKeeper 地址没设默认连了 localhost。按第 3 节的参数设置。ClassNotFoundExceptionjar 包没加全对照第 2 节的依赖清单逐个检查。NoSuchMethodError同名类导错包比如TableOutputFormat导成了 mapred 但代码用的是 mapreduce。6. 语义一致 CTA把通道和工程配置固定下来跑通一次之后把配置固定成模板。Spark 侧把 ZooKeeper 地址、端口、表名抽成配置文件模型侧把 Base URL、Key、Model ID 写进auth.json或客户端配置。这样下次新任务直接复用不用重新踩坑。需要长期跑编码任务或 Agent 的看 Coding Planhttps://taotoken.net/coding-plan 。只是偶尔验证模型或排查配置用模型对话页面就够https://taotoken.net/models 。接入文档和报错排查看 https://taotoken.net/doc API Key 管理在 https://taotoken.net/api-keys 。最后提醒一句HBase 表提前建好列族名和代码里保持一致classpath 里的 jar 包一个都别少。这三件事做到Spark 读写 HBase 基本不会卡住。
返回列表