前几天有个朋友找我,说他接了个人用户画像的项目,HDFS 里趴着几百亿条行为日志,业务方提了个要求:标签查询必须控制在两秒内返回。他第一反应是把 Spark 离线算完写到 MySQL,结果发现单表几个亿行之后 MySQL 根本扛不住,连索引都救不回来。我跟他讲,你缺的其实是 HBase。这种组合我这两年至少搭过四五个类似的系统:Spark 负责把脏数据、多源数据洗干净、聚合好,HBase 负责把结果按 RowKey 存起来并提供毫秒级点查,两者合在一起,才真正打得动"海量数据实时查询"这个需求。这篇把我自己的选型思考、搭建细节、三种读写实现和踩坑经验完整拆开讲,适合正在做实时查询方案、或者刚接触 Spark 和 HBase 集成的人直接对照着用。
1. 先想清楚再动手:Spark与HBase的分工边界与适用场景
1.1 它俩到底谁管什么
很多刚接触大数据的人容易把 Spark 和 HBase 搞混,以为都是"大数据技术",随便选一个就行。实际上它俩的分工完全不同,理解不了这一步,后面代码写得再好也会在架构上翻车。
Spark 是计算引擎,核心能力是把计算任务分发到几十上百个节点上并行执行。它能做复杂的清洗、关联、聚合、机器学习,但计算结果默认落回 HDFS、本地文件系统或者你指定的外部存储。问题就在这:如果结果数据量是亿级甚至十亿级,你要让业务方直接去扫 HDFS 里的文件做实时查询,那基本是灾难——文件列表扫描、小文件问题、格式解析开销,任何一个都足以让你的接口超时。
HBase 是分布式列式存储,它的设计目标非常纯粹:单行点查(Get)做到毫秒级,范围扫描(Scan)做到可控延迟。它底层用 HDFS 存数据,但通过 Region 切分、MemStore、BlockCache 这些机制把"随机点查"变快了。它不擅长做复杂的跨表关联和聚合计算,你要在 HBase 里做 Join 或者跑机器学习,那纯粹是和自己过不去。
所以这套组合的本质是:Spark 负责"算",HBase 负责"存和取"。Spark 把复杂的计算前置完成,HBase 把计算结果以最适合点查的格式存下来,业务侧通过 RowKey 直接拿数据。简单粗暴,但是极其有效。
1.2 典型场景:什么样的需求真的适合它
我拆过几个实际项目,符合这套方案特征的场景非常清晰,你可以对号入座。
第一个是用户画像。原始数据是行为日志,几百亿条,需要按用户维度聚合出年龄、性别、偏好、活跃度几十个标签。这类计算 Spark 跑得动,但结果是"用户ID -> 标签集合"这种 KV 结构,天然适合 HBase。查询侧拿到用户 ID,一个 Get 就能把全部标签拉回来,响应时间稳定在一两毫秒。
第二个是订单多维查询。订单数据量大,业务上需要按用户查最近订单、按订单号查详情、按时间倒序翻页。这里面有个技巧:订单号本身是天然的唯一键,但用户查最近订单需要 RowKey 以用户为前缀,所以通常用"userId反转 + 时间戳倒序"做 RowKey,Spark 负责把实时订单流和维度表关联好,再写入 HBase。
第三个是物联网设备状态。设备上报频率高,Spark Streaming 或 Structured Streaming 做实时清洗后,把设备最新状态写入 HBase,RowKey 用"设备ID + 时间戳",既能点查设备状态,也能按时间范围扫设备历史。
这三个场景有一个共同点:读多写少、查询模式明确、以点查为主。如果你的场景需要灵活的多维分析、复杂的 SQL 聚合、或者频繁的全表扫描,那 HBase 未必是最佳选择,要考虑 Elasticsearch、Doris 或者数仓方案。
1.3 什么情况建议你换个组合
我也见过不少把 Spark 和 HBase 硬凑在一起的项目,结果成本和复杂度翻倍。
数据量在几百 GB 到几 TB 量级、查询也不复杂,直接上 MySQL 分库分表或者 Elasticsearch 更划算。HBase 的运维成本和硬件要求都不低,小数据量上它属于杀鸡用牛刀。
另外,如果你只做离线分析,不要求实时查询,结果写回 HDFS 上以 Parquet 格式存着就够了,没必要引入 HBase 这个中间层,毕竟 Spark 自己读 Parquet 也很快。
还有一种情况:业务方说"我要在上面跑各种灵活的 SQL、还要 Join 多张表"。HBase 原生并不擅长这个,虽然可以借助 Phoenix 或者 Hive 映射表曲线救国,但性能往往要打折扣,而且会引入额外的调度和一致性开销。真遇到这种需求,建议直接考虑 Doris、ClickHouse 或 TiDB,别让 HBase 硬撑。
2. 环境搭建里的版本兼容、端口清单与配置坑
2.1 版本怎么配才不折腾
版本搭配是这套方案里第一个坑。网上教程很多停留在 Spark 2.x + HBase 1.x 时代,照抄下来你会发现 API 全是红的。我自己现在生产上用的组合是:Spark 3.x(3.2 ~ 3.5)+ HBase 2.x(2.4 ~ 2.5)+ Hadoop 3.x + Java 8 或 11。
为什么这么配?HBase 2.x 的客户端 API 相比 1.x 改动不小,推荐使用官方提供的hbase-client依赖,而 Spark 3.x 的newAPIHadoopRDD对 Hadoop 2.x 和 3.x 的兼容性已经做得很好。如果你还在用 HBase 1.x,虽然也能跑通,但 RegionServer 在热点和稳定性上的表现差距明显,BulkLoad 支持的完善程度也不同,没必要在旧版上吃苦头。
下表是我测过的几组版本组合,可以直接抄:
| Spark | HBase | Hadoop | 备注 |
|---|---|---|---|
| 3.5.x | 2.5.x | 3.3.x | 推荐,接口稳定 |
| 3.3.x | 2.4.x | 3.2.x | 老项目常用,资料多 |
| 3.2.x | 2.3.x | 3.2.x | 兼容性尚可,注意ZK版本 |
| 2.4.x | 1.4.x | 2.7.x | 不推荐,API差异大 |
Java 版本也要统一。Spark 3 在 Java 8 和 11 下都能跑,HBase 2.x 也要求 Java 8 以上,但注意如果你的 HBase 节点里还跑着老版本 Hadoop 组件,可能出现javax.xml.bind找不到这种 JDK 11 的典型报错,这时候要么切换到 Java 8,要么把缺失的模块显式加入依赖。
2.2 hbase-site.xml 里的关键配置项
HBase 部署好后,客户端要连上去,靠的就是hbase-site.xml里几个关键配置。我最低配置会放这几个:
<configuration> <property> <name>hbase.rootdir</name> <value>hdfs://hadoop-cluster/hbase</value> </property> <property> <name>hbase.zookeeper.quorum</name> <value>node1,node2,node3</value> </property> <property> <name>hbase.zookeeper.property.clientPort</name> <value>2181</value> </property> <property> <name>zookeeper.session.timeout</name> <value>120000</value> </property> <property> <name>hbase.client.write.buffer</name> <value>4194304</value> </property> </configuration>hbase.rootdir告诉 HBase 数据存在 HDFS 的哪个目录;hbase.zookeeper.quorum是 ZK 集群地址,客户端通过它找到 HBase 的元数据;zookeeper.session.timeout默认是 90 秒,但 Spark 任务跑长任务时偶尔会因网络抖动触发 ZK 会话超时,建议调到 120 秒以上,代价是异常恢复时间变长,你需要自己权衡。
hbase.client.write.buffer默认是 2MB,批写入时可以加大,减少 RPC 次数。
2.3 端口清单与网络放行
端口这块我踩过一次很惨的坑:集群内部一切正常,但是 Spark 任务从另一套机器上提交时连不上 HBase,日志反复报Connection refused。查了半天才发现安全组里漏开放了 RegionServer 的端口。这里直接把端口清单给你列全:
| 组件 | 端口 | 说明 |
|---|---|---|
| HMaster RPC | 16000 | HBase 1.x 以前是 60000,2.x 默认 16000 |
| HMaster Web UI | 16010 | 浏览器看集群状态 |
| RegionServer RPC | 16020 | Spark 客户端写入和读取主要走这个 |
| RegionServer Web UI | 16030 | 看单个 RegionServer 状态 |
| ZooKeeper | 2181 | 客户端通过 ZK 定位 HBase 元数据 |
如果你的 Hadoop 节点组在云上,记住主机名之间要能互相解析,安全组和防火墙要同时放行 2181、16000、16020,Web UI 端口按需开放。还要注意 HBase 内部会通过 ZK 分配一个hbase.master.info.port和hbase.regionserver.info.port,默认就是 16010 和 16030,如果你改过端口,客户端配置里要同步改。
2.4 Spark 这边需要做的三个小动作
很多人把 HBase 环境搭好了,Spark 里一跑就报ClassNotFoundException或者NoClassDefFoundError,大部分原因是 Spark 的 classpath 里没有 HBase 的客户端依赖和配置文件。
第一个动作:把 HBase 的hbase-site.xml放到 Spark 每台节点任务的 classpath 中。这不是玄学,HBase 客户端启动时会自动从 classpath 加载hbase-site.xml,如果缺失,它会从默认配置初始化,quorum 变成localhost,那肯定连不上。最省事的做法是把配置文件复制到$SPARK_CONF_DIR下,spark-submit 时会自动加入。
第二个动作:在代码里显式覆盖关键项。即使 classpath 里没有配置文件,也可以在代码中HBaseConfiguration.create()之后直接conf.set("hbase.zookeeper.quorum", "..."),这样至少保证关键连接信息不丢。
第三个动作:保证 HBase Client 依赖和 Spark 依赖不冲突。HBase 自带了一些 Hadoop 类和 Guava 类,直接丢进 Spark classpath 可能和 Spark 自带的版本冲突,典型表现是NoSuchMethodError。建议把 HBase 相关依赖打成 fat jar 的时候保留hbase-client、hbase-common、hbase-server这几个模块,但排除掉它的 Hadoop 依赖,让 Spark 用自己带的那份。
3. 从RDD到BulkLoad:三种读写集成方式的代码级拆解
3.1 基于TableInputFormat/TableOutputFormat的通用读写
这套方案的核心是 Spark 提供的newAPIHadoopRDD和saveAsHadoopDataset。它们把 HBase 的 InputFormat/OutputFormat 桥接到 Spark 的 RDD 模型上,逻辑最透明,适合大多数场景。
先看读侧。假设表里存的是用户画像,列族info下有age、gender两列:
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName} 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 val hbaseConf = HBaseConfiguration.create() hbaseConf.set("hbase.zookeeper.quorum", "node1,node2,node3") hbaseConf.set("hbase.zookeeper.property.clientPort", "2181") hbaseConf.set(TableInputFormat.INPUT_TABLE, "user_profile") // 限定扫描范围,避免扫全表 hbaseConf.set(TableInputFormat.SCAN, org.apache.hadoop.hbase.client.Scan .toString(Array.emptyByteArray, Array.emptyByteArray)) val rdd = spark.sparkContext.newAPIHadoopRDD( hbaseConf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result] ) val parsed = rdd.map { case (_, result) => val rowKey = Bytes.toString(result.getRow) val age = Bytes.toString(result.getValue(Bytes.toBytes("info"), Bytes.toBytes("age"))) val gender = Bytes.toString(result.getValue(Bytes.toBytes("info"), Bytes.toBytes("gender"))) (rowKey, age, gender) }读侧最关键的一点是:一定要把 Scan 范围限定到具体 RowKey 范围,而不是默认全表扫描。TableInputFormat.SCAN接受一个 Base64 编码的 Scan 对象,如果不想处理编码细节,可以直接用Scan.toString()生成范围,或者像我上面代码里那样传空范围先跑通,再根据业务需要加上setStartRow和setStopRow。
写侧稍微繁琐一点。要把任意 RDD 写成 HBase,核心是构造Put对象,再通过TableOutputFormat写回:
import org.apache.hadoop.hbase.mapreduce.TableOutputFormat import org.apache.hadoop.mapred.JobConf val jobConf = new JobConf(hbaseConf) jobConf.setOutputFormat(classOf[TableOutputFormat]) jobConf.set(TableOutputFormat.OUTPUT_TABLE, "user_profile") val writeRDD = parsed.map { case (rowKey, age, gender) => val put = new Put(Bytes.toBytes(rowKey)) put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("age"), Bytes.toBytes(age)) put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("gender"), Bytes.toBytes(gender)) (new ImmutableBytesWritable(), put) } writeRDD.saveAsHadoopDataset(jobConf)这个方式在中小数据量(百万到千万行)下很稳定,生产中也常见。它的缺点在于每条 Put 都要经过 HBase 写路径,会有 WAL 写日志的开销,数据量上千万级别后会明显变慢,这时候就要考虑下面的 BulkLoad 方案。
3.2 PySpark写入HBase:Thrift与JVM桥接两条路
很多团队习惯用 PySpark 做数据处理,但 HBase 的原生客户端是 Java 的,PySpark 写 HBase 有三种常见方式,我分别说下优缺点。
第一种是走 Thrift 服务,配合happybase客户端。这种方式代码最简单,也是网上教程里出现最多的:
from pyspark.sql import SparkSession import happybase spark = SparkSession.builder.appName("spark_to_hbase").getOrCreate() df = spark.read.json("hdfs://cluster/tmp/user_event.json") def write_partition(iterator): connection = happybase.Connection(host="hbase-thrift-host", port=9090) table = connection.table("user_profile") with connection.batch(batch_size=1000) as batch: for row in iterator: batch.put( str(row["user_id"]).encode(), { b"info:age": str(row["age"]).encode(), b"info:gender": str(row["gender"]).encode(), }, ) connection.close() df.foreachPartition(write_partition)batch对象会自动攒批批量提交,比单条 put 快得多。Thrift 方案适合数据量几百万到千万级别、对写入延迟不敏感的场景。缺点是 Thrift 吞吐上限不高,而且 HBase 需要单独开启 Thrift 服务。
第二种是走 Py4J 桥接,直接在 Python 里调用 Spark JVM 的 HBase Java API。这种方式不需要额外部署 Thrift,吞吐取决于 Java API 本身,但代码会显得很绕。核心思路是用SparkSession的_jvm属性拿到 JVM 里的类引用,然后像写 Java 一样逐行调用ConnectionFactory、Table、Put。这里建议封装一个 Java 工具类放到 classpath 里,PySpark 只传行数据给这个工具类,否则 Python 侧直接拼 Java API,代码会非常难维护。
第三种是干脆用 Scala 写完 HBase 再在 PySpark 里调用。实际上我见过的多数项目最终都走向这条:业务计算逻辑用 PySpark 负责,落 HBase 的写入模块单独用 Scala 编译成一个 jar,通过spark-submit --jars挂载进来,再用 PySpark 的rdd.mapPartitions调用 Scala 里暴露的写入函数。这样做的好处是,HBase 写入这种重操作能吃到 Java API 的全部性能,同时业务团队又能享受 PySpark 的便利。
3.3 BulkLoad:全量导入千万行级数据的最优解
如果你要做千万级、亿级甚至更大的全量初始化,或者周期性离线重建结果表,那直接调put写入基本不可行,因为每次写入都要走 WAL、走 MemStore、可能会触发 flush,RegionServer 的 CPU 和 IO 会被打满。生产上这个量级我一定是走 BulkLoad。
BulkLoad 的原理很巧妙:与其一条条往 HBase 里放数据,不如用 Spark 直接把要导入的数据在 HDFS 上生成 HBase 底层的存储文件(HFile),然后再把这些文件"无缝"挂载到表里。整个过程不经过 RegionServer 的写入路径,自然不写 WAL,也没有 MemStore flush,速度比批量 put 快一个数量级以上。
生成 HFile 的 Scala 代码大致是这样:
import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat2 import org.apache.hadoop.mapreduce.Job // 先创建表和 Region 预分区,务必提前把 Region 分好 val conn = ConnectionFactory.createConnection(hbaseConf) val admin = conn.getAdmin val tableName = TableName.valueOf("user_profile") val table = conn.getTable(tableName) val regionLocator = conn.getRegionLocator(tableName) // 会根据表的 Region 分布自动配置 HFile 输出格式 val job = Job.getInstance(hbaseConf) job.getConfiguration.set("mapreduce.output.fileoutputformat.outputdir", "/tmp/hfile_output") HFileOutputFormat2.configureIncrementalLoad(job, table, regionLocator) val ts = System.currentTimeMillis() val hfileRDD = dataRDD.map { row => val kv = new KeyValue( Bytes.toBytes(row.rowKey), Bytes.toBytes("info"), Bytes.toBytes("age"), ts, Bytes.toBytes(row.age) ) (new ImmutableBytesWritable(Bytes.toBytes(row.rowKey)), kv) } hfileRDD.saveAsNewAPIHadoopFile( "/tmp/hfile_output", classOf[ImmutableBytesWritable], classOf[KeyValue], classOf[HFileOutputFormat2], job.getConfiguration )生成完 HFile 后,把它挂到表上:
hbase org.apache.hadoop.hbase.mapreduce.LoadIncrementalHFiles /tmp/hfile_output user_profileBulkLoad 有个让人容易忽视的前提:表必须先预分区,而且分区边界要设计合理。HFile 生成完以后是按 RowKey 分布落到不同 HDFS 路径里的,如果表只有一个 Region,那所有 HFile 只能塞进同一个 Region,写入时不但没快,反而会引起单点热点和数据倾斜。这一点在下面第 4 节重点讲。
3.4 结果如何供业务实时查询
写进 HBase 只是第一步,业务侧要真正拿到毫秒级查询,还得设计一条完整的链路。我的做法是:Spark 算完的结果直接写在 HBase,业务服务启动时初始化一个 HBaseConnection单例,查询时用Get按 RowKey 拿数据。实际测试下来,单行 Get 平均延迟在 1-3 毫秒,接口整体耗时主要在网络传输和序列化上。
如果你觉得裸 HBase API 不够友好,可以在业务层包一层缓存,比如在查询服务前加 Redis,热点 Key 命中后直接返回,没命中再回源 HBase。这个设计可以大幅降低查询请求对 HBase 的压力,但要注意缓存和 HBase 的一致性,最稳妥的策略是设置合理的 TTL,允许短暂延迟。
另外一个容易被忽略的点:HBase 的表扫描(Scan)虽然也能支持范围查询,但它的性能和点查完全不是一个量级。之前有个项目把分页查询直接做成 Scan,每次扫描上千行数据,延迟飙到几百毫秒,而且大 Scan 还会占用 RegionServer 的堆内存,拖慢其他点查。后来改成两级设计:先用 Spark 把分页要展示的字段预聚合到一张"索引表"里,业务查询只做点查,问题才彻底解决。
4. 性能从分钟到秒级的关键:RowKey设计、预分区与Spark并行度对齐
4.1 RowKey三原则:散列、短、唯一
RawKey 设计是整个 HBase 方案里最重要的一环。面试中经常问"RowKey 设计要考虑哪些因素",项目里同样绕不开。我自己总结了三条原则,缺一不可。
第一,散列。RowKey 要尽量均匀分布在所有 Region 上,避免热点。假设你直接用自增 ID 做 RowKey,前几位都是000001、000002,那么它们全部落到同一个 Region,这个 Region 会被打爆,其他 Region 闲置。常用做法是取 RowKey 的 MD5 前缀做盐(salt),比如md5(userId).substring(0, 4) + userId,这样前四位是散列值,数据分布就均匀了。
第二,短。HBase 在每个 KeyValue 里都会冗余存储一份 RowKey,RowKey 越长,存储开销越大,网络传输和内存占用也越高。一份 1 亿行的表,RowKey 每多一个字节,光索引就多了近 100MB 的空间,这还不算多版本数据。能用 Int 就用 Int,能用 8 字节 Long 就别用 20 字节的字符串。
第三,唯一。RowKey 是全表唯一键,设计不好会造成数据互相覆盖。订单场景里如果直接用userId + orderId,那同一个用户同一订单只会保留一条记录,如果业务里允许同一订单有多个状态快照,你就需要把时间戳(或版本号)也拼进去。
4.2 预分区方案:一步到位还是动态调整
HBase 初始建表时默认只有一个 Region,写入全量数据会先打在一个 Region 上,直到文件大小超过阈值才触发拆分。更麻烦的是,拆分过程中 RegionServer 要做数据分裂和迁移,期间会影响读写。所以生产上建表第一步就是预分区。
你可以直接在 HBase Shell 里指定分区数,也可以让 HBase 用算法自动生成分区边界:
hbase shell create 'user_profile', 'info', {NUMREGIONS => 20, SPLITALGO => 'HexStringSplit'}HexStringSplit会按十六进制字符串范围把 RowKey 空间切成 20 段。这个方法适合 RowKey 前缀是散列值(比如 MD5 前四位)的场景。如果你用的是可读字符串前缀(比如userId本身就是字符串),更精确的做法是手动指定 split 点:
hbase shell create 'user_profile', 'info', {SPLITS => ['1000','2000','3000','4000']}手动 split 点意味着你必须清楚自己的 RowKey 分布。我的经验是:别拍脑袋定分区数。先拿一周真实数据在测试集群上跑一遍,统计 RowKey 前缀分布,找出实际的数据热点段,再决定 split 点。没有真实数据,至少也用抽样数据估算规模,不然你建出来的分区很可能和实际数据分布对不上,Region 有的空转、有的被压垮。
4.3 列族与表参数:影响读写的隐藏开关
列族(Column Family)是 HBase 里另一个性能开关。在生产环境,我基本坚持一个表只用一个列族,理由很简单:HBase 的 Region 会把不同列族的数据分开存储,同一个 RowKey 的多个列族数据可能落在不同的 HFile 里,扫描和点查时需要额外寻址;列族多于一个的情况下,MemStore 是分列族刷新还是统一刷新、Region 分裂时列族间的一致性处理,都会增加复杂度。除非你有特别强的隔离需求(比如冷热数据分离),否则一个列族足够了。
另外两个参数也要注意。VERSIONS控制保留几个版本的数据,默认是 1,如果你不需要保留历史值,就保持默认,版本越多,存储和查询开销越大。TTL控制数据过期时间,存储层会自动清理过期数据,适合那种"只保留最近 90 天"的业务,这类过期数据不用你写任务去删,省事很多。
还有压缩。HBase 表默认不压缩,但在海量数据场景下建议开启 Snappy 压缩:
hbase shell alter 'user_profile', {NAME => 'info', COMPRESSION => 'SNAPPY'}Snappy 压缩比较省 CPU,而且在 HBase 里解压是惰性的,点查只解压命中的块,开销很小。实测下来,用户画像这种高重复文本数据,开启压缩后存储量能降 40% 到 60%。
4.4 Spark并行度怎么和Region对齐
Spark 读写 HBase 时,并行度设置会直接影响性能。这里面有个比较容易理解但不小心就踩歪的原则:Spark 的分区数量要和 HBase 的 Region 数量匹配。
读场景里,Spark 每个分区对应一个或多个 HBase Region 的扫描任务。如果分区数远小于 Region 数,一个 task 要连续扫多个 Region,单点任务时间拉长;如果分区数远大于 Region 数,会产生大量空 task,浪费调度资源。我一般做法是:先看表的 Region 总数,再把 Spark 读 RDD 的分区数设为 Region 总数的 1 到 2 倍,让每个 Region 至少被一个 task 覆盖到。
写场景同理。如果每个分区里的数据量不均,很容易造成某些 RegionServer 写入压力远大于其他节点。配合上你 RowKey 里的盐值前缀,Spark 写入的 task 分布就能基本对齐到不同 Region 上。
另外两个读侧参数建议重点调:Scan.setCaching(500)和Scan.setBatch(1000)。setCaching控制客户端一次 RPC 拉取多少行到本地缓存,默认值太小(100 左右)会导致大量小 RPC,慢得离谱;setBatch控制一次 Scan 返回多少列,如果你的行比较宽,需要按需设置,避免一次拉回太多无用数据撑爆内存。
4.5 一条把查询从分钟压到秒级的完整链路
我完整复盘过一个真实案例:订单流水查询,数据 8000 万行,原始方案是 Spark 每天全量计算后把结果写到 MySQL,业务查询按用户 ID 模糊搜索,最慢的接口要 30 秒。后来改成 Spark + HBase 方案。
第一步,RowKey 设计成userId反转 + 下单时间戳倒序。为什么反转?因为自增 userId 的前缀集中,反转后前缀就散开了。为什么时间戳倒序?因为用户查"最近订单"时,最近的记录排在前面,配合 Scan 从头取就能做到"取最近 N 条"的目的,不用扫全量。
第二步,建表时按 userId 前缀分布做 40 个预分区,HFile 写入走 BulkLoad,每天凌晨把前一天计算结果增量加载进去。
第三步,查询接口改成 Get 或短 Scan(限定到单个用户的 RowKey 范围),实测接口 95% 请求在 5 毫秒内返回,最慢的也不超过 20 毫秒。这个量级,之前 30 秒到 5 毫秒的差距,就是 RowKey 和存储模型设计带来的质变。
5. 跑了一年的生产环境,这些坑我劝你提前避开
5.1 Connection被每个Task创建一次,RegionServer直接警告
这是新手最容易犯的错误:在 PySpark 或 Scala 的map函数里直接创建 HBaseConnection,每个 Task 执行一次就创建一个独立连接。你要知道,HBase 客户端Connection底层是重量级对象,内部维护了到 RegionServer 的长连接池和元数据缓存,一个几百分区的任务同时创建几百个连接,RegionServer 的线程池瞬间被打满,日志里会疯狂刷Blocked waiting for connection。
正确做法有两个,二选一:
一个是把Connection放在mapPartitions外部创建,一个分区只创建一个连接,处理完整个分区的数据后关闭。另一个更优:使用ConnectionFactory.createConnection创建一个全局单例,多个 Task 共享,HBase 客户端本身会复用内部的连接池。实测下来,单例方式最省资源,但要注意 HBaseConnection是线程安全的,而Table对象不是,所以每个 Task 里还是要通过connection.getTable(...)获取独立的Table实例。
5.2 序列化异常与Kryo注册的坑
Spark 任务里如果用 RDD 传递Put、Get、Result这类 HBase 对象,序列化是个大问题。HBase 很多内部类实现了Writable接口,但 Spark 默认的 Java 序列化器不一定认识它们,会出现ClassNotFound或DoNotRetryIOException。
解决方案有两个方向:如果数据量不大,直接把 HBase 对象改装成简单的 Scala/Python 对象(字符串、元组、case class)在 RDD 里流转,只在写 HBase 的最终步骤构造Put,这样序列化风险最低。如果数据量大、必须直接处理Result,那就在spark-submit里加上配置:
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.kryo.registrator=org.apache.hadoop.hbase.spark.HBaseKryoRegistrator这个HBaseKryoRegistrator是 HBase 官方提供的 Kyro 注册器,能处理大部分 HBase 内部类。还有一种隐藏坑:Kryo 注册器用到了org.apache.hadoop.hbase.spark包,如果依赖里没引入hbase-spark模块,会直接报NoClassDefFoundError,排查时容易一头雾水。
5.3 数据倾斜、Region热点与Salting
我遇到过线上 RegionServer 一台 CPU 100%、其他几台闲着的诡异景象。查下来发现是 RowKey 前缀设计有问题:某几个用户 ID 前缀覆盖了 30% 的数据,对应 Region 直接被写满读爆。
解决手段是 Salting(加盐):在 RowKey 前面加一个随机前缀或哈希前缀。比如原本 RowKey 是userId + timestamp,加工成(bucketId)+userId + timestamp,其中bucketId = hash(userId) % N。这样数据分布就均匀了 N 倍。代价是查询的时候你需要知道这个用户属于哪个 bucket,然后对每个可能的 bucket 前缀各发一次 Get。如果 bucket 数量是 16 或 32,这个代价可以忽略。
另外,如果你用的是 MurmurHash 这类哈希算法,注意Integer.MIN_VALUE取模的边界问题,算出来的 bucket 可能为负数,RowKey 前缀带上负号会影响排序,处理时先取绝对值。
5.4 RegionServer GC暂停把查询拖死
HBase 是 Java 应用,RegionServer 的 GC 是影响写入和查询延迟的最大隐形杀手。典型场景:大量写入触发 MemStore 频繁 flush,堆内存抖动,Full GC 暂停几秒,期间所有查询全部超时。
我自己的调优顺序是:先看 RegionServer 日志里的 GC 时间和 MemStore 占用率,如果频繁出现ms级别的 Full GC,优先调整 JVM 堆大小和hbase.regionserver.global.memstore.size(默认 0.4,即堆的 40% 用于 MemStore)。高写入场景我会把 MemStore 上限调低到 0.3,给 BlockCache 留更多空间,减少写多读时的缓存驱逐。
另外要警惕大 Scan 对堆内存的冲击。一个 Scan 如果缓存了太多行,会直接把 BlockCache 挤爆,其他点查的命中率直线下降。解决方式是严格控制 Scan 范围,给每个 Scan 设置合理的setBatch和setCaching,必要时限制单次 Scan 的总行数,做了这些之后,线上 GC 暂停次数肉眼可见地下降。
5.5 双写一致性与幂等设计
Spark 算完的结果写 HBase 往往不是"最后一步"——很多业务还要同时更新 Redis 缓存、Elasticsearch 索引或业务库。多系统双写最容易出一致性问题:HBase 写成功了,Redis 更新失败,客户端读到旧数据。
我这边比较稳妥的方案是:以 HBase 为唯一数据源,其他系统全部通过异步消费 HBase 的写入日志来同步。具体做法是让 Spark 任务先写 HBase,然后发送一条消息到消息队列(比如 Kafka),下游消费者拿到消息后去更新 Redis、ES,消费逻辑做幂等。这样即使某个下游写失败,也能通过消息重试追上进度。
如果实在不想引入消息队列,退而求其次的做法是定期回刷:每天凌晨用 Spark 重算全量结果,同时覆盖 HBase 和缓存。这个方案实现简单,但实时性差,只适合对一致性容忍度高的场景。反正记住一个原则:不要直接让业务接口去更新多个系统,那等于把一致性问题的复杂度全部堆到了调用链上。
最后再分享一点我自己的体会:这套组合上线到第 10 个月的时候,有一次凌晨数据突然写不进去,RegionServer 日志全是Blocked waiting,排查半天才发现是一个定时任务用错了 RowKey 模板,把 90% 的数据打到了同一个 Region 上。从那之后我每次改 RowKey 设计和预分区方案,都会先拿真实数据在测试集群压一遍分布图,确认没有热点再上生产。Spark 加 HBase 这套方案,"会写读写代码"只占 20% 的工作量,剩下 80% 的精力全花在数据分布、参数调优和运维细节上。你要是正准备搭这套方案,我的建议是先从 RowKey 和预分区开始设计,而不是急着写读写代码——顺序反了,后面全是填不完的坑。