On this page
TiKV is a distributed, transactional key-value database, and PD (Placement Driver) is its brain. TiDB, the MySQL-compatible SQL database, is built on them. Loam uses TiKV and PD directly, without TiDB, as the store for Loam Live, as a metadata backend for the engine, and as the target store for durable execution and jobs.
| Repositories | tikv/tikv, tikv/pd, tikv/client-rust |
| Docs | tikv.org/docs |
| License | Apache-2.0 (all three) |
| Versions checked | TiKV and PD v8.5.8; tikv-client 0.4.0 on crates.io; Loam pins a fork of client-rust by revision |
| In Loam | TiKV client layer and metastore Loam Live |
How TiKV works
Regions and Raft. TiKV splits the key space into contiguous ranges called regions, each a few tens of megabytes. Each region is its own Raft group, usually with three replicas on different stores. This is multi-Raft: thousands of small consensus groups rather than one big one, so leadership, load and recovery are spread across the cluster. A region that grows is split; small neighbours are merged. Data on each store sits in RocksDB.
PD. The Placement Driver holds the cluster's metadata (which regions exist, where their replicas and leaders are) and schedules: it moves replicas and leaders to balance load, and splits hot regions. It also runs the timestamp oracle (TSO), which hands out globally increasing timestamps: a physical millisecond part shifted left by 18 bits, plus a logical counter. PD is itself replicated with an embedded etcd.
Percolator transactions. TiKV's transactions follow Google's Percolator design. Each transaction gets a start timestamp from the TSO and reads a snapshot at it (MVCC keeps multiple versions per key). To commit, the client prewrites every key it wrote, placing a lock on each, with one key chosen as the primary. If any prewrite finds a newer committed write or another lock, the transaction conflicts. Otherwise the client gets a commit timestamp and commits the primary; that single write is the commit point. Secondary keys are committed afterwards, or resolved lazily by any reader that finds their locks by checking the primary. Two optimizations cut latency: async commit (the transaction is committed once all prewrites succeed, with the commit timestamp computed from them) and one-phase commit for transactions that touch a single region.
The result is snapshot isolation across regions, with optimistic and pessimistic modes. Pessimistic transactions take locks as they read (get_for_update), so contended keys queue instead of aborting in a loop.
API v2 and keyspaces
TiKV's original API let RawKV (plain key-value) and TxnKV (transactional) users share a cluster only with care. API v2 changes the key encoding: every key starts with a mode byte (r for raw, x for transactional) and a 3-byte keyspace id. RawKV gains MVCC too, which is what makes raw change capture possible. It is enabled with storage.api-version = 2 and enable-ttl = true, and a cluster should start on it: switching smoothly is only possible from an empty or TiDB-only cluster.
Keyspaces are why Loam uses TiKV at all. One cluster can hold up to 224 of them, isolated by prefix. PD creates them over an HTTP API, splits regions at their bounds, and moves them through states (enabled, disabled, archived, tombstone). A transactional client and a TiDB running in keyspace mode can share one cluster safely, each in its own keyspace; a spike confirmed that a key written through the Rust client was invisible from TiDB's keyspace and the default one.
Why we chose it
Loam's retrieval engine keeps everything in object storage, and that is the wrong shape for an application database or a job queue: they need millisecond transactions over small mutable records. We needed a store that is:
- Transactional across keys and ranges, with a clear isolation level, because a reactive database's mutation is "read some documents, write some documents, atomically".
- Multi-tenant at the storage layer, so millions of apps do not each need a cluster.
- Scale-out and self-hostable under a permissive license, since the same store must run in the managed cloud and in a customer's cluster.
- Proven, because it holds application data.
TiKV is Apache-2.0, a CNCF graduated project, runs TiDB's production deployments, and API v2 keyspaces give tenancy natively. FoundationDB was considered earlier for metadata and dropped. Postgres and DynamoDB remain metastore backends for deployments that already run them. For Live, TiKV's range scans over order-preserving keys map directly onto index queries, and PD's TSO gives Live one global clock to version query results.
How Loam uses it
Loam runs PD and TiKV unmodified from official releases: tiup playground in development and CI, and TiDB Operator v2 on Kubernetes (its component-group resources allow a cluster of PD and TiKV only).
operon-tikv, the shared client layer: keyspace bootstrap through PD's API, the TSO clock, a transaction runner that classifies errors and retries, fault-injection hooks for tests, commit tokens for unknown outcomes, the order-preserving tuple codec, and an MVCC garbage-collection loop. That last one matters: TiDB advances the GC safe point for its own keyspace, but nobody advances it for a plain transactional keyspace, so old versions would pile up forever unless Loam does it.operon-meta-tikv, the engine's metadata contract on TiKV. Each call is one transaction: optimistic for single-record writes (leases, pointer swaps), pessimistic with orderedget_for_updateon partition heads for WAL commits, so hot heads queue. Clock stamps come from the TSO, one monotonic clock with no hot row. After an error during commit, a commit-token row read at a fresh timestamp tells the caller whether the transaction landed, so outcomes are never unknown. It passes the same conformance suite and linearizability checks as the default openraft backend, plus its own fault matrix.- Loam Live: each mutation is one optimistic transaction, retried on conflict; document reads by id are promoted into the lock set to make read-modify-write serializable; every mutation writes a commit-journal entry in its own transaction for invalidation (platform part 3).
Both the metastore and Live currently commit with classic two-phase commit. Async commit and 1PC cut commit latency by 30 to 50% in our spike, and will be switched back on per component once their linearizability histories pass with them.
The change feed
TiKV has a change-data-capture service, used by TiCDC to stream TiDB's changes. Our spike subscribed to it from Rust for a plain transactional keyspace and received prewrites, commits, one-phase commits, deletes and resolved timestamps. So it works, with conditions we had to discover:
- The subscriber must declare
kv_api = TiDB, even for non-TiDB keys;TxnKVis refused with a misleading version error. - The client crate's generated CDC and PD protobuf code is private, so the subscriber needed its own stubs.
- The subscriber must track regions itself: scan the keyspace's regions, open one feed per store, and re-register after splits, merges and leader moves.
- Ordered delivery waits about a second, the interval at which resolved timestamps advance.
For Live's invalidation, one second is too slow, so Live writes its own commit journal inside each transaction. The change feed stays a validated secondary path, a candidate for the bridge from Live tables into search, where a second of lag is acceptable. (For RawKV, TiKV's separate TiKV-CDC project exists but is marked experimental and sees little activity.)
The client fork
The Rust client, tikv-client 0.4.0, says in its README that it is "not suitable for production use - APIs are not yet stable". It is maintained (the last functional commit was in September 2026), but our spike found problems we could not ship with:
- A PD stall killed the client's TSO stream permanently, so every later transaction failed until the client was rebuilt.
- Async commit unwrapped a value it should not and never set
max_commit_ts. - A reader that met a crashed async-commit writer's lock could not resolve it; only GC could, so the reader blocked.
- The generated CDC, PD and keyspace protobuf modules were private.
- The crates.io release pulled in old tonic and prost versions that fail our advisory checks.
So Loam depends on dina-kar/client-rust, branch loam, pinned by revision. It carries TSO stream reconnection, public proto modules, read-path resolution of async-commit locks, and max_commit_ts. Each fix is drafted as an upstream pull request. The pin moves only in a pull request that reruns the TiKV, metastore and Live suites.
Limits
- Not object storage. TiKV keeps data on local disks with Raft. Loam's design makes continuous log backup to the bucket (TiKV's backup-stream component plus BR snapshots) mandatory for every Live cluster; whether BR's point-in-time restore covers a non-TiDB keyspace is still to be verified. TiDB's next-generation kernel, which keeps storage in object storage, has a matching TiKV engine that is not public.
- Memory. One PD, one TiKV and one TiDB peaked at about 3.2 GB in our playground, most of it TiKV sizing its block cache from host RAM. Development configs cap it explicitly.
- Region overhead. Each keyspace costs at least one region with three replicas, which is why small Live apps share keyspaces.
- Hot keys. TiKV splits regions by load, but cannot split one key. A busy partition head or journal shard is a hot key by design, which pessimistic locking and sharding mitigate.
- Latency. Every transaction pays a TSO fetch plus prewrite and commit round trips: a few milliseconds per commit in our spike on a loaded machine.
- The client is pre-1.0. We carry a fork until the fixes land upstream.
The next post moves back to the bucket itself, and the object store Loam defaults to: RustFS.