Spark读写HBase操作(scala版):TaoToken统一Key下DataFrame与SparkSQL落地

发布时间:2026/10/8 22:00:09
Spark读写HBase操作(scala版):TaoToken统一Key下DataFrame与SparkSQL落地
1. Spark 读写 HBase 到底难在哪scala 版 DataFrame 落地场景拆解Spark 读写 HBase 这件事说简单也简单说坑多也是真的多。简单在于 HBase 官方提供了TableInputFormat和TableOutputFormatSpark 只要通过newAPIHadoopRDD和saveAsNewAPIHadoopDataset就能对接坑多在于返回的是RDD[(ImmutableBytesWritable, Result)]而业务侧想要的是带 schema 的 DataFrame中间这层转换、列类型映射、RowKey 处理每一步都能卡住人。这篇聚焦的场景很具体你用 scala 写 Spark 作业数据源是 HBase目标是把 HBase 表读成 DataFrame 做 SparkSQL 分析或者把 DataFrame 写回 HBase。同时作业里往往还要调用大模型能力做字段补全、语义标注、异常描述生成这时候访问凭据的管理就成了另一个隐形成本——每个作业里硬编码一堆 Key换环境就炸。所以我会把 TaoToken 统一 Key 这条线也串进来让 HBase 读写和模型调用共用一套凭据通道。适合谁看已经会写基础 Spark scala 作业、手上有一张 HBase 表、想把它接进 DataFrame/SparkSQL 管线的同学。如果你还没跑通过newAPIHadoopRDD跟着下面的步骤一步步来就行。先说清楚核心检索词Spark 通过 DataFrame 读写 HBase本质是「Hadoop InputFormat/OutputFormat RDD 转换 schema 映射」三件事。读的时候TableInputFormat把每个 row 变成一个Result对象你要自己决定 rowKey 放哪列、family:qualifier 怎么拍平成列名、字节数组怎么转成 String/Int/Long。写的时候反过来TableOutputFormat要求你提供(ImmutableBytesWritable, Put)的 RDD也就是每一行你得自己拼 Put 对象。我见过最多的翻车点有三个一是 schema 里列名没带family:qualifier读出来全是 null二是Bytes.toInt对空字节数组直接抛异常三是写回时 RowKey 用了String.getBytes但没指定字符集中文 RowKey 直接乱码。这些后面排障章节会逐个对照真实报错讲。另外提一句HBase 的 zookeeper 地址、表名、列族这些配置建议不要写死在代码里用Configuration从外部传入方便本地和集群切换。下面第二节先把 TaoToken 的凭据通道搭好因为后面读写示例里会顺带演示「读出来的 DataFrame 调模型做字段描述」这个组合动作。2. TaoToken 统一 Key 前置准备让 HBase 作业和模型调用共用一套凭据为什么一个 HBase 读写教程要讲 TaoToken因为真实工程里Spark 作业很少只干读写这一件事。你读出来的 DataFrame可能要调模型做地址标准化、做文本分类、做字段语义补全你写回 HBase 之前可能要用模型生成摘要列。这些调用如果每个作业都单独配 Key运维成本会很高而且 Key 散落在各个application.conf里轮换一次要改十几个地方。TaoToken 在这里的角色是「统一 Key / API 通道管理」你拿到一个 Key通过统一的 Base URL 访问不同模型Spark 作业里只需要维护一份凭据配置。官网入口是 https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content API 地址是 https://taotoken.net/api 这个不加 UTM直接用于代码里的 Base URL。具体要准备三样东西我把它叫「三件套」后面任何接入场景都按这个来配置项值说明Base URLhttps://taotoken.net/api所有模型请求的统一入口API Key在控制台生成形如sk-开头的一串Model ID按需选择比如对话模型、代码模型的标识生成 Key 的路径是控制台里的 API Keys 页面deep link 是 https://taotoken.net/console/api-keys?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。进去之后新建一个 Key复制出来存好页面关掉就看不到了。如果你用的是 Claude Code 这类编码工具或者 Cline 配 MCP配置方式略有不同但三件套不变。Claude Code 的接入文档在 https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 里面有针对 Anthropic 协议的说明。Coding Plan 适合长期跑编码 Agent 的场景入口是 https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。回到 Spark 作业本身我建议把凭据放在一个独立的配置文件里比如taotoken.conf用 Typesafe Config 读取taotoken { baseUrl https://taotoken.net/api apiKey sk-你的Key apiKey ${?TAOTOKEN_API_KEY} modelId 你的模型ID modelId ${?TAOTOKEN_MODEL_ID} }注意apiKey那行写了两次第二行是环境变量覆盖这样本地开发用配置文件集群提交时用--conf或环境变量注入不用改代码。这个模式在 Spark 里特别实用因为spark-submit可以直接传--conf spark.executorEnv.TAOTOKEN_API_KEYxxx。HBase 侧的配置也建议同样处理zookeeper 地址、端口、表名都从配置读hbase { zkQuorum 192.168.1.5 zkPort 2181 table table-001 family cf }这样一份作业配置里HBase 和 TaoToken 的凭据各占一块职责清晰。下一节进入可复制的配置和代码我会把 SparkSession 参数、HBase Configuration、读写示例完整给出。3. 可复制配置SparkSession 参数与 HBase 连接配置完整片段这一节给的是能直接抄进项目的配置。先看 SparkSession 的构建HBase 读写对序列化、内存、并行度都有要求参数不对会出现各种诡异问题。import org.apache.spark.sql.SparkSession import org.apache.hadoop.conf.Configuration import org.apache.hadoop.hbase.HBaseConfiguration import org.apache.hadoop.hbase.mapreduce.TableInputFormat import org.apache.hadoop.hbase.mapreduce.TableOutputFormat import org.apache.hadoop.hbase.io.ImmutableBytesWritable import org.apache.hadoop.hbase.client.Result val spark SparkSession.builder() .appName(SparkHBaseTaoTokenDemo) .master(local[*]) .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .config(spark.kryo.registrator, com.example.HBaseKryoRegistrator) .config(spark.sql.shuffle.partitions, 8) .config(spark.hadoop.hbase.zookeeper.quorum, 192.168.1.5) .config(spark.hadoop.hbase.zookeeper.property.clientPort, 2181) .getOrCreate()几个关键点解释一下。KryoSerializer比默认的 Java 序列化快很多HBase 的Result对象在 shuffle 时会被序列化用 Kryo 能明显减少开销。spark.sql.shuffle.partitions默认 200本地跑小表会开一堆空 task设成 8 或按数据量调。spark.hadoop.前缀的配置会透传给底层的 Hadoop Configuration这样 HBase 客户端也能读到 zookeeper 地址。然后是 HBase 的 Configuration 构建这是读写的核心def buildHBaseConf( zkQuorum: String, zkPort: String, table: String, startRow: Option[String] None, stopRow: Option[String] None, columns: Seq[String] Seq.empty ): Configuration { val conf HBaseConfiguration.create() conf.set(hbase.zookeeper.quorum, zkQuorum) conf.set(hbase.zookeeper.property.clientPort, zkPort) conf.set(TableInputFormat.INPUT_TABLE, table) startRow.foreach(r conf.set(TableInputFormat.SCAN_ROW_START, r)) stopRow.foreach(r conf.set(TableInputFormat.SCAN_ROW_STOP, r)) if (columns.nonEmpty) { conf.set(TableInputFormat.SCAN_COLUMNS, columns.mkString( )) } conf }SCAN_COLUMNS的格式是family:qualifier用空格分隔比如cf:name cf:age。这个参数能显著减少网络传输因为 HBase 只返回你指定的列。注意SCAN_ROW_START和SCAN_ROW_STOP是左闭右开stopRow 本身不包含在内。写 HBase 的配置稍微不同用的是TableOutputFormatdef buildHBaseWriteConf( zkQuorum: String, zkPort: String, table: String ): Configuration { val conf HBaseConfiguration.create() conf.set(hbase.zookeeper.quorum, zkQuorum) conf.set(hbase.zookeeper.property.clientPort, zkPort) conf.set(TableOutputFormat.OUTPUT_TABLE, table) conf.set(mapreduce.outputformat.class, org.apache.hadoop.hbase.mapreduce.TableOutputFormat) conf }mapreduce.outputformat.class这行容易被漏掉漏了会报OutputFormat类型不匹配。写的时候还要注意saveAsNewAPIHadoopDataset要求 RDD 的元素类型是(ImmutableBytesWritable, Put)RowKey 放在ImmutableBytesWritable里列数据放在Put里。TaoToken 的配置读取用 Typesafe Configimport com.typesafe.config.ConfigFactory val ttConf ConfigFactory.load(taotoken.conf).getConfig(taotoken) val baseUrl ttConf.getString(baseUrl) val apiKey ttConf.getString(apiKey) val modelId ttConf.getString(modelId)这样三件套就齐了。下一节进入实际的读写代码和验证动作。4. 读写示例与验证从 HBase 读成 DataFrame 再写回并核对 RowKey先看读。目标是把 HBase 表读成两种形态一种是通用形态rowKey 一个 Map 列装所有 family:qualifier另一种是强 schema 形态每列有明确类型。通用形态的代码import org.apache.spark.sql.Row import org.apache.spark.sql.types._ import org.apache.hadoop.hbase.util.Bytes import org.apache.hadoop.hbase.client.Result def readHBaseGeneric(spark: SparkSession, conf: Configuration): DataFrame { val hbaseRDD spark.sparkContext.newAPIHadoopRDD( conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) val rowRDD hbaseRDD.mapPartitions { iter iter.map { case (_, result) val rowKey Bytes.toString(result.getRow) val familyMap result.getNoVersionMap var columns Map[String, String]() familyMap.forEach { (family, qualifierMap) qualifierMap.forEach { (qualifier, value) val colName Bytes.toString(family) : Bytes.toString(qualifier) columns (colName - Bytes.toString(value)) } } Row(rowKey, columns) } } val schema StructType(Seq( StructField(HBaseRowKey, StringType, nullable false), StructField(columns, MapType(StringType, StringType, valueContainsNull true), nullable true) )) spark.createDataFrame(rowRDD, schema) }跑通之后df.show(false)应该能看到类似这样的结果------------------------------------------------------ |HBaseRowKey |columns | ------------------------------------------------------ |row-001 |{cf:name - 张三, cf:age - 28} | |row-002 |{cf:name - 李四, cf:age - 32} | ------------------------------------------------------强 schema 形态需要你指定列和类型核心是把Result.getValue拿到的字节数组按类型转换def readHBaseTyped( spark: SparkSession, conf: Configuration, columns: Seq[(String, DataType)] ): DataFrame { val hbaseRDD spark.sparkContext.newAPIHadoopRDD( conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) val rowKeyField StructField(HBaseRowKey, StringType, nullable false) val dataFields columns.map { case (name, dt) StructField(name, dt, nullable true) } val schema StructType(rowKeyField : dataFields) val rowRDD hbaseRDD.mapPartitions { iter iter.map { case (_, result) val rowKey Bytes.toString(result.getRow) val values columns.map { case (name, dt) val parts name.split(:) val raw result.getValue(Bytes.toBytes(parts(0)), Bytes.toBytes(parts(1))) if (raw null || raw.isEmpty) null else dt match { case StringType Bytes.toString(raw) case IntegerType Bytes.toInt(raw) case LongType Bytes.toLong(raw) case DoubleType Bytes.toDouble(raw) case FloatType Bytes.toFloat(raw) case _ Bytes.toString(raw) } } Row.fromSeq(rowKey : values) } } spark.createDataFrame(rowRDD, schema) }注意raw null || raw.isEmpty这个判断HBase 里不存在的列返回 null空值返回长度为 0 的数组两种情况都要处理成 null否则Bytes.toInt会抛IllegalArgumentException。写回 HBase 的代码import org.apache.hadoop.hbase.client.Put def writeHBase(df: DataFrame, conf: Configuration, family: String): Unit { val putRDD df.rdd.mapPartitions { iter iter.map { row val rowKey row.getAs[String](HBaseRowKey) val put new Put(Bytes.toBytes(rowKey)) row.schema.fields.foreach { field if (field.name ! HBaseRowKey) { val value row.getAs[Any](field.name) if (value ! null) { put.addColumn( Bytes.toBytes(family), Bytes.toBytes(field.name), Bytes.toBytes(value.toString) ) } } } (new ImmutableBytesWritable(Bytes.toBytes(rowKey)), put) } } putRDD.saveAsNewAPIHadoopDataset(conf) }验证动作分三步。第一步读出来的 DataFrame 注册成临时视图跑一条 SparkSQLdf.createOrReplaceTempView(hbase_table) spark.sql(SELECT HBaseRowKey, columns[cf:name] AS name FROM hbase_table).show()第二步写回后重新读一遍核对 RowKey 数量val before df.count() writeHBase(df, writeConf, cf) val after readHBaseGeneric(spark, readConf).count() println(sbefore$before after$after)第三步用 HBase shell 的scan命令抽查几行确认列值没乱码。如果 RowKey 是中文Bytes.toBytes默认用 UTF-8一般没问题但如果你用了String.getBytes不带参数在某些 JVM 默认字符集下会出问题建议统一写Bytes.toBytes。5. 常见报错排查401、local proxy failed、reading choices 与 OAuth 对照这一节把真实踩过的报错列出来对照着改。报错一java.lang.IllegalArgumentException: offset (0) length (4) exceed the capacity of the array这个通常出现在Bytes.toInt或Bytes.toLong上原因是 HBase 里存的字节数组长度不够。比如你 schema 声明是 IntegerType但实际存的是字符串 abc4 字节读不出来。解决办法是读之前先判断长度或者统一按 String 读出来再在 SparkSQL 里 cast。报错二org.apache.hadoop.hbase.TableNotFoundException: table-001表名拼错或者 zookeeper 地址连到了另一个集群。检查TableInputFormat.INPUT_TABLE的值以及hbase.zookeeper.quorum是否指向正确的集群。本地跑的时候如果 hosts 没配zookeeper 返回的 region server 主机名解析不了也会报连接超时这时候要在本地 hosts 里加上映射。报错三401 Unauthorized这个一般不是 HBase 的错而是调 TaoToken 模型接口时 Key 不对。检查三件套Base URL 是不是https://taotoken.net/apiKey 是不是从控制台复制的完整串Model ID 是不是当前 Key 有权限的模型。如果 Key 里带了空格或换行也会 401。建议在代码里打印 Key 的前 6 位和后 4 位做核对不要打印全量。报错四local proxy failed或连接被拒绝这类报错通常出现在本地网络环境对出站请求有限制的时候。检查你的请求是否走了正确的 Base URL以及本地是否有防火墙拦截。如果是公司内网确认https://taotoken.net/api在允许列表里。不要尝试用任何非正规的网络工具绕过合规环境下应该走正常的网络申请流程。报错五reading choices相关解析失败这个报错一般出现在模型返回的 JSON 结构和你代码里解析的字段对不上时。比如你按choices[0].message.content解析但实际返回结构不同。解决办法是先打印原始响应体确认字段路径再写解析逻辑。用spark.read.json解析模型响应时建议先println(response)看一眼。报错六OAuth 相关错误如果你用的是 Claude Code 或某些需要 OAuth 授权的工具接入报 OAuth 错误通常是回调地址或 token 过期问题。这时候回到接入文档 https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 对照配置。Claude Code 的 Anthropic 协议接入有专门的说明页deep link 是 https://taotoken.net/claude-code-anthropic?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。报错七Task not serializableSpark 里把SparkSession或Configuration直接闭包进 map 里会报这个。解决办法是把需要的配置在 driver 侧读成普通变量或者用broadcast广播出去。HBase 的Configuration不要直接传进 executor 闭包而是在每个 partition 内部重新构建。排查顺序建议先看异常栈最底层的Caused by那才是根因再看是 driver 侧还是 executor 侧报错最后对照上面几条定位。6. 把 HBase 读写接进你的 Spark 管线下一步动作到这里读和写的代码都能跑了。接下来可以做的几件事把readHBaseGeneric和readHBaseTyped封装成一个工具类按配置决定用哪种把 TaoToken 的模型调用封装成一个 UDF在 SparkSQL 里直接SELECT enrich(columns[cf:name]) FROM hbase_table把写回逻辑加上幂等判断避免重复 RowKey 覆盖。如果你要长期跑这类作业Coding Plan 会比按次调用更划算入口在 https://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。需要先验证模型返回效果可以用模型对话页面快速试地址是 https://taotoken.net/models?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。Key 管理和新建在 https://taotoken.net/console/api-keys?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 。最后留一个实用技巧HBase 读出来的 DataFrame 如果列很多df.show()会截断用df.show(false)显示全量或者df.printSchema()先看结构。写回之前先df.rdd.getNumPartitions看一下分区数分区太少写入会慢太多会产生大量小文件一般按 region 数量对齐比较合适。