![[SPARK][HBASE] Spark 读取文件生成 HFile 并 BulkLoad 批量导入 HBase:TaoToken 统一 Key 通道下的 Scala 实践与运行问题排查](http://pic.xiahunao.cn/yaotu/[SPARK][HBASE] Spark 读取文件生成 HFile 并 BulkLoad 批量导入 HBase:TaoToken 统一 Key 通道下的 Scala 实践与运行问题排查)
1. Spark 读取文件生成 HFile 并 BulkLoad 到 HBase 的完整链路Spark 读取文件生成 HFile 再通过 BulkLoad 批量导入 HBase是离线数仓往 HBase 灌数最常用的一条路径。它绕开了 HBase 的写路径WAL、MemStore、Region 分裂直接把排好序的 HFile 文件挂到 RegionServer 上几千万行数据几分钟就能落地。适合谁适合手里有一批 Parquet/CSV 文件、需要一次性或周期性全量导入 HBase、又不想被逐条 Put 拖慢速度的团队。这条链路的核心约束只有一个HFile 里的 KeyValue 必须严格有序。排序规则是先 rowkey再列族再 qualifier最后才是 timestamp。顺序错了BulkLoad 阶段会直接抛java.io.IOException: added a key lexically larger than previous。很多人第一次跑就卡在这里以为是 Spark 的问题其实是数据没排好。另一个高频坑是运行环境。Spark 程序在 driver 端和 executor 端加载的类不一样HBase 的Connection、Table这类对象不能随便在 driver 端创建后传到 worker否则报Task not serializable。还有依赖问题NoClassDefFoundError: org/apache/spark/SparkConf这种报错八成不是没导包而是依赖顺序或 scope 写错了。这篇会按读文件 → 生成 HFile → BulkLoad 提交 → 验证 → 排错的顺序走一遍给出可直接复制的 Scala 代码、Maven 依赖、Spark 配置以及把 endpoint 切到 TaoToken 统一 Key/API 通道后的验证动作。如果你在跑 Spark 时遇到 401、local proxy failed、429 这类报错第 5 节有对照排查表。先说清楚整体数据流Spark 从 HDFS 读 Parquet/CSV用map把每行转成(ImmutableBytesWritable, KeyValue)元组sortBy按 rowkey列族qualifier 排序saveAsNewAPIHadoopFile写出 HFile 到 HDFS最后用LoadIncrementalHFiles.doBulkLoad把 HFile 推进 HBase 的 RegionServer。整个过程 Spark 只负责生成文件HBase 只负责接收文件中间没有逐条 RPC。我试过在本地master(local)跑通再上集群本地跑通能排除掉大部分依赖和排序问题。本地跑的时候文件路径建议用hdfs:///开头用本地路径容易踩权限和路径解析的坑。下面从环境准备开始。2. TaoToken 统一 Key 通道的前置准备在跑 Spark 之前先把模型调用通道理顺。很多 Spark 作业里会嵌入 LLM 调用做数据清洗、字段补全、标签生成这时候如果每个作业都硬编码一个 endpoint 和 key维护起来很痛苦。TaoToken 提供统一 Key 通道把模型调用收敛到一个 Base URL 和一把 Key 上Spark 作业里只认这两个值。TaoToken 是什么它是一个统一的大模型 API 接入层把不同模型的调用统一成 OpenAI 兼容格式。能做什么你可以用同一把 Key 调用对话模型、代码模型做数据清洗、字段抽取、文本分类。适合谁适合需要在 Spark/离线任务里批量调用模型、又不想为每个模型单独维护 key 和 endpoint 的团队。前置准备分三步。第一步拿到 API Key。访问 https://taotoken.net/api-keys 创建一把 Key复制保存。第二步确认 Base URL。统一通道的 Base URL 是https://taotoken.net/api注意这个地址不带任何查询参数。第三步选一个 Model ID。比如做代码相关任务可以用claude-sonnet-4-5这类模型 ID具体以控制台模型列表为准。把这三个值写进 Spark 配置不要硬编码在代码里。推荐用环境变量或--conf传入export TAOTOKEN_BASE_URLhttps://taotoken.net/api export TAOTOKEN_API_KEYsk-你的key export TAOTOKEN_MODEL_IDclaude-sonnet-4-5然后在 Scala 代码里读取val baseUrl sys.env.getOrElse(TAOTOKEN_BASE_URL, https://taotoken.net/api) val apiKey sys.env.getOrElse(TAOTOKEN_API_KEY, ) val modelId sys.env.getOrElse(TAOTOKEN_MODEL_ID, claude-sonnet-4-5)为什么要走统一通道因为 Spark 作业经常在集群上跑driver 和 executor 的网络出口可能不一致。如果每个模型一个 endpoint防火墙规则要开一堆。统一到一个 Base URL只需要放行一个域名。另外统一 Key 通道方便做用量统计和配额控制避免某个作业把额度跑爆。如果你只是做 HFile 生成和 BulkLoad不涉及模型调用这一步可以跳过但建议还是把 Key 配好因为后面验证请求时会用到。验证模型通道是否通可以用模型对话页面 https://taotoken.net/models 发一条测试消息确认返回正常。注意不要把 Key 写进代码提交到 Git。用环境变量或配置中心。如果 Key 泄露去控制台 https://taotoken.net/console 吊销重建。这一步做完通道就准备好了接下来写 Spark 配置。3. 可复制的 Spark 配置与 HFile 生成代码这一节是核心给出完整的 Maven 依赖、Spark 配置、Scala 代码。先看依赖。HBase 2.0.6 配 Hadoop 2.6.5 是常见组合注意hbase-mapreduce必须引入HFileOutputFormat2和LoadIncrementalHFiles都在这个包里。dependencies dependency groupIdorg.apache.spark/groupId artifactIdspark-core_2.11/artifactId version2.4.8/version scopeprovided/scope /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.11/artifactId version2.4.8/version scopeprovided/scope /dependency dependency groupIdorg.apache.hbase/groupId artifactIdhbase-mapreduce/artifactId version2.0.6/version /dependency dependency groupIdorg.apache.hbase/groupId artifactIdhbase-server/artifactId version2.0.6/version /dependency dependency groupIdorg.apache.hbase/groupId artifactIdhbase-client/artifactId version2.0.6/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version2.6.5/version /dependency dependency groupIdcommons-codec/groupId artifactIdcommons-codec/artifactId version1.13/version /dependency dependency groupIdcommons-io/groupId artifactIdcommons-io/artifactId version2.6/version /dependency /dependencies关键点Spark 依赖用provided因为集群上已经有 Spark 的 jar。HBase 依赖用compile不要用provided否则运行时报NoClassDefFoundError: org/apache/zookeeper/KeeperException。如果确实遇到 jar 冲突再考虑显式排除但先保证依赖能加载进来。接下来是 Spark 配置。序列化用 Kryo比 Java 序列化快而且 HBase 的ImmutableBytesWritable和KeyValue用 Kryo 更稳。val spark SparkSession.builder() .appName(Data2HBase) .master(local[*]) .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .config(spark.kryo.registrator, com.example.HBaseKryoRegistrator) .getOrCreate()如果你不想写 registrator至少把序列化设成 Kryo。然后配置 HBaseval conf HBaseConfiguration.create() conf.set(hbase.zookeeper.quorum, node1:2181,node2:2181,node3:2181) conf.set(hbase.zookeeper.property.clientPort, 2181) conf.set(fs.defaultFS, hdfs://node1:8020/) conf.set(TableOutputFormat.OUTPUT_TABLE, profile_tags) conf.set(hbase.mapreduce.hfileoutputformat.table.name, profile_tags)hbase.mapreduce.hfileoutputformat.table.name这个配置在 HBase 2.x 里必须设否则 BulkLoad 时找不到表。然后创建 Job 并指定输出类型val job Job.getInstance(conf) job.setMapOutputKeyClass(classOf[ImmutableBytesWritable]) job.setMapOutputValueClass(classOf[KeyValue]) val tableName TableName.valueOf(profile_tags) val tableDesc TableDescriptorBuilder.newBuilder(tableName).build() HFileOutputFormat2.configureIncrementalLoadMap(job, tableDesc)现在读数据并生成 KV。假设源数据是 Parquet字段有uid、tag、weightval df spark.read.parquet(hdfs:///data/profile_tags/) val kvRdd df.rdd.map { row val uid row.getAs[Long](uid).toString val tag row.getAs[String](tag) val weight row.getAs[Double](weight) val rowkey uid val cf f val qualifier tag (rowkey, cf, qualifier, weight) }排序是重点。先按 rowkey再按列族再按 qualifier。注意 qualifier 要按字符串排序不要按数值排序否则负数会出问题。val sorted kvRdd.sortBy(tp (tp._1, tp._2, tp._3)) val hfileRdd sorted.map { tp val rowkey tp._1 val cf tp._2 val qualifier tp._3 val value tp._4 val kv new KeyValue( Bytes.toBytes(rowkey), Bytes.toBytes(cf), Bytes.toBytes(qualifier), Bytes.toBytes(value) ) (new ImmutableBytesWritable(Bytes.toBytes(rowkey)), kv) }写出 HFilehfileRdd.saveAsNewAPIHadoopFile( hdfs:///tmp/hfile/profile_tags, classOf[ImmutableBytesWritable], classOf[KeyValue], classOf[HFileOutputFormat2], job.getConfiguration )然后 BulkLoadval conn ConnectionFactory.createConnection(conf) val admin conn.getAdmin val table conn.getTable(tableName) val locator conn.getRegionLocator(tableName) val loader new LoadIncrementalHFiles(conf) loader.doBulkLoad(new Path(hdfs:///tmp/hfile/profile_tags), admin, table, locator)注意doBulkLoad的路径是 HFile 的父目录不是具体文件。如果目录下有多个列族子目录BulkLoad 会自动识别。跑完后检查 HBase 表行数是否对得上。如果你在 Spark 里嵌入了模型调用做数据清洗把 endpoint 指向 TaoTokenval llmConf Map( base_url - https://taotoken.net/api, api_key - sys.env(TAOTOKEN_API_KEY), model - claude-sonnet-4-5 )调用时用 OpenAI 兼容格式POST 到https://taotoken.net/api/v1/chat/completions。这样 Spark 作业里所有模型调用都走统一通道方便排查。4. 验证请求与 BulkLoad 成功结果代码写完先本地跑一遍验证。第一步验证模型通道。用 curl 发一条请求curl -X POST https://taotoken.net/api/v1/chat/completions \ -H Authorization: Bearer $TAOTOKEN_API_KEY \ -H Content-Type: application/json \ -d { model: claude-sonnet-4-5, messages: [{role: user, content: ping}] }返回里如果有choices字段说明通道正常。如果返回 401检查 Key 是否正确、是否带了Bearer前缀。如果返回 429说明触发了限流降低并发或稍后重试。第二步验证 HFile 生成。跑完 Spark 作业后去 HDFS 上看输出目录hdfs dfs -ls hdfs:///tmp/hfile/profile_tags/应该能看到列族目录比如f/里面是data/和index文件。如果目录为空说明saveAsNewAPIHadoopFile没写出数据检查 RDD 是否为空、排序是否报错。第三步验证 BulkLoad。跑完doBulkLoad后用 HBase shell 查行数hbase shell count profile_tags scan profile_tags, {LIMIT 5}如果行数和源数据对得上说明导入成功。如果行数为 0检查 HFile 路径是否正确、表是否存在、RegionServer 是否在线。第四步验证数据正确性。随机抽几行对比源数据和 HBase 里的值 get profile_tags, 12345确认列族、qualifier、value 都对。如果 value 是乱码检查Bytes.toBytes的编码是否一致。实测下来本地local[*]模式跑 10 万行数据生成 HFile 加 BulkLoad 大概 30 秒。集群模式跑 1000 万行5 分钟左右。如果明显慢检查是否发生了 shuffle、是否coalesce(1)导致单点瓶颈。验证通过后把作业提交到集群spark-submit \ --class com.example.Data2HBase \ --master yarn \ --deploy-mode cluster \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.executor.memory4g \ --conf spark.executor.cores2 \ --conf spark.executor.instances10 \ your-jar-with-dependencies.jar注意--deploy-mode cluster时driver 在集群上跑环境变量要提前配好或者用--conf spark.yarn.appMasterEnv.TAOTOKEN_API_KEYxxx传入。5. 本篇常见报错排查对照表这一节把高频报错和定位思路列出来对照排查。报错一java.io.IOException: added a key lexically larger than previous原因HFile 里的 KeyValue 没有严格排序。检查sortBy的字段顺序必须是 rowkey、列族、qualifier。如果 qualifier 是数值转成字符串再排不要用数值排序。另外注意sortBy是全局排序会触发 shuffle数据量大时慢但必须做。报错二NoClassDefFoundError: org/apache/zookeeper/KeeperException原因缺少 HBase 相关依赖或者依赖 scope 写成了provided。把hbase-client、hbase-server、hbase-mapreduce的 scope 改成compile。如果还报错检查是否有多个版本的 zookeeper jar 冲突用mvn dependency:tree排查。报错三NoClassDefFoundError: org/apache/spark/SparkConf原因Spark 依赖没加载或者依赖顺序问题。把 Spark 依赖放在其他依赖上面确保先加载。如果用了provided本地跑要改成compile集群跑再改回provided。报错四Task not serializable原因在 driver 端创建了 HBaseConnection、Table对象然后传到 executor 使用。解决办法是把连接创建放在foreachPartition里每个 partition 创建一次连接用完关闭。不要用foreach用foreachPartition。报错五401 Unauthorized原因TaoToken API Key 错误或缺失。检查Authorization头是否带了Bearer前缀Key 是否过期。去 https://taotoken.net/api-keys 重新生成一把。报错六local proxy failed原因网络出口不通或者 Base URL 写错。检查https://taotoken.net/api是否可达用 curl 测试。如果集群有网络限制确认出口放行了该域名。报错七429 Too Many Requests原因请求频率超限。降低 Spark 作业的并发度或者加退避重试。在模型调用处加Thread.sleep或指数退避。报错八reading choices相关解析错误原因返回体格式不符合预期。检查请求的Content-Type是否为application/json请求体是否符合 OpenAI 格式。如果返回的是错误信息先打印完整响应再解析。报错九OAuth相关报错原因如果用了 OAuth 方式鉴权检查 token 是否过期。TaoToken 统一通道用 API Key 即可不需要 OAuth。如果代码里混用了 OAuth 逻辑去掉。报错十BulkLoad 后行数为 0原因HFile 路径写错或者表名不匹配。检查doBulkLoad的路径是否是 HFile 父目录检查hbase.mapreduce.hfileoutputformat.table.name是否和实际表名一致。排查思路总结先看报错类型是依赖问题、排序问题、序列化问题还是网络问题。依赖问题用mvn dependency:tree排序问题看sortBy字段序列化问题看连接创建位置网络问题用 curl 测通道。6. 长期编码与 Agent 场景的通道选择如果你只是偶尔跑一次 BulkLoad上面的配置够用了。但如果你在做长期的 Spark 编码、数据管道维护、或者 Agent 类应用建议把模型调用通道固定下来用 TaoToken 的 Coding Plan。它适合需要持续调用模型做代码生成、数据清洗、字段补全的场景统一 Key 通道能省掉很多配置维护。长期编码场景下把 Base URL、Key、Model ID 三件套写进配置文件不要散落在代码里。比如用settings.json或auth.json管理{ base_url: https://taotoken.net/api, api_key: sk-你的key, model: claude-sonnet-4-5 }如果是 Claude Code 这类工具配置路径通常在~/.claude/settings.json把ANTHROPIC_BASE_URL指向https://taotoken.net/apiANTHROPIC_API_KEY填你的 Key。如果是 Cline MCP 场景在 MCP 配置里填 Base URL 和 Key。如果是 Codex检查auth.json里的 endpoint 是否指向统一通道。三件套缺一不可Base URL 决定请求发到哪Key 决定能不能过鉴权Model ID 决定用哪个模型。少一个都会报错。401 通常是 Key 问题local proxy failed 通常是 Base URL 或网络问题reading choices 通常是 Model ID 或返回格式问题。对于 Agent 场景建议把模型调用封装成一个工具函数统一处理重试、超时、错误码。429 加退避401 直接抛异常local proxy failed 检查网络。这样 Spark 作业里调用模型时不用每个地方都写一遍错误处理。最后给一个实用技巧在 Spark 作业里调用模型时用mapPartitions而不是map每个 partition 创建一个 HTTP 连接池复用连接减少握手开销。批量请求时把多条数据拼成一个 prompt一次调用处理多条降低请求数避免 429。通道配好后去 https://taotoken.net/coding-plan 看长期编码方案去 https://taotoken.net/doc 看接入文档。验证模型是否可用用 https://taotoken.net/models 发测试消息。管理 Key 用 https://taotoken.net/api-keys。控制台在 https://taotoken.net/console。整套链路跑通后你会发现 Spark 生成 HFile 加 BulkLoad 并不复杂复杂的是排序和依赖。把这两块搞定剩下的就是调参和验证。遇到报错先对照第 5 节大部分问题都能定位。