6.1 协处理器:把计算搬到数据旁边 本节摘要:协处理器(Coprocessor)是运行在 RegionServer 进程内的用户代码钩子,分 Observer(观察者,类数据库触发器)与 Endpoint(端点,类存储过程)两类。本节实现一个二级索引 Observer 的骨架,讲加载与卸载方式,并给出它带来的风险清单——服务端代码没有隔离,写得差等于给集群埋雷。 为什么需要下推 第 5 章的二级索引靠应用双写。但写入方很多(多个服务、多个任务都在写订单表)时,每个写入方都得记得维护索引,漏一个就是数据不一致。换个思路:能否在 RegionServer 处理写入的环节自动挂钩子,主表一写,索引表跟着写,写入方无感?这就是 Observer 的用途。
本节摘要:协处理器(Coprocessor)是运行在 RegionServer 进程内的用户代码钩子,分 Observer(观察者,类数据库触发器)与 Endpoint(端点,类存储过程)两类。本节实现一个二级索引 Observer 的骨架,讲加载与卸载方式,并给出它带来的风险清单——服务端代码没有隔离,写得差等于给集群埋雷。
第 5 章的二级索引靠应用双写。但写入方很多(多个服务、多个任务都在写订单表)时,每个写入方都得记得维护索引,漏一个就是数据不一致。换个思路:能否在 RegionServer 处理写入的环节自动挂钩子,主表一写,索引表跟着写,写入方无感?这就是 Observer 的用途。
更广义的场景是"聚合下推":按 Region 求局部和、再汇总全局,数据不动逻辑动——Endpoint 的用途。两类用一个关系数据库的类比即可区分:Observer 像触发器,Endpoint 像存储过程。
Observer 挂在 RegionServer 处理请求的关键路径上(-region 生命周期、客户端读写、WAL、Master 操作等钩子点)。实现第 5 章的订单二级索引:
public class OrderIndexObserver implements RegionCoprocessor, RegionObserver { private Connection conn; // 复用 4.1 节的连接生命周期规则 @Override public Optional<RegionObserver> getRegionObserver() { return Optional.of(this); // 声明自己是 Region 级观察者 } @Override public void start(CoprocessorEnvironment env) throws IOException { conn = ConnectionFactory.createConnection(env.getConfiguration()); } @Override public void postPut(ObserverContext<RegionCoprocessorEnvironment> c, Put put, WALEdit edit, Durability durability) { // 只挂钩主表 orders 的写入 if (!"orders".equals(c.getEnvironment().getRegion().getTableName().getNameAsString())) return; byte[] row = put.getRow(); // 主表行键 rev(uid)+inv(ts)+oid byte[] oid = extractOid(row); // 业务函数 取订单号段 Put idx = new Put(oid); idx.addColumn(Bytes.toBytes("f"), Bytes.toBytes("rk"), row); // 索引指向主表 try (Table t = conn.getTable(TableName.valueOf("orders_by_oid"))) { t.put(idx); } catch (IOException e) { // 关键决策点:索引写失败怎么办?记日志进补偿队列,绝不能抛出 // 抛异常会连带主表写入失败,把可用性问题扩大 } } }
代码里有三个值得停留的设计点:连接在 start 里建一次(4.1 节规则在服务端同样成立);postPut 在主写完成后触发,失败不回滚主写;异常必须吞掉并补偿——这是协处理器最重要的纪律。
打包为 jar 后加载,三种方式按需选:
hbase:110:0> alter 'orders', METHOD => 'table_att', hbase:111:1* 'coprocessor' => hbase:112:1* 'hdfs:///libs/order-index-v3.jar|org.example.OrderIndexObserver|1001|'
四段式:jar 路径、类全名、优先级、可选参数。卸载用 METHOD => 'table_att_unset', NAME => 'coprocessor$1'。此外还有配置文件全局加载(hbase.coprocessor.region.classes,全表生效、需重启)与 HBase 2.x 的动态加载接口。
验证:向 orders 写一条,随后查索引表:
hbase:113:0> put 'orders', '...rev(u1001)...5678', 'cf:status', 'PAID' hbase:114:0> get 'orders_by_oid', 'ORD20240819005678' COLUMN CELL f:rk timestamp=..., value=...rev(u1001)...
写入方完全无感,索引自动生成。
Endpoint 定义一个客户端可调用的服务接口(Protobuf 定义加实现),在每个 Region 上并行执行局部计算,客户端汇总。适合行数统计、分桶聚合这类"数据不动逻辑动"的场景。不过实践中 Endpoint 的开发成本(Protobuf 代码生成、版本兼容)常高于收益——同样的聚合交给 Spark(6.3 节)扫描计算更省心,Endpoint 只在小规模、强实时、低延迟聚合里还有不可替代性。HBase 自带的 IntegrationTestBigLinkedList 等工具与一些计数场景是它的残余阵地。
第 5 章用应用双写实现二级索引,本节用 Observer 实现同一件事,两者到底怎么选?把决策维度摆开:
| 维度 | 应用层双写 | Observer 索引 |
|---|---|---|
| 一致性 | 两次独立写,中间宕机会漏 | 与主写同一 RegionServer,时序更近 |
| 覆盖面 | 只覆盖走该代码路径的写入 | 覆盖一切写入(含 Shell、MapReduce、其他服务) |
| 故障半径 | 索引写失败只影响该业务 | 钩子异常可能拖累宿主 RegionServer |
| 升级成本 | 改代码即生效 | 重编译重打包,滚动卸载再加载 |
| 可观测性 | 业务日志直接可见 | 要到 RegionServer 日志里翻 |
关键差异在第二行。双写方案最怕的就是"某个写入方忘了维护索引"——Shell 里临时修一条数据、历史数据回灌任务、第三方系统直写,这些路径都不经过你的业务代码,索引就此出现黑洞。Observer 挂在表上,不管谁来写都会触发,这是它真正不可替代的地方。反过来看第五行:双写出问题,业务同学看自己的日志五分钟定位;Observer 出问题,值班同学要在 RegionServer 堆栈里找你的 jar,排障成本完全不同量级。
所以我的建议保持不变但加一条例外:当写入路径无法收敛到一个代码入口时(多团队共用表、有批处理回灌),Observer 从"最后手段"升级为"必要手段"——此时宁可承受它的排障成本,也不能接受索引悄悄烂掉。上线前把灰度期的对账任务准备好:定时抽样比对主表与索引表的行数差,差值超过阈值就告警,这是协处理器方案的安全网。
协处理器跑在 RegionServer 进程内,没有沙箱隔离:
工程上的四条纪律:钩子逻辑尽量短平快;一切外部调用设超时;异常只记不抛;灰度先挂单表观察一周再推广。我在实际项目里的排序永远是:能应用层双写就不用 Observer,能 Spark 算就不用 Endpoint——协处理器是最后手段,不是第一反应。
⚠️ 常见坑:协处理器 jar 放在本地目录而 RegionServer 没有同步,出现"有的机器索引生效有的不生效"的幽灵故障。jar 必须放 HDFS 全集群可见,或走动态加载接口统一分发。
写侧的下推讲完了,6.2 转向读侧的两台加速器:布隆过滤器与双层读缓存。