Apache Flink and Apache HBase meet in two places. HBase is a sink: Flink computes a result, such as a running aggregate, a feature vector or the latest state of a device, and writes it to HBase so that online services can read it by key in milliseconds. HBase is also a lookup source: a stream of events is enriched with reference data stored in HBase, a customer profile or a product record, using a lookup join. The Flink HBase connector supports both, through Flink SQL and the Table API.
This article explains how the connector maps Flink rows onto HBase's column-family model, what its upsert sink guarantees across failures, how lookup joins and their cache behave, and how to tune and operate the pair. A sibling article, HBase streaming ingest patterns, covers idempotent row design and counters in depth; here the focus is the connector and the join. Option names below were checked against the connector documentation for current Flink; the connector is released separately from Flink, so match its version to yours using the compatibility table on its page.
How Flink rows map onto HBase
Start with the data model mismatch. A Flink table is a set of typed columns. An HBase table is a sorted map from row key to column families, each holding any number of qualifiers whose values are bytes. The connector bridges this with one convention: the single atomic-typed column is the row key, and each column family is declared as a ROW type whose fields are the qualifiers. Values are encoded with HBase's Bytes utility: strings as UTF-8, integers as fixed-width big-endian, timestamps as milliseconds since epoch in a long. ARRAY, MAP and nested ROW values inside a family are not supported.
That encoding matters for interoperability. If another application wrote the HBase table with, say, integers stored as decimal strings, Flink will misread them; declare such columns as STRING and cast in SQL. And because the row key is encoded the same way, a key declared as INT must have been written as four big-endian bytes for lookups to match.
Declaring an HBase table
A table declaration names the HBase table, the ZooKeeper quorum used to find it, and the family layout:
CREATE TABLE customers (
rowkey STRING,
info ROW<name STRING, tier STRING, country STRING>,
limits ROW<credit_limit DECIMAL(12,2), updated_at TIMESTAMP(3)>,
PRIMARY KEY (rowkey) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
'table-name' = 'crm:customers',
'zookeeper.quorum' = 'zk1:2181,zk2:2181,zk3:2181',
'zookeeper.znode.parent' = '/hbase',
'lookup.cache' = 'PARTIAL',
'lookup.partial-cache.max-rows' = '500000',
'lookup.partial-cache.expire-after-write' = '10 min',
'lookup.async' = 'true',
'lookup.max-retries' = '3'
);Notes on each choice. table-name accepts the namespace:table form. The connector identifier depends on the connector release: the documentation for current stable Flink uses hbase-2.2, while the connector's development branch documents hbase-2.6, so read the page that matches the jar you deploy. Kerberos and other client settings are passed through with the properties.* prefix, for example properties.hbase.security.authentication. The connector jar is not part of the Flink distribution; add it to the job or the cluster's lib directory.
The PRIMARY KEY must be the row key, and it is NOT ENFORCED because Flink cannot check uniqueness in an external store; it tells the planner the table is keyed, which is what makes upsert writes and key lookups possible.
Writing: the upsert sink
The HBase sink works only in upsert mode. The planner hands it a changelog keyed by the row key: inserts and updates become HBase Put mutations for that row, and deletes become Delete mutations. Because the sink is keyed, the planner does not need to send the old version of an updated row, only the new one, which keeps traffic proportional to changes.
CREATE TABLE customer_stats (
rowkey STRING,
s ROW<orders BIGINT, revenue DECIMAL(14,2), last_order TIMESTAMP(3)>,
PRIMARY KEY (rowkey) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
'table-name' = 'crm:customer_stats',
'zookeeper.quorum' = 'zk1:2181,zk2:2181,zk3:2181',
'sink.buffer-flush.max-size' = '4mb',
'sink.buffer-flush.max-rows' = '2000',
'sink.buffer-flush.interval' = '1s',
'sink.parallelism' = '8'
);
INSERT INTO customer_stats
SELECT customer_id,
ROW(COUNT(*), SUM(amount), MAX(order_ts))
FROM orders
GROUP BY customer_id;The sink buffers mutations and sends them in batches when any of three thresholds is reached: buffered bytes (sink.buffer-flush.max-size, default 2mb), buffered rows (sink.buffer-flush.max-rows, default 1000) or elapsed time (sink.buffer-flush.interval, default 1s). Pending mutations are also handled at every checkpoint, so a completed checkpoint means everything before it is either in HBase or in the checkpointed sink state.
What this guarantees. Delivery is at least once. After a failure, Flink restores state from the last checkpoint and replays input from there, so some mutations are written twice. For an upsert keyed by row key that is harmless: writing the same final value for a key twice leaves the same row. That is why the pattern works without two-phase commit. It stops being harmless when the written value is not a pure function of the key's state, for example if you write a client-side increment or derive the HBase cell timestamp from processing time. Keep sink values deterministic, and let Flink's checkpointed state hold counters, as the ingest article explains.
Two options change the row semantics. sink.ignore-null-value (default false) skips null fields instead of writing them, which turns a full-row upsert into a partial update and lets two jobs write different qualifiers of the same row. null-string-literal sets the string stored for a null string, default null; readers in other languages must agree on it.
Reading: lookup joins and the cache
A lookup join enriches each event with the current value of a key in an external table. In SQL it is a temporal join on processing time: the probe table needs a processing-time attribute, and the join uses FOR SYSTEM_TIME AS OF.
CREATE TABLE orders (
order_id STRING,
customer_id STRING,
amount DECIMAL(12,2),
order_ts TIMESTAMP(3),
proc_time AS PROCTIME()
) WITH ('connector' = 'kafka', 'topic' = 'orders', ...);
SELECT o.order_id, o.amount, c.info.tier, c.info.country,
o.amount > c.limits.credit_limit AS over_limit
FROM orders AS o
LEFT JOIN customers FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.customer_id = c.rowkey;For each order, the operator issues a Get for the row key. The join condition must be an equality on the row key; a join on a qualifier would need a scan, which is not how lookups work. Use LEFT JOIN unless you want orders for unknown customers silently dropped.
The cache. With lookup.cache set to PARTIAL, each parallel subtask keeps a bounded cache of recent results, sized by lookup.partial-cache.max-rows and aged by expire-after-write or expire-after-access. Misses are cached too by default (lookup.partial-cache.caching-missing-key is true), which saves HBase from repeated lookups of absent keys but means a customer created a minute ago may still read as missing until the entry expires. The cache trades freshness for load: with a 10-minute write expiry, an updated credit limit can take up to 10 minutes to affect the join.
Async lookups. lookup.async issues lookups without blocking the operator thread, so one subtask can have many Gets in flight. It helps when HBase latency, not CPU, bounds throughput. The documentation limits async lookups to the HBase 2.x connector (hbase-2.2 in the stable docs), so check that the identifier you deploy supports it. The general mechanics of async enrichment are in Flink async I/O.
Note that a processing-time lookup join is not reproducible. Replaying the same orders tomorrow joins them with tomorrow's customer data. If you need the value as of the event time, model the dimension as a changelog stream and use an event-time temporal join instead.
Worked example: enriching orders
Put the two halves together. Orders arrive at 20,000 per second across 8 Kafka partitions. Each must be enriched with the customer's tier and credit limit, and per-customer totals are written back for a dashboard that reads HBase by customer ID.
Lookup load. Suppose 2 million active customers with a skewed distribution where the hottest 300,000 cover 90% of orders. With parallelism 8, Flink does not partition the lookup by key unless you hash the stream first, so each subtask may see every hot customer. A 500,000-row cache per subtask holds the hot set. Cold-key misses then cost about 2,000 Gets per second, but expiry adds refreshes: 300,000 hot keys cached in each of 8 subtasks, re-fetched every 600 seconds, is another 4,000 per second. About 6,000 Gets per second in total, against 20,000 with no cache, is well within a few region servers' capacity. A longer expiry cuts the refresh share at the cost of staler joins.
Sink load. The aggregate emits an update per order, 20,000 upserts per second, but many target the same hot customers. Flink's mini-batch aggregation (table.exec.mini-batch.enabled with a latency such as 1 s) collapses repeated updates per key before they leave the operator, often reducing writes several-fold. At a 1-second flush interval and 2,000-row batches, each of 8 sink subtasks sends a few batches per second.
Row-key design. Customer IDs that are sequential would send all new customers to one region. Salt or hash the key on write (for example a two-character hash prefix), and apply the same transformation in the lookup's ON clause so the Gets still match.
When to use the DataStream API
SQL covers most needs. Drop to the DataStream API when you need control the connector does not expose: writing to several tables from one stream, choosing explicit cell timestamps, using Increment or check-and-mutate operations, or sharing an HBase connection pool across operators. The usual pattern is a RichAsyncFunction for enrichment that wraps the HBase 2 AsyncConnection, and a custom sink that holds a BufferedMutator and flushes it at every checkpoint, which gives the same at-least-once guarantee as the connector. Flink 2.0 removed the old SinkFunction API, so write it against the Sink V2 interfaces, whose writer's flush runs before each checkpoint.
// Flush-on-checkpoint sink skeleton (Java, Flink Sink V2, HBase 2 client)
public class StatsSink implements Sink<Stat> {
@Override public SinkWriter<Stat> createWriter(WriterInitContext ctx) throws IOException {
return new StatsWriter();
}
}
class StatsWriter implements SinkWriter<Stat> {
private final Connection conn;
private final BufferedMutator mutator;
StatsWriter() throws IOException {
conn = ConnectionFactory.createConnection(HBaseConfiguration.create());
mutator = conn.getBufferedMutator(
new BufferedMutatorParams(TableName.valueOf("crm:customer_stats"))
.writeBufferSize(4 << 20));
}
@Override public void write(Stat s, Context ctx) throws IOException {
Put put = new Put(Bytes.toBytes(s.customerId));
put.addColumn(CF, ORDERS, Bytes.toBytes(s.orders)); // value, not delta
mutator.mutate(put);
}
@Override public void flush(boolean endOfInput) throws IOException {
mutator.flush(); // checkpoint completes only after HBase acks
}
@Override public void close() throws Exception { mutator.close(); conn.close(); }
}The critical line is the flush. Without it, a checkpoint can complete while mutations sit in the client buffer; a crash then loses them, and Flink will not replay them because the checkpoint says they were processed. Check the exact interface signatures against the Flink version you build with.
Failure modes
| Failure | What you see | Response |
|---|---|---|
| Region server slow or splitting | Sink back-pressure, checkpoint duration climbs, lookup latency spikes | Watch checkpoint alignment time; presplit tables; keep region count balanced |
| Region server dies | Retries, then task failure and restart from checkpoint | At-least-once replay rewrites the same upserts; make sure values are deterministic |
| Stale enrichment | Joins use old tier or limit after a change | Lower expire-after-write, or switch to an event-time temporal join on a changelog |
| Missing-key cache | New customers look unknown for minutes | Disable caching-missing-key, or shorten expiry |
| Encoding mismatch | Garbage numbers, lookups that never match | Align types with whoever writes the table; declare as STRING and cast |
| Hot region | One sink subtask far slower than others | Salt row keys; check that the key distribution is not dominated by a few values |
| Delete then replay | A deleted row reappears after recovery | Replays re-apply the full changelog in order from the checkpoint; check for out-of-order writers outside Flink |
The single best health signal is checkpoint duration. HBase trouble shows up there first, because a slow flush delays the barrier. See Flink checkpoint design for setting intervals and timeouts that tolerate an HBase region move.
Trade-offs
HBase is a strong partner for Flink when results must be served by key at low latency and high volume, and you already run Hadoop-ecosystem storage. It is a weaker fit when reads need secondary indexes or ad hoc filters, in which case a search engine or an OLAP store serves better, and when the dimension data is small enough to broadcast into Flink state, which removes the network hop entirely. Compared with keeping the dimension in Flink state via a changelog join, an HBase lookup uses less job memory and lets other systems update the data, at the cost of freshness controlled by the cache and a dependency on HBase availability for every event. For telemetry-shaped workloads, the table layouts in HBase for IoT pair well with this job design.
What to do next
- Declare one existing HBase table in Flink SQL and run a
SELECTlookup join against a test stream; verify a known key returns the expected values. - Confirm the connector identifier and options against the documentation page for the exact connector version you deploy.
- Turn on
PARTIALcaching, measure hit rate and HBase request rate, and pick an expiry that matches how stale the business can tolerate. - Make every sink value a deterministic function of key state, then kill a TaskManager mid-run and check that HBase rows end up identical to a clean run.
- Alert on checkpoint duration and sink back-pressure, and presplit result tables before launch.
- If you write a custom sink, add a flush in the checkpoint hook and a test that crashes between buffer and flush.