0%

HBase协处理器

HBase 协处理器:在 RegionServer 上跑自定义代码,实现触发器、存储过程和二级索引

HBase 原生不支持触发器、不支持存储过程、不支持二级索引。但协处理器可以补上这些缺口。

协处理器是运行在 RegionServer 上的用户自定义代码——数据在哪儿,计算就在哪儿,不用把数据拉到客户端再算。

两种类型:

类型 类比 触发方式
Observer 数据库触发器 事件驱动(写数据前、写数据后、读数据时)
Endpoint 存储过程 主动调用(客户端发 RPC 到 RegionServer)

Observer:事件触发,自动执行

Observer 在特定事件发生时自动执行,适合”顺便做点事”的场景。

三种 Observer:

类型 监听什么 典型用途
RegionObserver Region 上的读写操作(Put/Get/Delete/Scan) 写入校验、二级索引同步、数据脱敏
MasterObserver Master 上的 DDL 操作(建表/删表/改表) 权限校验、操作审计
WALObserver WAL 日志写入事件 日志加密、自定义存储

一个最简单的 RegionObserver:在 Put 之后打印 RowKey

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
import org.apache.hadoop.hbase.Coprocessor;  
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.coprocessor.*;
import org.apache.hadoop.hbase.util.Bytes;
import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.wal.WALEdit;
import java.io.IOException;
import java.util.Optional;

// 实现 RegionCoprocessor 和 RegionObserver 接口
public class RegionObserverExample implements RegionCoprocessor, RegionObserver {

// 返回当前类作为 RegionObserver 实例
@Override
public Optional<RegionObserver> getRegionObserver() {
return Optional.of(this); // 必须返回当前实例,否则无法触发
}

// 重写 postPut 方法:数据写入后执行
@Override
public void postPut(
ObserverContext<RegionCoprocessorEnvironment> c,
Put put,
WALEdit edit,
Durability durability
) throws IOException {
// 从 Put 对象中获取 RowKey 并打印
String rowKey = Bytes.toString(put.getRow());
System.out.println("[PostPut] 写入的 RowKey: " + rowKey);
}
}

常用 Observer 钩子方法:

方法 触发时机
prePut 写入前(可做校验,抛异常可阻止写入)
postPut 写入后(做索引同步等)
preGet 查询前
postGet 查询后
preDelete 删除前
postDelete 删除后

Endpoint:主动调用,服务端计算

Endpoint 是”存储过程”——客户端主动调用,在 RegionServer 上执行自定义逻辑,只把结果返回给客户端。

典型场景: 统计一张表有多少行数据。如果用客户端扫全表,数据量大的时候网络传输巨大。用 Endpoint 在服务端计数,只返回一个数字。

Endpoint 实现步骤(简版):

  1. 定义协议接口(继承 CoprocessorProtocol
  2. 实现服务端逻辑
  3. 客户端通过 Table.coprocessorService() 调用

抽象理解:

1
2
3
4
5
6
// 客户端调用
Map<byte[], Long> results = table.coprocessorService(
MyProtocol.class, // 协议接口
scan, // 扫描范围
(protocol) -> protocol.countRows() // 每个 Region 上执行
);

每个 Region 上执行一次 countRows(),结果汇总到客户端。

Observer vs Endpoint 怎么选?

Observer Endpoint
触发方式 自动(事件驱动) 手动(客户端调用)
典型用途 写入校验、索引同步 聚合计算、自定义查询
性能影响 每次操作都执行 只有调用时才执行
开发复杂度 低(实现接口即可) 高(需定义协议)

协处理器怎么加载到 HBase?

两种方式:静态加载(全局生效)和动态加载(表级生效)。

方式1:静态加载——改配置,重启集群

适合全局性功能(比如所有表的写入都做校验)。

  1. JAR 包放到所有 RegionServer 的 $HBASE_HOME/lib/ 目录
  2. 修改 hbase-site.xml
1
2
3
4
5
6
7
8
9
10
11
<!-- 配置 RegionObserver 协处理器 -->  
<property>
<name>hbase.coprocessor.region.classes</name>
<value>com.example.RegionObserverExample</value> <!-- 替换为实际全类名 -->
</property>

<!-- 若需配置 MasterObserver 或 WALObserver,添加对应配置 -->
<property>
<name>hbase.coprocessor.master.classes</name>
<value>com.example.MasterObserverExample</value>
</property>
  1. 重启 HBase 生效

静态加载的影响范围是全局的,所有表都生效。慎用。

方式2:动态加载——只对指定表生效,不用重启

适合特定业务表的功能扩展。

  1. JAR 包上传到 HDFS:
1
hdfs dfs -put my-coprocessor.jar /user/hbase/
  1. 给指定表加载协处理器:
1
2
# 语法:alter '表名', METHOD => 'table_att', 'Coprocessor' => 'JAR路径|全类名|优先级|参数'  
alter 'test', METHOD => 'table_att', 'Coprocessor' => 'hdfs://namenode:9000/user/hbase/coprocessor.jar|com.example.RegionObserverExample|1001|arg1=value1'
  1. 验证是否加载成功:
1
2
describe 'test_table'
# 输出中应该能看到 COPROCESSOR 行

动态加载不需要重启 HBase,加完立即生效。

协处理器卸载

  • 静态加载:需修改 hbase-site.xml 移除配置,重启 HBase。

  • 动态加载:通过alter命令删除表属性:

    1
    2
    # 卸载第 1 个协处理器($1 表示第一个协处理器)  
    alter 'test', METHOD => 'table_att_unset', NAME => 'coprocessor$1'

协处理器的典型应用场景

场景 用什么 怎么实现
二级索引同步 RegionObserver.postPut 主表写入后,自动写索引表
写入数据校验 RegionObserver.prePut 校验字段格式,不合法抛异常阻止写入
数据脱敏 RegionObserver.preGet 查询时对敏感字段做脱敏处理
表级行数统计 Endpoint 服务端统计行数,只返回数字
DDL 操作审计 MasterObserver 建表/删表时记录操作日志

什么时候该用协处理器,什么时候不该用

该用:

  • 需要”顺便做点事”(写入后自动同步索引)
  • 需要服务端聚合(计数、求和)
  • 需要自定义权限校验

不该用:

  • 逻辑可以用 HBase 原生 API 实现(直接用就行了,别加协处理器)
  • 逻辑复杂、耗时长(会影响读写性能)
  • 逻辑需要频繁修改(协处理器改一次要重新加载)

注意事项

  1. 性能影响:Observer 在读写路径上执行,逻辑太慢会拖慢整个 RegionServer
  2. 异常处理要严谨:协处理器抛异常,整个操作会失败
  3. 版本兼容性:协处理器依赖 HBase 内部 API,升级 HBase 可能需要重新编译
  4. 调试困难:代码跑在服务端,本地 debug 不方便,靠日志排查

欢迎关注我的其它发布渠道