datavisorethanqiu commented on issue #66787:
URL: https://github.com/apache/doris/issues/66787#issuecomment-5862173264
Design document for this proposal. The implementation is in #68531.
# GLOBAL_POINT index: cross-tablet pruning for point queries
## Metadata
| Item | Value |
|---|---|
| Status | Proposal with a working implementation, for review |
| Issue | apache/doris#66787 |
| Code | #68531 (9 commits, based on `master` `bdb165d693`) |
| Audience | Doris storage, planner and cloud-mode maintainers |
| Prerequisites | Rowset and segment layout, cloud mode (meta service, file
cache), Nereids planning, memtable-on-sink-node loads |
| Wording | MUST / SHOULD / MAY as in RFC 2119 |
## 1. Context
A point query `WHERE c = X` scans only the tablets that can hold `X`. FE
finds them in two ways today:
- partition pruning, when `c` is a partition column;
- bucket pruning (`HashDistributionPruner`), when `c` is the distribution
column.
For any other column, the query reaches every tablet the planner selected.
Every existing index (zone map, bloom filter, inverted, short key, primary key)
lives inside a segment. It can reject data only after the tablet is already
part of the scan. None of them tells FE which tablets to leave out.
In cloud mode, each tablet in the scan costs a rowset meta sync, a scanner,
segment opens, and remote reads of footers and index metadata. So the cost of a
point query grows with the number of tablets, even when only one row matches.
The workload behind this proposal: an event table of about 5 billion rows in
566 tablets, `DISTRIBUTED BY HASH(user_id)`, queried with `WHERE eventId =
'Y'`. `eventId` is nearly unique and is neither a partition nor a distribution
column. Each query fans out to 566 tablets. The inverted index on `eventId`
rejects 565 of them quickly, but only after all 566 are in the scan path.
Adding compute does not remove this fan-out.
## 2. Goals and non-goals
### Goals
| Goal | Measured by |
|---|---|
| G1. A point query on an indexed column scans only the tablets that may
contain the value | EXPLAIN `globalPointIndex: <col> -> kept/total`; kept
should be close to the number of tablets that really hold the value |
| G2. Query results never change because of the index, whatever its state |
Regression suite compares results with pruning on and off, for every write path
|
| G3. A table without the index plans and runs exactly as before | The prune
step returns before any RPC when the table has no GLOBAL_POINT index |
| G4. Write cost is about one hash and one insert per indexed value, with no
extra read IO | Code review of the write path; load benchmarks (to be measured)
|
### Non-goals
| Non-goal | Reason |
|---|---|
| Locating rows inside a tablet | That is the job of the existing segment
indexes |
| Useful pruning for low-cardinality columns | A value present in most
tablets cannot prune anything |
| Zero false positives | An exact structure costs far more (see
Alternatives); a false positive only means one extra tablet is scanned |
| Changing data placement (partitioning, bucketing) | Out of scope; the
index works on any layout |
| Plan-time pruning in local (non-cloud) mode | It relies on the cloud
per-query snapshot version (`setVisibleVersionForOlapScanNodes`). The scan-time
gate works in both modes |
## 3. Overview
**Core idea:** each rowset carries an immutable bloom filter per indexed
column. At planning time, FE asks the BEs which tablets have a bloom that may
contain the probed values, and drops the rest.
Design invariants:
1. **Only a definite miss drops data.** Every uncertainty (no descriptor,
unreadable or unusable bloom, version not synced, RPC error or timeout,
unsupported type or predicate) keeps the tablet or rowset. A broken index can
make a query slower. It MUST NOT change the result.
2. **The probe covers at least what the query reads.** The probed rowsets
MUST be the rowsets visible at the query's snapshot version. A BE whose view
has not reached that version keeps the tablet instead of probing an older view.
3. **Pruning only removes tablets.** The result is intersected with the
planner's selection. It never adds a tablet.
4. **Blooms are never mutated or merged across rowsets.** A rewrite
(compaction, schema change, BUILD INDEX) builds a new bloom from its own output.
## 4. Detailed design
The implementation is 9 commits. Each one maps to a section below.
| # | Commit | Section |
|---|---|---|
| 1 | Add the GLOBAL_POINT index type and its DDL | 4.1 |
| 2 | Add the GLOBAL_POINT bloom file format | 4.2 |
| 3 | Build GLOBAL_POINT blooms on the rowset write path | 4.3 |
| 4 | Skip rowsets at scan time with GLOBAL_POINT blooms | 4.6 |
| 5 | Build GLOBAL_POINT blooms on memtable-on-sink-node loads | 4.4 |
| 6 | Rebuild GLOBAL_POINT blooms in compaction, schema change and BUILD
INDEX | 4.5 |
| 7 | Prune tablets at planning time with GLOBAL_POINT indexes | 4.7 |
| 8 | Keep GLOBAL_POINT index files in the BE file cache | 4.8 |
| 9 | Add regression tests for the GLOBAL_POINT index | 7 |
### 4.1 DDL
```sql
CREATE TABLE t_evt (
id BIGINT NOT NULL,
ev INT NULL,
name VARCHAR(64) NULL,
INDEX idx_ev (ev) USING GLOBAL_POINT PROPERTIES ("fpp" = "0.01")
)
DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id) BUCKETS 32;
CREATE INDEX idx_name ON t_evt (name) USING GLOBAL_POINT;
ALTER TABLE t_evt ADD INDEX idx_name (name) USING GLOBAL_POINT;
BUILD INDEX idx_name ON t_evt; -- cloud mode: index historical rowsets
```
| Decision | Reason |
|---|---|
| A new index type, not a property of the bloom filter index | Different
granularity (rowset, not segment), storage (own file) and consumer (FE). No
coupling with existing index code paths |
| Types: integers, LARGEINT, date types, strings. Rejected: float, double,
decimal, nested, JSON, variant | The bloom is probed with the value's bytes, so
equality must be byte equality. A decimal's bytes depend on scale; a float has
several equal representations |
| One column per index, one GLOBAL_POINT index per column | Keeps the probe
a single-column membership test |
| `fpp` is the only property, in [1e-6, 0.5], default
`global_point_index_default_fpp` (0.01) | See 4.2 for what fpp means |
| ADD INDEX and DROP INDEX are light index changes in both modes | New
rowsets get blooms from then on. Old rowsets have no descriptor and are scanned
(invariant 1) until compaction or BUILD INDEX rewrites them |
| BUILD INDEX is supported in cloud mode only | Cloud BUILD INDEX is an
index-change compaction through the normal rowset writer, which builds the
blooms. The local index builder writes per-segment index files only, so FE
rejects it there |
### 4.2 Bloom file and sizing
One bloom per indexed column per rowset, covering all segments of the
rowset, in a file `<rowset_id>_<col_unique_id>.gpidx` next to the segment
files. The rowset meta carries a descriptor `ColumnPointIndexPB` (column, index
id, fpp, size, number of values, crc32) in the new field
`RowsetMetaPB.point_query_indexes` (1019; 118 in `RowsetMetaCloudPB`).
The file is a 48-byte header (magic `GPIX`, format version, hash strategy,
number of bits, number of values, body crc32) followed by the bloom body. The
reader checks all of it against the descriptor. Any mismatch returns no bloom,
never an error.
| Decision | Reason |
|---|---|
| Bitmap in its own file, only a descriptor in the rowset meta | The meta
service stores rowset meta in FDB, where a value is limited to about 100 KB.
Blooms are MB-sized per tablet |
| One bloom per rowset, not per segment | A query probes every bloom of a
tablet, and the false-positive rate grows with their number (below). Large
compaction output can have tens of segments |
| One bloom per rowset, not one mutable bloom per tablet | A rowset is
immutable once committed, so its bloom is too. A per-tablet aggregate would
have to be updated at every commit, and a query would have to know whether it
already covers the rowsets it reads; any lag is a wrong result. Per-rowset
blooms make stale state impossible: the query probes exactly the rowsets it
reads |
| Block-split bloom filter, murmur3 | Testing one value touches one 32-byte
block. Reuses `BloomFilter` |
**What fpp means.** A query keeps a tablet when any of its blooms answers
"maybe". With `B` blooms of fpp `p`, the tablet's false-positive rate is about
`B * p`. With `p` = 1% and `B` = 50, it is about 40%, and pruning is nearly
useless. So the DDL fpp is a budget for the whole tablet, and each bloom gets
`p = fpp / global_point_index_expected_blooms_per_tablet` (default 5).
**Sizing.** Compaction, schema change and BUILD INDEX know the exact row
count and size exactly. Loads do not, so they size from
`global_point_index_write_path_estimated_rows` (default 1M) and cap the bloom
at `global_point_index_max_write_path_bloom_bytes` (default 256 KB). The cap
exists because a load holds one bloom per indexed column for every tablet it
writes. An underestimate only raises the fpp. It never causes a false negative.
**Sizing self-check.** A bloom far too small for its row count saturates and
answers "maybe" for everything. Queries stay correct and nothing is marked
degraded, so this failure would be invisible. The warm-up sweep (4.8) checks
each descriptor's bits per key against its own fpp and reports undersized
blooms.
### 4.3 Write path
`BaseBetaRowsetWriter` owns one `GlobalPointIndexBuilder` per indexed column
and registers it in `RowsetWriterContext`. `VerticalSegmentWriter` hands the
builder to the column writer, which feeds each non-null value and records
nulls. Segments of one rowset may flush concurrently, so the builder is
thread-safe. `build()` writes the files after every segment is closed.
`BetaRowsetWriter` and `CloudRowsetWriter` both do this.
| Case | Behavior | Reason |
|---|---|---|
| Cloud packed files | The `.gpidx` packed location is recorded like the
segments' | Small files are often packed; without the location, reads fail |
| File cache | New `FileType::GLOBAL_POINT_INDEX_FILE`, cached like index
files under `enable_file_cache_write_index_file_only` | Blooms are read at
planning time |
| Recycler | Deletes the `.gpidx` files of a rowset | Otherwise they leak |
| MOW partial-update publish (`merge_rowset_meta`) | Drops the descriptors
of non-key columns; key columns keep theirs. The transient writer builds no
bloom | The appended segments have no bloom, and two blooms of different sizes
cannot be merged. Partial update supplies every key column, so appended rows
repeat keys that are already indexed |
| Row binlog rowsets | No bloom | They are read through the binlog wrapper
tables, which plan-time pruning skips |
### 4.4 Memtable-on-sink-node loads
With `enable_memtable_on_sink_node` (default true), the sender sees the
values and the receiver writes the rowset, and one rowset can be written by
several senders.
- Each sender (`BetaRowsetWriterV2`) builds partial blooms. After the last
flush and before `CLOSE_LOAD`, it sends each one to every replica with the new
opcode `ADD_POINT_QUERY_INDEX`: a `PGlobalPointIndexPart` in the header, and
the body as the attachment.
- The receiver builds no local bloom, because it never sees values and an
empty bloom would answer "absent" for everything. It ORs the parts byte by byte.
- **All or nothing:** if any sender that wrote a segment of the rowset sent
no part, the receiver drops every descriptor of that rowset. A partial bloom
would give false negatives.
- A part with a bad crc, a wrong size, or a different shape or index id from
the other senders removes that column's index. None of this fails the load.
- **Mixed versions:** `PStreamHeader.opcode` is a proto2 enum, so an older
receiver would read the unknown opcode as `APPEND_DATA` and append the bloom to
a segment. The receiver therefore advertises
`POpenLoadStreamResponse.supports_point_query_index`, and the sender only sends
to receivers that did. A rowset written by an older receiver simply has no
descriptor.
### 4.5 Rewrites
Compaction (local and cloud) and schema change pass the input row count as
`exact_row_count_for_global_point_index`. Their blooms are exact and not
capped, and a base compaction output tends to become the single large bloom of
the tablet.
Ordered data compaction links segment files instead of rewriting them. Its
output has no descriptor and is scanned until the next rewrite. It is logged.
A query probes every bloom of a tablet, so cloud cumulative compaction
raises the score of tablets with this index once their delta count exceeds
`global_point_index_expected_blooms_per_tablet`, by at most 4x. Other tablets
are not affected.
### 4.6 Scan-time gate
Before a rowset reader opens any segment, it tests the rowset's blooms
against the pushed-down EQ/IN predicates on indexed columns. On a definite
miss, it reads nothing from the rowset. The gate works in both modes, and also
covers types FE does not probe (CHAR, LARGEINT, dates) and tablets plan-time
pruning kept.
It walks the rowset's own descriptors, so it is not affected by schema
changes made after the rowset was written. Profile counters:
`GlobalPointIndexGateTime`, `GlobalPointIndexRowsetsFiltered`,
`GlobalPointIndexRowsetsProbed`, `GlobalPointIndexBytesRead`,
`GlobalPointIndexDegraded`. Switch: BE `enable_global_point_index_scan_gate`.
### 4.7 Plan-time pruning
FE (`GlobalPointIndexPruner`, called from
`OlapScanNode.applyGlobalPointIndexPrune`):
1. Runs in `ScanNode.setVisibleVersionForOlapScanNodes`, after the plan is
translated (the conjuncts are attached) and the snapshot versions are pinned
(invariant 2).
2. Takes the first GLOBAL_POINT column with an EQ or IN filter, and encodes
the values exactly as the BE inserts them: raw UTF-8 for strings, little-endian
two's complement for integers. Other types are not probed, because a wrong
encoding would be a false negative.
3. Groups the tablets by BE and sends one `prune_global_point_index` RPC per
BE, with each tablet's snapshot version.
4. Intersects the answer with the planner's tablets and rebuilds the scan
ranges only if something was pruned. EXPLAIN prints `globalPointIndex: <col> ->
kept/total tablets (probes=N, unchecked=M)`.
It is skipped for rollups, wrapper tables, TABLESAMPLE, short-circuit point
queries and the query cache.
BE (`be/src/cloud/cloud_global_point_index.cpp`):
- Uses cached tablet metadata only. A tablet that is not cached, or whose
cached rowsets have not reached the snapshot version, is kept. It does not sync
from the meta service on the planning path.
- Reads every (tablet, rowset) bloom concurrently on a shared pool
(`global_point_index_prune_io_max_threads`). The first hit of a tablet skips
its remaining reads.
- Answers with the tablets to keep, not the ones to drop, so an omission can
only cause extra scanning.
| Decision | Reason |
|---|---|
| FE asks the BEs, instead of each BE gating at scan time only | Only FE
decides which tablets are scanned. A tablet that enters the plan already pays
scan-range scheduling, fragment instances, rowset meta sync, scanners and
segment opens; the gate saves only the last part |
| One RPC per BE, not per tablet | The RPC count grows with the number of
BEs |
| Timeout (`global_point_index_prune_timeout_ms`, default 3000) keeps the
tablets of a BE that does not answer | Pruning is an optimization and must not
fail or stall a query. The default is generous for cold caches; with warm
caches a healthy cluster answers well within 500 ms |
| Switches: FE `enable_global_point_index_prune`, session variable of the
same name; `global_point_index_max_probe_values` (32) | Per-query A/B
comparison; long IN lists are not worth probing |
### 4.8 Keeping the files cached
After a BE restart or eviction, each probed tablet would pay remote reads on
the planning path. A master-only FE daemon (`GlobalPointIndexWarmUpDaemon`,
every 60 s) sends each alive BE the tablets of indexed tables that map to it:
- an initial warm-up when it sees a new BE process (its `lastStartTime`
changed): the BE syncs rowsets and caches every file;
- a repair sweep at most every `global_point_index_repair_interval_sec`
(600) afterwards, using cached rowsets only.
The BE (`warm_up_global_point_index` RPC) queues one task per tablet on its
own pool (`global_point_index_warmup_io_max_threads`) and answers at once. Each
task checks residency with `BlockFileCache::probe()`, which touches no LRU
state. Fully cached files are left alone. Cached blocks in the NORMAL or
DISPOSABLE queue are moved to the INDEX queue. Only missing files are read. The
same walk runs the sizing self-check. Results come back in the next response
and as BE metrics:
| Metric | Meaning |
|---|---|
| `global_point_index_warmup_{checked,resident,repaired,failed}_total` |
Files looked at, already cached, brought into the cache or re-queued, failed |
| `global_point_index_blooms_undersized_total` | Blooms far too small for
their fpp: pruning is lost for those rowsets |
| `global_point_index_blooms_empty_total` | Blooms with no value: a bloom
never fed, or an all-NULL column |
| `global_point_index_blooms_missing_total` | Non-empty rowsets without a
descriptor; normal after ADD INDEX until rewritten |
Blocks written at load time land in the NORMAL queue, like every other
written file today. The read paths and the sweep put them in the INDEX queue.
## 5. Capacity
Reference scale: N = 5e9 rows, T = 566 tablets, about 8.8M rows per tablet,
a near-unique indexed column.
| Quantity | Formula | Reference value |
|---|---|---|
| Lower bound for any exact cross-tablet index | `N * log2(T) / 8` bytes |
about 5.4 GB |
| Target per-tablet false-positive rate | `p_tablet = 1 / T` (about one
extra tablet per query) | 0.18% |
| Bits per key at `B` = 5 blooms per tablet | `-ln(p_tablet / B) / (ln
2)^2`, block-split filters need more | about 16.5, so about 18 MB per tablet,
about 10 GB in total |
| Worst case after rounding the bitmap up to a power of two | x2 | plan for
10–20 GB in object storage |
| Planning cost | one RPC per BE; one small file read per (tablet, rowset)
on a cache miss | to be measured with and without warm-up |
Reducing `B` matters more than reducing `p`: from 50 to 5 blooms per tablet
saves about 23% of storage, but divides the files read per probe by 10.
The INDEX queue of the file cache defaults to 5% of the capacity and SHOULD
be sized to hold the blooms (`index_percent` in `file_cache_path`).
## 6. Alternatives
| Alternative | What it is | Why not chosen |
|---|---|---|
| Keep the current behavior | Rely on the segment indexes | The fan-out and
its cost stay. This is the problem statement |
| Exact `value -> tablet` postings in the meta service (FDB) | One key per
row | About 5e9 writes and 200–300 GB in FDB, on the critical path of all
metadata. Exact, but the cost is in the wrong place |
| A router table maintained by the application | A Doris table `value ->
tablet`, then `SELECT ... TABLET(...)` | No engine change, exact. But it needs
dual writes with a freshness window, and tablet ids change after schema change,
restore or rebalancing. Not transparent to users |
| One transposed bloom per partition (`hash slot -> tablet bitmap`) | One
structure answers the probe in a few range reads, with no fan-out | A global
mutable aggregate: every commit updates it, and a query must know whether it
covers the latest rowsets. That is the correctness risk invariant 4 avoids.
Worth revisiting at much larger tablet counts |
| An in-memory prune directory in FE or a separate service | Blooms held in
memory | Gigabytes of memory, and a second copy of rowset visibility to keep
consistent |
| Per-segment blooms | Reuse the segment bloom filter index | `B` becomes
tens per tablet and pruning loses most of its power (4.2) |
## 7. Testing
| Level | Coverage | Result on `master` |
|---|---|---|
| FE unit tests (20) | DDL parsing and validation, light-change rules, probe
encoding against the BE byte layout, prune gating, warm-up modes | All pass,
checkstyle clean |
| BE unit tests (23) | File format round trip and every kind of damage;
sizing and health check; write path over several segments with NULLs; transient
writer; `merge_rowset_meta`; sink-node receiver merge rules and sender
lifecycle | All pass (RELEASE build; ASAN left to CI) |
| Regression | `index_p0/global_point_index`: DDL, and results equal with
pruning on and off across both load paths, compaction, and an index added after
data; EXPLAIN in cloud mode | Both suites pass on a single-node cloud cluster
(FE, BE and meta service built from this branch) |
Manual check on the same cluster, 8 buckets, 2000 rows loaded through the
default memtable-on-sink-node path: `WHERE ev = <absent value>` plans
`globalPointIndex: ev -> 0/8 tablets (probes=1, unchecked=0)`; `WHERE name =
<present value>` plans `1/8` and returns the row.
## 8. Risks
| Risk | Mitigation | Signal |
|---|---|---|
| A saturated bloom silently stops pruning | Exact sizing on rewrites;
sizing self-check | `global_point_index_blooms_undersized_total`, EXPLAIN
kept/total |
| Too many blooms per tablet before compaction | Per-tablet fpp budget;
cumulative compaction score bias | EXPLAIN kept/total, rowset count |
| Cold file cache puts remote reads on the planning path | Warm-up and
repair sweep; timeout keeps tablets | Prune latency, `unchecked=` in EXPLAIN,
warm-up metrics |
| An encoding mismatch between FE and BE would be a false negative | Only
strings and fixed-width integers are probed; unit test pins the byte layout |
Regression results with pruning on vs off |
| Mixed-version loads | Capability flag in the open-stream response |
Rowsets without descriptors (`blooms_missing`) |
## 9. Rollout
The index is opt-in per column. A cluster without GLOBAL_POINT indexes runs
the same code paths as before.
| Phase | Acceptance |
|---|---|
| Create the index on one table, new data only | Regression suite passes;
EXPLAIN shows pruning on new data; results equal with the session variable on
and off |
| Backfill with BUILD INDEX or let compaction rewrite |
`blooms_missing_total` trends to zero |
| Watch | `blooms_undersized_total` stays flat; prune latency within budget |
Rollback: set `enable_global_point_index_prune` to false (FE) and
`enable_global_point_index_scan_gate` to false (BE), or drop the index. Stored
data is unaffected either way.
## 10. Open questions
1. Should the file cache route `.gpidx` writes directly to the INDEX queue?
Today written files always land in NORMAL; the sweep moves them later.
2. Should the planner pick the most selective indexed column when several
have predicates (needs NDV statistics), instead of the first one?
3. Should FE probe CHAR, LARGEINT and date columns once their BE byte layout
is pinned by shared test vectors?
4. Should plan-time pruning also run in local mode, with a snapshot
mechanism equivalent to the cloud one?
## Appendix A. Glossary
| Term | Meaning |
|---|---|
| bloom | The GLOBAL_POINT bloom filter of one column for one rowset |
| descriptor | `ColumnPointIndexPB` in the rowset meta, pointing to the
bloom's file |
| definite miss | Every visible non-empty rowset has a usable bloom, and
none of them may contain any probe value |
| degraded / unchecked | Kept because it could not be checked, not because
of a hit |
| candidate | A tablet the BE says to keep |
| B | Number of blooms of a tablet at a snapshot, about its number of
non-empty rowsets |
| plan-time pruning | Dropping tablets in FE before scan ranges are sent
(4.7) |
| scan-time gate | Skipping a rowset in the BE reader (4.6) |
## Appendix B. Decision log
| Date | Decision |
|---|---|
| 2026-09-27 | Target apache/master; submit as one PR of 9 commits |
| 2026-09-27 | `enable_global_point_index_sink_build` defaults to true, made
safe for mixed versions by the open-stream capability flag |
| 2026-09-27 | Fixed on the way: `CREATE INDEX ... USING GLOBAL_POINT`
created an inverted index, because the statement builder did not map the type |
| 2026-09-27 | Dropped from the first version: write-time INDEX-queue
routing in the io layer (see open question 1) |
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]