Issue 44619: data index inconsistent after DDL owner change
An index built in ingest mode while the owner resigned twice, followed by an upgrade, finished with under half the rows indexed. Severity critical against the 7.1 support line.
PingCAP has written down nearly every substantial change to TiDB and TiKV since 2018, in 136 dated design documents and 46 accepted RFCs, and has shipped a release note for every version since October 2017. Read in order, that record shows the binding constraint moving three times, in a sequence most shared-database teams will repeat, and it shows which of the grand redesigns actually reached users.
One transactional database, many applications, a decade of its own engineering paperwork. The question this guide answers is not how the system works but in what order its constraints arrived, and what that order predicts for anyone building the same thing.
The problem, stated without any product in it: you operate one strongly consistent, horizontally scaled transactional store, several teams write to it, and you cannot stop it to rebuild it. Over ten years, what actually forces redesign? TiDB is an unusually good place to answer that, because PingCAP runs a design-document process that requires a dated file in the repository before substantial work starts, and states the bar in writing: a document "is accepted or rejected when at least two committers reach consensus". The result is a decade of decisions with dates on them, next to a release-note corpus of 208 files that records what shipped, when, and what broke afterwards.
Three answers come out of that record, and the first one is the surprise. The most ambitious storage redesign of the second half of the decade did not reach users as an architecture. RFC 0082 proposed replacing the 96 MiB region with regions up to "10GiB (before compression)" as "the first step that we try to support PiB scale cluster", and RFC 0093 followed it by giving every region its own RocksDB instance. That engine was introduced as experimental in v6.6.0 in February 2023, promising to "expand the storage capacity of the cluster from TB to PB", was last mentioned in v7.4.0 in October 2023 as "part of the near-term GA of the architecture", and appears in no release note after that, through v8.5.8 of 27 August 2026. What users did get is in the shipped source: the default split size moved from 96 MB to 256 MB "in version >= 8.3.0", August 2024. The redesign landed as a constant.
The second answer is a repeated move rather than a single decision. Once quorum replication is paying for durability, every local write that duplicates it becomes negotiable, and the record shows that cashed in four times: pessimistic locks no longer replicated, a key-value write-ahead log no longer written, a raft log no longer waited on, and flow control taken out of the storage engine. Each deletion is paired with a compensation, and the compensations are where the complexity now lives. The third answer is where the work went instead. Resource governance is absent from the design record before 2021 and dominates it afterwards, from request units and keyspaces to runaway-query kills, background-task demotion and, in 2026, a per-statement limit on in-flight requests to one store. The cluster stopped being a scaling problem and became a neighbourhood.
This guide covers what the public repositories of one company say about one system between October 2017 and August 2026: design documents, accepted and rejected RFCs, release notes, shipped configuration defaults and severity-critical bug reports. It does not cover TiDB Cloud internals, which are not in the open repositories; it does not compare TiDB with its competitors; and it contains no engineering-blog posts, conference talks or customer incident reports, because this session's network reached code hosts only. Where the public record stops, the guide says so rather than filling the gap.
Four planes, one of which did not exist for the first five years. Every box below is attributable to a document or a source file, and the late arrival of the fourth plane is the main architectural fact of the decade.
The serving shape is the one PingCAP has described since the beginning and it has not changed. A stateless SQL tier parses and plans; a control plane called PD hands out timestamps and moves data; a row store replicates each key range as its own Raft group and runs two-phase commit over it. The TiKV readme names the lineage directly: the design "is inspired by some great distributed systems from Google, such as BigTable, Spanner, and Percolator", the transaction model "is similar to Google's Percolator with some performance improvements", and "Similar to Google's Spanner, TiKV supports externally-consistent distributed transactions". Spanner's own paper describes that property as "externally consistent reads and writes, and globally-consistent reads across the database at a timestamp". Hold on to that sentence; the last decision in this guide walks away from it.
Two components are variation points rather than parts of the core. Columnar replicas are attached to the same Raft groups as learners, which is how one system serves analytical queries without a separate pipeline, and they are optional: a cluster with no analytical workload simply has none. The storage engine under the row store is also a variation point, and an unresolved one. The original arrangement is one RocksDB instance per store holding every region's data; RFC 0093 proposed one instance per region, calling each a tablet, and the release notes carried it as experimental for four versions before going quiet. Treat the per-region engine as something to ask a vendor about rather than something to plan on.
The fourth plane is the interesting one, because it is young. Before 2021 there is nothing in the design record that bounds what one application may consume. The global resource control document of November 2022 says so plainly: "TiDB lack of global admission control limits the request from the SQL layer to the storage layer, and different users or applications may influence each other in the shared environment". What was added is a token bucket at the SQL entry, denominated in request units, and a fair queue in the storage read pool using "the dmclock algorithm to ensure the fairness between different resource groups". The same document names the commercial reason in one line: the mechanism "can be used by Resource Manage in the future of TiDB Cloud. Users can pay for the actual usage."
Resource groups carry a rate in request units, a priority and a burst flag, bound to a user, a session or a single statement by hint. The quota is counted across both the SQL and the storage layer, so one number covers work the planner cannot predict.
Specified in: resource control design, 2022-11-25; GA in v7.1.0
Priority queueing inside TiKV, with one rule worth copying: background work is pinned to the lowest priority whatever its group says, because "Background tasks (GC, compaction, statistics) use LOW group_priority regardless of their resource group's configured priority". A full queue rejects with a busy error rather than growing.
Specified in: queue fairness design doc
A keyspace is a prefix on every key plus a namespace in the metadata store, so many applications share one storage cluster while each SQL instance serves exactly one keyspace. The ceiling is explicit: "The max keyspace id is 16777216, no new keyspaces can be created after the maximum keyspace id has been reached."
Specified in: keyspace design, 2022-12-07
Four forks, each with the rejected option and the stated reason, taken from the documents that argued them. The last column is the one to steal: the condition under which the rejected option is the right one for your system.
tidb_txn_mode variable from "" to "pessimistic"".15".| Decision | Chosen | Rejected | Because | Evidence |
|---|---|---|---|---|
| Concurrency control | Pessimistic locking by default | Optimistic only | Migrated MySQL applications cannot retry whole transactions | v3.0.8 release notes, 2019-12-31 |
| Lock durability | Locks in leader memory, not replicated | Replicated, persisted locks | 20% of disk write bandwidth and half the lock latency, against a "low probability that a pessimistic lock will be lost" | RFC 0077, v6.0.0-DMR |
| Local persistence ordering | Apply at committed, not at persisted | Wait for the local raft write | Cloud disk latency jitter on one node slows every transaction that touches it | RFC 0112 |
| Write flow control | Throttle in the scheduler, disable the engine's own stall | Keep RocksDB write stall | A default delayed write rate of 16 MB/s "doesn't take real disk ability into account" and produces latency spikes | RFC 0067; contrast Pebble |
| Unit of scale | Larger regions, buckets for concurrency; shipped as a 256 MB default | Keeping 96 MiB regions | Region count, not region size, drove RPC fan-out, heartbeat CPU and slow tooling | RFC 0082, split-check config |
| Isolation granularity | Tenant and statement limits | Region-level fairness | Deprioritising a hot region stretches the tail of every query that touches it | RFC PR 121, closed 2026-04-14 |
| Cross-region writes | Two clusters, eventual consistency, last write wins | One cluster stretched across regions | A cross-region quorum puts inter-region latency in every commit | active-active design, 2025-11-05 |
PingCAP does not publish customer postmortems, so the failure record has to be read from the severity-critical bug reports and from the release notes that fixed them. Read that way it is remarkably concentrated: the serving path is not where the correctness incidents are.
Across all 208 release-note files there are 148 bullets fixing an inconsistency involving data or indexes. Classifying them by trigger puts the largest group, 36 bullets, in the online schema-change path, ahead of outbound replication (21) and the transaction layer itself (21). The storage and replication core, which is the part everyone worries about in review, accounts for four. That classification is mine, by keyword, and the boundaries are arguable; the imbalance is not.
Inside that group the trigger repeats with unusual consistency. An index build is a long-running background job that reads every row, sorts, and ingests the result, and it goes wrong when the ownership of the job changes underneath it: the owner node resigns, the owner is partitioned, the cluster is upgraded mid-build, the placement driver leader is killed, or a dependency is unavailable longer than a fixed retry budget. The serving path tolerates all five events by design. The maintenance path, which was distributed later and more quickly, did not.
REPLACE INTO executed between the import and merge stages of a unique-index build left index and row out of agreement. It was found by a schema-change fuzzer with failpoint injection, not by a user.None of these failures pages anyone. The symptom is error 8223 from ADMIN CHECK,
a command a user has to decide to run, and the
2021
design document that added permanent in-path assertions says why that matters:
inconsistency causes "Data Loss: Data written cannot be read", and "It is very difficult
for support engineers to locate data index inconsistencies when they occur". Two of the
four reports above were produced by PingCAP's own fuzzing and fault injection rather than
by a customer, which is a credit to the practice and a warning about the record: the
public corpus counts the bugs that were caught, not the ones that ran.
Every figure below comes from a design document, a release note or a shipped default, with the date it was published. Nothing here is a benchmark run for this guide.
| Metric | Value | At | Context | As of | Source |
|---|---|---|---|---|---|
| Default region split size | 256 MB | TiKV | 96 MB before v8.3.0; region max size is 1.5 times the split size | 2024-08 | split-check config |
| Region size proposed by the redesign | 10 GiB | TiKV | With 128 MiB logical buckets for concurrency; shipped bucket default is 50 MB and buckets are off by default | 2021 | RFC 0082 |
| Keys per region at the old default | 960,000 | TiKV | Split triggers on size or key count, whichever comes first | 2021 | RFC 0082 |
| Update-index latency, async commit on | 7.01 ms | TiDB v5.0 | Down from 12.04 ms, a 41.7% reduction, Sysbench at 64 threads | 2021-04 | v5.0.0 release notes |
| Clustered index effect | +39% | TiDB v5.0 | TPC-C tpmC, vendor's own test | 2021-04 | v5.0.0 release notes |
| Analytical engine against Greenplum 6.15 and Spark 3.1.1 | 2 to 3× | TiDB v5.0 | Vendor benchmark, same cluster resources, some queries 8 times faster | 2021-04 | v5.0.0 release notes |
| Raft log store replacement | -40% I/O | TiKV v6.1 | Also 10% less CPU, about 5% more foreground throughput, 20% lower tail latency "under certain loads" | 2022-06 | v6.1.0 release notes |
| In-memory pessimistic locks, projected | -20% write bandwidth | TiKV | And 50% lower lock latency, "according to preliminary tests on the TPC-C workload" | 2021 | RFC 0077 |
| In-memory pessimistic locks, as shipped | -10% latency | TiDB v6.0 | And 10% more queries per second, under the bottleneck the feature targets | 2022-04 | v6.0.0-DMR release notes |
| RocksDB default delayed write rate | 16 MB/s | TiKV | The number that made engine-level write stalls unacceptable | 2021 | RFC 0067 |
| Coprocessor request deadline | 60 s | TiKV | Per request, which is why it never bounds a long query | 2023-06 | runaway queries design |
| In-flight coprocessor requests per statement per store | 15 | TiDB | New system variable; zero disables the statement-level limit | 2026-07 | per-store limiter design |
| Buffered transaction size limit | 100 MB | TiDB | Default txn-total-size-limit, the constraint pipelined writes were designed to remove | 2024-01 | pipelined DML design |
| Maximum keyspaces per cluster | 16,777,216 | PD | Identifiers are never reused, so this is a lifetime budget rather than a concurrent one | 2022-12 | keyspace design |
| Leader lease | 9 s | TiKV | Shipped raft_store_max_leader_lease; bounds how long a stale leader can serve a local read | 2026-10 | raftstore config |
| Unpersisted apply gap | 1,024 logs | TiKV | Shipped max_apply_unpersisted_log_limit; the RFC's own default for the feature switch was zero, so whether it is on by default cannot be read from this file alone | 2026-10 | raftstore config |
| Dependency outage that truncates an index build | 810 s | TiDB | Observed during a rolling restart, before the fix | 2024-09 | issue 55808 |
| Stated horizontal scale of the row store | 100+ TB | TiKV | Readme claim; the per-region engine promised "from TB to PB" and has not reached general availability in any release note | 2026-10 | TiKV readme |
Measured by the vendor, not independently: every percentage in the table comes from PingCAP's own Sysbench, TPC-C or TPC-H runs, published without the harness. Treat them as the direction of an effect, not its size on your workload. Projected versus shipped: the in-memory lock rows are the same feature described twice, once as a design estimate of 20% write bandwidth and 50% lock latency, once as a shipped claim of 10% and 10%. When a vendor publishes both, the second number is the one to plan with. Derived by this guide: the 148 and 36 bullet counts, the 0 to 12 governance-document count, and the comparison between a 10 GiB proposal and a 256 MB default; the method is in the ledger and the classification is keyword-based. Nobody has published: the price of a request unit, how many clusters run the per-region engine, how often these inconsistencies reached customers, or any latency figure for the active-active design.
One pattern explains most of the middle of the decade, and one sequence explains the whole of it. Both transfer to systems that are nothing like this one.
The pattern first. A replicated system pays once for durability at the quorum, and then usually pays again on every node, because the components it is built from were designed to be durable on their own. Three of the decade's measured wins come from noticing that and deleting the second payment: locks that are not replicated, a write-ahead log that is not written, a local disk write that is not waited for. The reported gains are not marginal, with the log-store replacement alone claiming 40% less write traffic and 20% lower tail latency. The discipline that makes this safe is visible in every one of those documents: each states the event that would expose the deletion, and adds a compensation for it. RFC 0077 ships locks to peers before a voluntary leader transfer; RFC 0093 migrates every piece of store state into the raft log because without the log "every state can be lost"; RFC 0112 refuses the optimisation for merge preparation and bounds how far ahead applying may run. An architect copying the move should copy the enumeration, not just the deletion. If you cannot list the events that make the duplicate load-bearing, you are not removing redundancy, you are removing a safety net you have not found yet.
Now the sequence, which is the part worth taking to a review. For five years the binding constraint was correctness and compatibility, and the design record is full of SQL surface and transaction semantics. For the next three it was the local cost of the guarantee, and the record fills with storage layout and durability. Since 2021 it has been the neighbours: request units, keyspaces, runaway kills, background demotion, a memory arbitrator, and per-statement limits on fan-out. Nothing moves back. Each era's work remains, and each new constraint arrives while the previous one is still being paid for.
The last era is the one everybody builds last and needs first. Admission control was retrofitted into a serving path that was already mature: first design document late 2022, tenant-level half generally available in May 2023, statement-level half still arriving in 2026. Meanwhile the change that would have spared the second era, a genuinely bigger unit of scale, is the one that did not ship: absent from release notes after October 2023, with the workaround it was meant to retire still on by default.
If you are building a shared transactional system, the order of your own constraints is probably correctness, then the cost of the guarantee, then isolation between tenants. The first two you will be forced into. The third you can either design in while the system is young, or retrofit later into a serving path that has no idea who is asking. The claim this guide is prepared to defend, from the record above
Twenty-six sources, graded. Every one was read in this session from a
repository clone or a fetched page. The full ledger, with the quote supporting each claim
and the derived counts, is in sources.md beside this file.
An index built in ingest mode while the owner resigned twice, followed by an upgrade, finished with under half the rows indexed. Severity critical against the 7.1 support line.
A REPLACE INTO landing between the import and merge stages of a unique-index
build corrupted the index. Found by a schema-change fuzzer with failpoint injection,
reported against four support lines at once.
The distributed backfill framework plus an injected network partition at the owner produces the user-visible data inconsistency error. Backported into four release lines, which is the best available proxy for exposure.
A storage node down longer than the ingest retry budget during a rolling restart left key-value pairs unwritten, and the index silently short.
Proposes raising the region ceiling from 96 MiB to 10 GiB with logical buckets for concurrency, naming the three costs of region count: extra RPCs per transaction, per region driving cost on each node, and tooling that loops over every region.
Gives every region its own RocksDB instance to stop cross-region compaction and lock contention, and deletes the key-value write-ahead log as a consequence, moving every store state into the raft log.
Stops replicating pessimistic locks, on the stated ground that losing one is safe because the transaction fails and retries. Lists the compensations: ship locks before a voluntary leader transfer, move them on split and merge, bound the memory per region.
Raises the apply bound from the minimum of committed and persisted to committed, because cloud disk latency jitter on one node slows every distributed transaction that touches it. Excludes merge preparation and log compaction, and bounds the gap.
Turns off RocksDB's own write stall and throttles at the scheduler with a token bucket, because the engine's default delayed write rate of 16 MB/s ignores the real disk and converts pressure into latency spikes.
The same question answered the other way by another storage team: Pebble removes artificial write delays entirely, arguing they "increase write latencies without a clear benefit" in an open-loop system. TiKV's RFC cites this document and keeps throttling anyway, one layer up.
Replaces offline recovery tooling with a control-plane driven procedure, and states exactly which invariants an operator trades: committed writes, transaction atomicity, and the agreement between a table and its indexes.
Admits the absence of global admission control, then introduces request units as the single unit of account across the SQL and storage layers, with a fair queue in the storage read pool. Names the billing motive in the same document.
Tenancy as a key prefix plus a metadata namespace, with each SQL node bound to exactly one keyspace until restarted, and a hard ceiling of 16,777,216 never-reused identifiers.
Two or more clusters, each committing locally, replicated by change data capture, with conflicts resolved by last write wins and a soft-delete column so a delete keeps a comparable timestamp. States that global transactional consistency is not provided.
Explains why per-request deadlines do not bound a query, then adds rules on elapsed execution time with three escalating actions, the strongest being to kill the statement and watch for others like it.
Replaces after-the-fact memory accounting with a subscribe-then-allocate model, ranks the three problems it is solving (out of memory, heavy garbage collection, killed sessions) and reserves killing for genuine exhaustion risk.
A proposal to deprioritise hot regions, argued over from December 2025 and closed on 14 April 2026. Three reviewers reject it in favour of query-level control, on the ground that deprioritising a region stretches the tail of every query touching it.
The shipped answer to the region-size redesign, in a comment next to the constant: the default split size was 96 MB before v8.3.0 and 256 MB after. Buckets default to 50 MB and are off unless the newer raftstore is in use.
Hibernation of idle regions is still on by default in 2026, five years after the RFC that expected larger regions to make it unnecessary. The unpersisted apply gap is capped at 1,024 logs and the leader lease is nine seconds.
States the lineage (BigTable, Spanner, Percolator, Raft), the scale claim of more than 100 TB, externally consistent distributed transactions, and graduated status in the Cloud Native Computing Foundation.
One line in a patch release changes the concurrency model for every new cluster: the default transaction mode becomes pessimistic, six months after the feature shipped as an experiment.
The clearest set of published numbers in the corpus: async commit cutting update-index latency from 12.04 ms to 7.01 ms, clustered index worth 39% on TPC-C, and the analytical engine benchmarked against Greenplum and Spark.
The per-region storage engine enters as experimental with a promise to take clusters "from TB to PB", is described eight months later as approaching general availability, and is then absent from every subsequent release note through August 2026.
Resource control reaches general availability and is framed as consolidation: combine applications from different systems into one cluster, "reduce the number of clusters", and lay "the foundation for multi-tenancy".
Read from a mirrored copy in a public repository. The sentence that underwrites a decade of local optimisation: "Raft guarantees that committed entries are durable and will eventually be executed by all of the available state machines."
The system TiKV names as its model, offering "externally consistent reads and writes, and globally-consistent reads across the database at a timestamp". The 2025 active-active design declines to provide that property between clusters.
No engineering-blog posts, no conference talks, no customer case studies and no independent benchmarks: this session's network reached code hosts only, so anything published on a company blog, a documentation site, a conference archive or a paper repository was unreachable and has been left out rather than cited from memory. The practical effect is that every performance figure here is the vendor's own, and that the failure section is built from severity-critical bug reports instead of incident reviews with timelines and customer impact. Where a claim needed one of those, this guide says that nobody has published it.
Six rungs on a three-node cluster you can run on one machine. The interesting ones are the middle two, where the exercise stops being a tutorial and starts producing evidence about your own system.
Run one SQL node, one control-plane node and three storage nodes locally, load a few gigabytes, then list the regions and their leaders. Change the split size from the default and reload the same data.
Done when: you can state the region count for a given data size at two different split sizes. Teaches: the metadata unit and the replication unit are the same object, which is the constraint behind RFC 0082.
Start a steady insert and update workload, add an index on a column the workload touches, and run the table check afterwards. Time the check.
Done when: you can quote how long verification takes relative to the build. Teaches: detection of the most common correctness failure in this corpus is a command somebody has to decide to run.
Repeat the previous rung, and while the backfill runs, restart or isolate the node holding the schema-change ownership. Then check the table again, and repeat with the distributed execution framework enabled and disabled.
Done when: you have run both configurations at least three times and recorded any check failure. Teaches: why four separate severity-critical reports in this guide share one trigger, and whether your version still shares it.
Run a well-behaved transactional workload and a full-table scan from a second account. Measure the first workload's latency. Now put the scan in a resource group with a request-unit cap and a low priority, and add a runaway rule that kills statements over an elapsed-time threshold. Measure again.
Done when: you can show the protected workload's tail latency with and without the group. Teaches: what admission control is worth, and that a kill action is part of the design rather than an operator's improvisation.
With a write-heavy workload, measure disk write bandwidth and commit latency with in-memory pessimistic locks on and off. Then kill the leader of a hot region during the run and count the transactions that fail.
Done when: you have both numbers, the saving and the failure count. Teaches: the trade in RFC 0077 as a measurement rather than a claim, and whether your application can absorb the retries.
On a throwaway cluster, stop two of three replicas of some regions, then run the online unsafe recovery procedure. Afterwards, check the tables and compare against what you wrote.
Done when: the cluster serves again and you can name which rows and indexes are wrong. Teaches: that recovery from majority loss is a data-loss event with a documented shape, which is a different conversation with your risk owner than "we have three replicas".
The commands and filters that produced this page. They work on any company that keeps its design record in public, which is the point: the technique outlives this particular database.
git clone --depth 1 --filter=blob:none --sparse https://github.com/pingcap/tidb.git && git sparse-checkout set docsls docs/design | sort | awk -F- '{print $1}' | uniq -cgrep -l "Motivation" docs/design/202[4-6]-*.md | xargs grep -A4 -h "^## Motivation"is:pr is:closed is:unmergedrepo:tikv/rfcs is:pr is:closed is:unmerged rfcrepo:pingcap/tidb is:issue label:severity/critical "data inconsistency"grep -ri "inconsisten" releases/*.md | grep -iE "data|index" | wc -lgrep -ho "deprecated starting from[^.]*" releases/*.md | sort -ugrep -l "Partitioned Raft KV" releases/*.md | sort -Vgrep -rn "impl Default for Config" -A60 components/raftstore/srcgrep -rn "ReadableSize::mb\|Default value: 0" components/ | head -40pdftotext paper.pdf - | tr '\n' ' ' | grep -o "guarantees that[^.]*\."