在线批量写和 Bulk Load 解决的是同一个输入问题:大量记录要进入同一张 HBase 表。二者的成本模型完全不同。BufferedMutator 仍然走 RegionServer 写入协议,写入进入 WAL、MemStore、flush 和 compaction;Bulk Load 先离线生成合法 HFile,再让运行中的 RegionServer 接管这些 StoreFile。

本篇只抓一个核心问题:什么时候该用在线批量写,什么时候该先生成 HFile;两条路径各自把 CPU、网络、WAL、MemStore 和 compaction 成本放在哪里。

两条批量写路径图

flowchart LR
  A[应用批量记录] --> B[BufferedMutator]
  B --> C[RegionServer RPC]
  C --> D[WAL]
  C --> E[MemStore]
  E --> F[flush]
  F --> G[HFile]
  A --> H[MapReduce/Spark 预排序]
  H --> I[HFileOutputFormat2]
  I --> J[HFile]
  J --> K[BulkLoadHFiles/completebulkload]
  K --> L[RegionServer adopt StoreFile]

BufferedMutator 适合在线流量:记录持续产生、表仍然服务读写、每批数据规模不大到足以单独跑离线任务。它把多次 mutation 缓冲起来,降低 RPC 次数,但每条变更仍然进入正常写路径。

Bulk Load 适合离线导入:数据已经在 HDFS 或对象存储兼容文件系统上,能先按表的 Region 边界和 row key 排序,导入窗口可控。它绕过普通 client API 的逐条 RPC 成本,但不能绕过 HBase 对 HFile 合法性、Region 边界和 hbase:meta 一致性的要求。

这一区分很重要。Bulk Load 不是“直接把文件丢进 HDFS”。官方 Bulk Loading 文档明确把流程拆成两步:先用 HFileOutputFormat2 生成内部 StoreFile,再用 completebulkloadBulkLoadHFiles 把文件加载进运行中的集群。

BufferedMutator 的语义边界

BufferedMutator 是 HBase 2.6 public API。API 索引把它描述为面向单表的批量、异步写接口,BufferedMutatorParams 用于创建参数,BufferedMutator.ExceptionListener 用于接收异步异常。

关键对象如下:

对象 角色 需要记住的边界
BufferedMutator 单表写缓冲 不是事务容器
BufferedMutatorParams 缓冲区、listener 等参数 只影响客户端侧行为
ExceptionListener 异步异常回调 不能把回调缺失当成功证明
flush() 推送当前缓冲 不等于 Region flush 到 HFile
close() 关闭并释放资源 应处理关闭时暴露的异常

BufferedMutator 的核心收益是批量合并 RPC。应用连续生成 Put 时,不必每一条都调用 Table.put 形成一个同步往返;客户端可以把多条 mutation 合并发送。服务端收到后仍按普通 mutation 处理,满足当前 durability 配置,进入 WAL 和 MemStore。

这意味着 BufferedMutator 不改变行级原子性边界。单行 mutation 仍按 HBase 的行级语义执行,多行批量没有自动事务。一个批次里部分行成功、部分行失败是必须处理的正常形态。ExceptionListener 的存在就是这个边界的 API 证据。

最小 Java 形状如下,目标为 HBase 2.6 public API。当前本地环境没有 HBase client 依赖,本段标记为 UNVERIFIED_RUNTIME,API 名称按官方 2.6 User API 核对。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
Configuration conf = HBaseConfiguration.create();
TableName table = TableName.valueOf("lab:bulk_write");
try (Connection connection = ConnectionFactory.createConnection(conf)) {
BufferedMutatorParams params = new BufferedMutatorParams(table).listener((exception, mutator) -> {
for (Row row : exception.getFailedOperations()) {
System.err.println("failed row: " + Bytes.toStringBinary(row.getRow()));
}
});
try (BufferedMutator mutator = connection.getBufferedMutator(params)) {
Put put = new Put(Bytes.toBytes("u#0001"));
put.addColumn(Bytes.toBytes("cf"), Bytes.toBytes("score"), Bytes.toBytes(42L));
mutator.mutate(put);
mutator.flush();
}
}

这段示例只表达接口关系,不提供虚构输出。真实运行时应创建 lab:bulk_write 表,并让 cf 列族存在。

Bulk Load 的对象和状态

Bulk Load 把写入成本从“在线服务端逐条处理”前移到“离线作业预处理”。官方流程里的三个对象要分清:

对象 输入 输出 状态约束
HFileOutputFormat2 已排序的 KeyValue/Cell HBase 内部 HFile 每个 HFile 应落在单个 Region 范围
TotalOrderPartitioner 当前 Region 边界 分区文件 输出 key range 对齐 Region
BulkLoadHFiles / completebulkload HFile 目录和表名 被 RegionServer 接管的 StoreFile 依赖健康的 hbase:meta

HFileOutputFormat2.configureIncrementalLoad() 会基于当前 Region 边界设置分区,使输出文件的 key range 和 Region 对齐。这个动作不是优化细节,而是正确性边界:一个 HFile 横跨多个 Region 时,导入阶段需要拆分文件;官方文档明确提醒,Region 边界在准备和加载之间变化会导致自动拆分,效率不高。

导入阶段由 completebulkload 或程序化 BulkLoadHFiles 完成。工具会遍历准备好的 HFile,判断每个文件属于哪个 Region,联系对应 RegionServer,让 RegionServer 接管文件并放入表的存储目录。文件成为 StoreFile 后,读路径会像读取普通 flush 出来的 HFile 一样读取它。

这条路径绕过 WAL 和 MemStore 的逐条写入成本,因此 CPU 和网络消耗低于普通 API 批量导入。但这也带来三个工程后果。

第一,导入前必须保证输出数据已经按 row key 排序,并且 family 目录结构符合目标表。排序错误不是“慢一点”,而是生成的 HFile 不符合 HBase 读取假设。

第二,目标表 Region 边界变化会影响导入效率。准备作业和加载命令之间间隔越长,越容易遇到 split 后边界不一致。

第三,复制语义要单独配置。官方 Bulk Loading Replication 文档说明,bulk loaded HFile 的复制需要开启 hbase.replication.bulkload.enabled,默认不是自动复制普通 WAL edit 的同一路径。

最小实验

本地环境没有可用 HBase 2.6.6 集群,以下实验是 UNVERIFIED_RUNTIME。它只列出可复现步骤和观察点,不给出伪造输出。

准备在线批量写表:

1
hbase shell

Shell 内创建 namespace 和表:

1
2
create_namespace 'lab'
create 'lab:bulk_write', 'cf'

在线写入实验观察三类指标:

1
hbase hbtop

需要关注对应表或 Region 的 #WRITE/SMEMSTORE#SFBufferedMutator 写入期间 #WRITE/SMEMSTORE 会变化;显式 mutator.flush() 只推送客户端缓冲,不应被误读成 Store flush。

Bulk Load 实验的最小路径如下:

1
hbase org.apache.hadoop.hbase.mapreduce.ImportTsv -Dimporttsv.columns=HBASE_ROW_KEY,cf:score -Dimporttsv.bulk.output=/tmp/hbase-bulk-out lab:bulk_write /tmp/input.tsv

加载生成的 HFile:

1
hbase completebulkload /tmp/hbase-bulk-out lab:bulk_write

观察点不是终端输出里的某个固定字符串,而是三类状态:目标表是否读到导入行、HDFS 上 bulk 输出目录是否被移动或清空、RegionServer 上 StoreFile 数是否增加。

失败恢复

BufferedMutator 的失败恢复重点在客户端。批量异步异常可能延迟暴露,应用必须在 listener、flush()close() 三个位置处理失败。重试策略应以 row key 幂等为前提:同一 row、family、qualifier、timestamp 的重复写会覆盖同一 Cell;未指定 timestamp 的 Put 重试会生成新版本或覆盖最新版本,具体可见性受列族版本配置影响。

Bulk Load 的失败恢复重点在文件和元数据。导入前失败时,可以修正离线作业重新生成 HFile。导入中失败时,不能手工把 HFile 挪进 Region 目录假装完成;应重新运行 completebulkload,让工具根据 hbase:meta 和 Region 位置判断文件归属。官方文档的“adopting stray data”也要求先确认 hbase:meta 健康。

Region split 是 Bulk Load 的常见干扰。离线作业按旧边界生成的 HFile,在加载前目标表 split,导入工具会拆分文件再加载。这个过程可恢复,但浪费 I/O;生产导入通常会先预分区,导入期间控制 split 或缩短准备到加载的时间窗口。

工程迁移

场景 选择 原因
实时写入、持续小批量 BufferedMutator 保留在线语义和普通恢复路径
离线历史回灌、TB 级数据 Bulk Load 把排序和编码成本放到离线作业
需要立即复制到异地 普通 API 或显式配置 bulkload replication Bulk Load 复制不是普通 WAL edit 的自然副作用
数据需要复杂校验 先离线校验再 Bulk Load 错误 HFile 进入表后回滚成本高
行级幂等 upsert BufferedMutator 重试模型更简单

可迁移模式是“把在线路径的随机成本前移成离线路径的顺序成本”。LSM 系统里的 SSTable ingestion、搜索引擎的 segment merge、对象存储上的 manifest commit 都是同类模式。前提是导入文件必须满足服务端内部格式和分区边界。

常见误解

误解一:BufferedMutator 是事务。它只是客户端写缓冲和异步批量接口,不改变 HBase 的行级原子性边界。

误解二:Bulk Load 绕过 HBase。Bulk Load 绕过普通写入 RPC 的逐条路径,但加载阶段仍由 HBase 工具、RegionServer、表描述和 hbase:meta 共同完成。

误解三:Bulk Load 一定更快。预排序、生成 HFile、边界变化拆分、复制配置和导入窗口都会成为成本。小规模实时写入通常不值得启动一条离线链路。

误解四:bulk loaded HFile 会天然进入普通复制流。官方文档要求为 bulk load replication 单独启用和配置。

练习

  1. 设计一个 10 亿行历史回灌方案,列出预分区、离线排序、HFile 输出、导入和校验五个阶段。
  2. 把同一批输入分别设计成 BufferedMutator 写入和 Bulk Load 写入,标出 WAL、MemStore、HFile 和 compaction 成本分别在哪里发生。
  3. 检查一个现有批量写程序是否处理了 listener、flush()close() 三处异常。

系列导航

参考资料

证据等级:API 和机制为 VERIFIED_SOURCE;本地没有 HBase 2.6.6 集群,实验为 UNVERIFIED_RUNTIME