A coprocessor is Java code that HBase loads into the Master, a RegionServer or a region and calls at defined points in its own request processing. Observers react to operations (a Put, a Get, a region opening, a WAL append) and Endpoints add new RPC methods that run next to the data. Apache Phoenix secondary indexes, HBase's own access control and visibility labels, and many audit and validation layers are built this way.
The interfaces are small, which makes coprocessors look easy. The difficulty is in the architecture around them: which host loads your class and when, what order it runs in relative to other coprocessors, which locks are held when your hook fires, what happens to the server when your code throws, and what happens to the cluster when your hook makes an RPC of its own. This article walks through that machinery for HBase 2.x, builds a small validating observer, and then looks at the patterns (and the anti-patterns) that the machinery implies.
Four hosts and the coprocessor lifecycle
Every coprocessor lives in a host, and the host decides its lifetime and its blast radius. HBase 2.x has four: MasterCoprocessorHost for MasterObserver hooks on DDL and cluster operations, RegionServerCoprocessorHost for server-wide events such as log rolls and replication, RegionCoprocessorHost for RegionObserver and Endpoint code, and WALCoprocessorHost for WALObserver. Region coprocessors are by far the most common, and they are instantiated per region: a RegionServer hosting 400 regions of a table holds 400 instances of your class, each created when its region opens and stopped when it closes, moves or splits.
In 2.x a coprocessor implements a host-specific interface such as RegionCoprocessor and exposes its observers through getters like getRegionObserver(), which return an Optional. This replaced the 1.x style of extending BaseRegionObserver, and code ported from older tutorials often compiles but never fires because the getter was not overridden. The host calls start(env) once per instance, passing a RegionCoprocessorEnvironment that gives access to the region, the server configuration merged with any properties supplied at load time, and a shared cluster connection.
Four hosts and the coprocessor lifecycle
Every coprocessor lives in a host, and the host decides its lifetime and its blast radius. HBase 2.x has four: MasterCoprocessorHost for MasterObserver hooks on DDL and cluster operations, RegionServerCoprocessorHost for server-wide events such as log rolls and replication, RegionCoprocessorHost for RegionObserver and Endpoint code, and WALCoprocessorHost for WALObserver. Region coprocessors are by far the most common, and they are instantiated per region: a RegionServer hosting 400 regions of a table holds 400 instances of your class, each created when its region opens and stopped when it closes, moves or splits.
In 2.x a coprocessor implements a host-specific interface such as RegionCoprocessor and exposes its observers through getters like getRegionObserver(), which return an Optional. This replaced the 1.x style of extending BaseRegionObserver, and code ported from older tutorials often compiles but never fires because the getter was not overridden. The host calls start(env) once per instance, passing a RegionCoprocessorEnvironment that gives access to the region, the server configuration merged with any properties supplied at load time, and a shared cluster connection.
System and table loading, jars and priority
There are two ways to load a coprocessor, and they behave differently. System coprocessors are listed in hbase-site.xml under keys such as hbase.coprocessor.region.classes, hbase.coprocessor.master.classes and hbase.coprocessor.regionserver.classes. They must be on the server classpath, apply to every table, and change only with a rolling restart. Table coprocessors are attributes of a table descriptor. They can name a jar on HDFS, carry an explicit priority and key-value properties, and load when the table's regions open, so altering the table reopens its regions with the new definition.
// HBase 2.x: attach a table coprocessor with a versioned jar path and properties
TableDescriptor current = admin.getDescriptor(TableName.valueOf("events"));
TableDescriptor updated = TableDescriptorBuilder.newBuilder(current)
.setCoprocessor(CoprocessorDescriptorBuilder
.newBuilder("com.example.hbase.TenantGuardObserver")
.setJarPath("hdfs:///hbase/coprocessors/tenant-guard-1.4.0.jar")
.setPriority(Coprocessor.PRIORITY_USER)
.setProperty("tenant.guard.max.value.bytes", "1048576")
.build())
.build();
admin.modifyTable(updated); // regions of 'events' reopen with the coprocessor loadedJars loaded from HDFS go through a dedicated coprocessor classloader that copies the jar locally and caches it by path. Overwriting a jar in place and expecting servers to pick it up is a classic mistake: some regions keep the cached class, others reload, and the cluster ends up running two versions. Put the version in the file name and change the path on every release.
Priority decides order. System coprocessors default to Coprocessor.PRIORITY_SYSTEM and table coprocessors to Coprocessor.PRIORITY_USER; lower values run first, so built-in security observers see a request before your code does. Within the same priority, order follows the order of definition. Two switches matter for safety: hbase.coprocessor.enabled turns all loading off, and hbase.coprocessor.user.enabled turns off table-level loading while keeping system coprocessors, which is a reasonable default on a shared cluster where tenants can alter their own tables.
Where hooks fire on the write and read paths
The value of an observer depends on where its hook sits. Per the RegionObserver documentation, prePut and postPut run once per mutation, outside the row-lock window. preBatchMutate runs after the row locks for the whole batch have been acquired and server timestamps assigned, and it may still add or change mutations. postBatchMutate runs after the edits have been written to the WAL and applied to the MemStore but before the locks are released and the MVCC read point advances, so readers cannot see the data yet. postBatchMutateIndispensably runs after the batch whether it succeeded or failed, which makes it the place to release resources you took in a pre-hook.
That ordering creates the main rule of observer design: work inside the lock window must be fast and local. Each millisecond spent in preBatchMutate holds row locks that block other writers to the same rows and holds an RPC handler thread that cannot serve anyone else. Validation, enrichment computed from the mutation itself, and adding cells to the same row belong there. Network calls do not.
On the read side, preGetOp can serve a result without touching the store, and scanner hooks such as preScannerOpen and postScannerNext can filter or transform rows. A post-scan filter runs after HBase has already read the cells, so pushing a real Filter into the scan is usually cheaper than filtering in a hook.
Worked example: a tenant guard observer
Here is a complete observer that enforces a multi-tenant row-key convention and a value size limit. It throws DoNotRetryIOException so the client fails fast instead of retrying a request that can never succeed:
public class TenantGuardObserver implements RegionCoprocessor, RegionObserver {
private long maxValueBytes;
@Override
public Optional<RegionObserver> getRegionObserver() {
return Optional.of(this); // without this, no hook ever fires
}
@Override
public void start(CoprocessorEnvironment env) {
// properties set on the CoprocessorDescriptor arrive in the environment's configuration
maxValueBytes = env.getConfiguration().getLong("tenant.guard.max.value.bytes", 1 << 20);
}
@Override
public void prePut(ObserverContext<RegionCoprocessorEnvironment> ctx, Put put,
WALEdit edit, Durability durability) throws IOException {
byte[] row = put.getRow();
if (row.length < 9 || row[8] != '|') {
throw new DoNotRetryIOException("row key must be <8-byte tenant id>|<rest>");
}
for (List<Cell> cells : put.getFamilyCellMap().values()) {
for (Cell cell : cells) {
if (cell.getValueLength() > maxValueBytes) {
throw new DoNotRetryIOException("value over " + maxValueBytes + " bytes");
}
}
}
}
}The four-argument prePut shown is the long-standing 2.x signature; later 2.x releases add a variant without Durability and deprecate this one, so check the javadoc for the minor version you build against. Test it with the HBase testing utility in a mini-cluster rather than with mocks: the bugs that matter are about loading, ordering and exceptions, which mocks do not reproduce.
Exceptions, aborts and bypass
Exceptions are where coprocessors most often hurt production. An IOException thrown from a hook is part of the contract: the operation fails and the exception travels back to the client. Any other Throwable (a NullPointerException, a ClassCastException, an OutOfMemoryError) is treated as a broken coprocessor. With hbase.coprocessor.abortonerror at its default of true, the hosting server aborts. Its regions move to other RegionServers, load your coprocessor there, hit the same input and abort those servers too. A single malformed row can walk across a cluster this way.
Setting abortonerror to false makes the server unload the faulty coprocessor and return an error to the client instead. That keeps the server up, but the coprocessor is now running on some regions and not others, which for an access-control or indexing observer is a silent correctness failure. Neither setting is safe on its own; the fix is to catch everything inside your hooks, convert it to an IOException deliberately, and emit a metric when you do.
The second sharp edge is ctx.bypass(). Since 2.0 it is honoured only on a subset of hooks, mostly pre-hooks on RegionObserver, and when set it skips the core operation and every coprocessor after yours in the chain. A user-priority observer that bypasses can therefore hide an operation from another coprocessor that expected to see it. Treat bypass as a last resort and document which hooks rely on it.
Writing to other tables from a hook
The tempting pattern is to write to another table from inside a hook: maintain an index table, append an audit row, increment a counter elsewhere. The environment hands you a connection, so it is one line of code. It is also the most common way coprocessors take clusters down.
Do the arithmetic. A RegionServer has a fixed pool of RPC handlers, set by hbase.regionserver.handler.count (30 in recent defaults). Suppose a write-heavy table's observer makes one synchronous Put to an index table that takes 5 ms. Each handler can now complete at most about 200 writes per second, so the server tops out near 6,000 writes per second regardless of CPU. Worse, the index Put lands on another RegionServer and needs one of its handlers, while that server's own observers are making index writes back to the first. Under load every handler on both servers ends up waiting on a handler on the other, nothing completes, and clients see timeouts on two servers with idle CPUs. That is a distributed deadlock built from a thread pool.
There are three defensible answers. Keep cross-region writes off the client handler pool, as Phoenix does by installing its own RPC scheduler with separate handlers for index traffic. Make the secondary write asynchronous, queued and replayable, and accept an index that lags. Or move the work out of HBase entirely with replication or change data capture into a downstream system. What you should not do is a synchronous remote write on the shared handler pool.
Secondary indexes and endpoints
Secondary indexes deserve their own note because they are the canonical coprocessor use case and the most misunderstood. A coprocessor cannot make a write to the data table and a write to an index table atomic: they are different regions, often on different servers, with different WALs. Any design has to choose an order and a recovery story.
Phoenix's current global indexes are a good model of what that takes. Index rows are written first in an unverified state, then the data row, then the index rows are marked verified. A reader that hits an unverified index row checks the data table and repairs or ignores the index entry. The result is consistent for readers without distributed transactions, but it needed read-side logic, a background tool to rebuild and verify indexes, and careful handling of concurrent updates to the same row. Treat that as the bar any home-grown index observer has to clear; most teams are better off using Phoenix or keeping the index in another system.
Endpoints are the other half of the API: a protobuf Service registered by a region coprocessor and invoked by clients over a key range, so that aggregation runs next to the data and only results cross the network. The calling side, including fan-out and partial failure, is covered in calling coprocessor endpoints from the client.
Rolling out a coprocessor safely
A team adds the tenant guard above to a 40-server cluster. A safe rollout follows the architecture:
- Build the jar with HBase dependencies marked provided, version it in the file name, and upload it to an HDFS directory that only the HBase service user can write, since anyone who can replace the jar can run code as HBase.
- Load it on a canary table first with
modifyTable, then send malformed and oversized writes and confirm the client receivesDoNotRetryIOExceptionand no server logs a coprocessor abort. - Measure write latency before and after on the canary. A validation hook should add microseconds; anything measured in milliseconds means it is doing I/O or heavy allocation.
- Roll to production tables one at a time. Each alter reopens that table's regions, so schedule it away from peak and watch region-in-transition counts.
- To roll back, alter the table to remove the coprocessor or point it at the previous jar path. Never delete the old jar until no table references it.
Failure modes
- Cascading aborts. An unchecked exception on one input aborts each server that opens the region. Catch everything in hooks and keep
abortonerrordeliberately chosen. - Handler starvation. Synchronous remote calls in hooks cap throughput and can deadlock pairs of servers. Look for RPC queue growth with low CPU.
- Version skew. In-place jar replacement leaves mixed classes across regions. Use versioned paths.
- Silent no-op. Ported 1.x code with no
getRegionObserver()override loads cleanly and never runs. Assert in tests that the hook fired. - Memory growth. Per-region instances multiply any per-instance cache by the region count. Share state per server deliberately, or not at all.
What to do next
- Inventory every coprocessor on your clusters: host, priority, load method, jar path and owner.
- Read each hook and mark any network or disk I/O inside the lock window; plan to move it out.
- Wrap every hook body in a catch-all that converts unexpected errors to
IOExceptionand counts them. - Decide
hbase.coprocessor.user.enabledandabortonerrorper cluster and write the reason down. - Switch all jar references to versioned HDFS paths with restricted write access.
- Add a mini-cluster test per coprocessor that loads it the way production does and asserts that hooks fire.
Related reading: what coprocessors can and cannot do, Apache Phoenix on HBase and RegionServer internals.