A coprocessor endpoint is code you deploy into the RegionServers. You call it like a remote procedure, and it runs next to the data in every region you aim it at. Instead of pulling ten million cells across the network to count them, the client sends one small request per region and gets back one small answer per region. The server side of this, meaning how endpoints and observers are written, loaded and isolated, is covered in HBase Coprocessors. This article covers the other half: the code that calls them.
Client calls are where most production problems with endpoints show up. Results come back keyed by region and change shape when regions split. A row argument looks as if it is sent to the server and is not. The blocking and async APIs disagree about whether the stop row is inclusive. Errors arrive per region, so a call can half-succeed. The APIs below were checked against the HBase 2.6 javadoc and source; the code uses the 2.x client.
What the client is actually calling
From the client's point of view an endpoint is a protobuf service: a named set of RPC methods with request and response messages. The server registers an implementation of that service with each region. The client gets a stub for the same service and points it at an RPC channel that HBase connects to one region. Everything the endpoint needs must be in the request message. The endpoint sees its region's data and nothing else, so cross-region work (sums, merges, top-k) is always finished on the client.
So a client call has three parts: which regions (one row, or a key range), what to send each region (a request message), and how to combine the answers (a map, a callback, or a reducer you write). Here is the service used through the rest of this article:
// tenant_stats.proto -- compile with the protoc that matches the protobuf-java
// version your HBase client exposes for coprocessor services (com.google.protobuf).
syntax = "proto2";
option java_package = "com.example.hbase.cp";
option java_outer_classname = "TenantStatsProtos";
option java_generic_services = true;
message CountRequest {
required bytes start_row = 1; // inclusive
required bytes stop_row = 2; // exclusive
optional int64 min_ts = 3;
}
message CountResponse {
required int64 rows = 1;
required int64 cells = 2;
}
service TenantStatsService {
rpc count(CountRequest) returns (CountResponse);
}java_generic_services = true makes protoc generate the abstract service class with newStub and newBlockingStub. In HBase 2.x, the client-side coprocessor APIs use plain com.google.protobuf types (RpcController, RpcCallback), not the shaded copy HBase uses internally. Keep this generated code in a small shared jar that both the endpoint and its clients depend on, and change it only in wire-compatible ways: add optional fields and never renumber.
How a range call fans out
The fan-out is per region, which explains most of the client-side surprises. Region boundaries rarely line up with your logical range, so the endpoint in region B also sees rows below t1|, and region C sees rows past the range. That is why CountRequest carries start_row and stop_row, and why the endpoint must clamp its scan to them. Without that, the tenant count silently includes neighbours' rows from the edge regions.
The fan-out also follows region count. A table that has split from 40 regions to 400 makes the same call ten times more expensive in RPCs, with ten times more chances for one region to be moving at that moment. Regions and splits explains why the region count changes under you.
Blocking calls with Table
The Table interface has two forms. coprocessorService(byte[] row) returns a CoprocessorRpcChannel connected to the region that contains row. The javadoc is explicit that the row does not have to exist and is not passed to the endpoint; it is only used to find the region. The range form, coprocessorService(Class<T> service, byte[] startKey, byte[] endKey, Batch.Call<T,R> callable), calls callable once per region, from the region containing startKey up to and including the region containing endKey. A null key means the first or last region. It returns a Map<byte[],R> keyed by region name.
import com.example.hbase.cp.TenantStatsProtos.*;
import com.google.protobuf.ByteString;
import org.apache.hadoop.hbase.client.*;
import org.apache.hadoop.hbase.ipc.CoprocessorRpcUtils;
import org.apache.hadoop.hbase.ipc.ServerRpcController;
import org.apache.hadoop.hbase.util.Bytes;
import java.util.Map;
static long countTenant(Table table, String tenant) throws Throwable {
byte[] start = Bytes.toBytes(tenant + "|");
byte[] stop = Bytes.toBytes(tenant + "}"); // '}' sorts right after '|'
CountRequest req = CountRequest.newBuilder()
.setStartRow(ByteString.copyFrom(start))
.setStopRow(ByteString.copyFrom(stop))
.build();
Map<byte[], CountResponse> perRegion = table.coprocessorService(
TenantStatsService.class, start, stop,
stub -> {
ServerRpcController controller = new ServerRpcController();
CoprocessorRpcUtils.BlockingRpcCallback<CountResponse> done =
new CoprocessorRpcUtils.BlockingRpcCallback<>();
stub.count(controller, req, done);
CountResponse r = done.get();
controller.checkFailed(); // rethrows the endpoint's IOException
return r;
});
long rows = 0;
for (CountResponse r : perRegion.values()) rows += r.getRows();
return rows;
}The checkFailed() line is easy to forget. An endpoint reports an error by setting it on the controller, and done.get() then returns null rather than throwing. Without the check, you get a NullPointerException in your reducer, or worse, a null silently skipped. The overload with a Batch.Callback<R> streams each region's result to update(region, row, result) as it arrives instead of building a map. Use it when results are large or when you want to reduce as you go; note that the callback can run on several threads at once.
For a single region, wrap the channel in a blocking stub: TenantStatsService.newBlockingStub(table.coprocessorService(row)). This suits a per-entity call, for example "compute this user's aggregate", where the row key is the entity.
Async calls with AsyncTable
AsyncTable has the same two shapes, without blocking threads. The single-row form returns a CompletableFuture<R>. The range form takes a stub maker, a ServiceCaller and a CoprocessorCallback, and returns a builder that you finish with fromRow, toRow and execute().
AsyncConnection conn = ConnectionFactory.createAsyncConnection(conf).get();
AsyncTable<AdvancedScanResultConsumer> t = conn.getTable(TableName.valueOf("events"));
LongAdder rows = new LongAdder();
CompletableFuture<Long> total = new CompletableFuture<>();
t.<TenantStatsService.Stub, CountResponse>coprocessorService(
TenantStatsService::newStub,
(stub, controller, done) -> stub.count(controller, req, done),
new AsyncTable.CoprocessorCallback<CountResponse>() {
public void onRegionComplete(RegionInfo region, CountResponse resp) { rows.add(resp.getRows()); }
public void onRegionError(RegionInfo region, Throwable err) { total.completeExceptionally(err); }
public void onComplete() { total.complete(rows.sum()); }
public void onError(Throwable err) { total.completeExceptionally(err); } // e.g. locate failed
})
.fromRow(start) // inclusive by default
.toRow(stop, true) // default is EXCLUSIVE -- unlike the blocking endKey
.execute();Two traps are hidden in that snippet. First, toRow(byte[]) defaults to exclusive, while the blocking API's endKey includes the region that contains it. If you port blocking code and pass a stop row that is the first key of a region, the async version skips that region. Here the endpoint clamps to the request's own range anyway, so including one extra region costs one wasted RPC, not a wrong answer. Second, the callbacks run on HBase's RPC threads. Keep them to counters and future completion, and never block or call HBase again inside them.
AggregationClient: the built-in endpoint
HBase ships one general-purpose endpoint, AggregateImplementation, with a client-side helper, org.apache.hadoop.hbase.client.coprocessor.AggregationClient, that does the fan-out and merge for row count, sum, min, max, average, standard deviation and median. The table must have the endpoint loaded. Every method except rowCount requires the scan to name exactly one column family; rowCount also accepts a scan with none.
try (AggregationClient agg = new AggregationClient(conf)) { // owns a Connection: close it
Scan scan = new Scan()
.withStartRow(Bytes.toBytes("t1|"))
.withStopRow(Bytes.toBytes("t1}"))
.addFamily(Bytes.toBytes("d"));
long n = agg.rowCount(TableName.valueOf("events"), new LongColumnInterpreter(), scan);
}The constructor you pick decides who owns the connection. new AggregationClient(conf) opens its own connection, which close() shuts down, so treat it as a long-lived singleton rather than creating one per request. new AggregationClient(connection) borrows yours. The no-argument constructor only works with the methods that take a Table. Values must be stored in the encoding the interpreter expects. LongColumnInterpreter reads 8-byte longs and returns null for any other length, so string-encoded numbers are silently skipped rather than reported as errors. These methods declare throws Throwable, so catch and wrap them at the boundary.
Worked example: counting one tenant's day
A table events has 40 regions and row keys tenant|yyyyMMddHH|eventId. The dashboard needs today's event count for tenant t1. A plain scan would stream every matching row to the client. With the endpoint, the client sends CountRequest{start=t1|2026100300, stop=t1|2026100400} over the range those keys cover.
Suppose t1 is a large tenant whose day spans two regions: one ends at t1|2026100314 and the next starts there. The blocking call sends two RPCs and gets a map of two entries, for example 4,180,226 rows and 3,902,117 rows, and returns 8,082,343. The network carries two small requests and two answers of a few bytes each instead of eight million rows.
Next week the first region splits at t1|2026100307. The same call now sends three RPCs, the map has three entries, and the total is unchanged. Code that assumed one entry per tenant, or cached results by region name, breaks at that point. Always reduce over values() and never key anything durable on region names.
Failure modes
- Endpoint not loaded. Calling a service that is not registered on the region fails with a server-side exception on every region, not an empty result. Check table descriptors in CI, not in production.
- Partial failure. One region throwing aborts the blocking range call, but other regions may already have done their work. Endpoints that write must be idempotent so the client can safely retry the whole call.
- Region moves mid-call. The client retries regions that move or split under its normal retry settings, which stretches latency. Set
hbase.rpc.timeoutandhbase.client.operation.timeoutfor the slowest region you are willing to wait for, not the average one. - Unclamped ranges. An endpoint that scans its whole region double-counts data from the edge regions. Always send the range and clamp to it.
- Thread exhaustion. The blocking range call fans out on the connection's batch pool. Several concurrent calls on a 400-region table can starve ordinary gets and puts. Limit concurrency, or move to
AsyncTable. - Proto drift. A client compiled against a newer
.protothan the deployed endpoint sends fields the server ignores. Version the service and roll servers out first.
Trade-offs
Endpoints win when the answer is much smaller than the data and the computation is simple: counts, sums, sketches, top-k per region. They lose when the per-region work is heavy, because it runs inside the RegionServer with no sandbox and competes with serving traffic. They also add deployment coupling: the client now depends on a jar inside your servers. A server-side filter is a lighter way to push predicates down when you still want the rows. Phoenix, Spark or a MapReduce job are better when the computation is large, and coprocessor patterns compares the shapes. For client wiring in general, including connection lifecycle, see HBase client libraries.
What to do next
- Put the service
.protoin its own versioned jar shared by endpoint and clients; only add optional fields. - Always send the logical row range in the request and clamp the endpoint's scan to it; add a test with a range that starts and ends mid-region.
- In blocking code, call
controller.checkFailed()after everydone.get(), and reduce over the map's values, never its keys. - In async code, set
toRow(stop, inclusive)explicitly, and keepCoprocessorCallbackmethods non-blocking. - Make writing endpoints idempotent so a whole-call retry after partial failure is safe.
- Hold one
AggregationClientper process (or pass it yourConnection) and close it on shutdown. - Load-test the call at your expected region count, then at ten times that, and set RPC and operation timeouts from the results.