Skip to content

[Feature] Support GLOBAL_POINT index for cross-tablet point-query pruning #66787

Description

@datavisorethanqiu

Search before asking

  • I had searched in the issues and found no similar issues.

Description

Description

Problem

For a point query such as WHERE c = X, Doris can prune partitions/tablets at planning time when the predicate sufficiently constrains the partition or distribution keys. For a non-partition, non-distribution column, however, the query may still fan out to every tablet selected by the planner.

Existing indexes (zonemap, bloom filter, inverted index, short key, primary key, etc.) are local pruning structures. They can reduce the work within an already selected tablet, but they do not provide the FE with a cross-tablet membership view for deciding which tablets need to be scanned in the first place.

This is particularly expensive in storage-compute-separated mode, where a wide fan-out amplifies scan-range scheduling, tablet/rowset metadata work, and remote I/O.

The production workload that motivated this proposal is a ~5B-row event table with 566 tablets:

DISTRIBUTED BY HASH(user_id)

and queries such as:

WHERE eventId = 'Y'

eventId is near-unique and is neither a partition nor distribution column. A query returning one matching row therefore fans out across 566 tablets. An inverted index on eventId can efficiently reject data inside 565 of them, but only after those tablets have already entered the scan path.

The dominant cost is remote I/O and fan-out overhead rather than predicate evaluation itself, so adding compute does not remove the underlying I/O amplification.


Proposal

Introduce a new index type, tentatively named GLOBAL_POINT, that adds a cross-tablet membership layer.

For each indexed column, every rowset stores an immutable bloom filter in an independent index file on remote storage next to the rowset's segment files. RowsetMeta stores only a small descriptor.

At planning time, the FE extracts supported equality predicates, asks the BEs to probe the bloom filters for the rowsets visible at the query snapshot, and removes tablets for which all relevant blooms are definite misses.

Example:

CREATE TABLE t_evt (
    id BIGINT NOT NULL,
    dt DATE NOT NULL,
    ev INT NULL,
    name VARCHAR(64) NULL,
    INDEX idx_ev (ev) USING GLOBAL_POINT
)
DUPLICATE KEY(id)
DISTRIBUTED BY HASH(id)
BUCKETS 32;

The index can also be added later:

ALTER TABLE t_evt ADD INDEX idx_name (name) USING GLOBAL_POINT;

Adding the index is a metadata-only light schema change. New rowsets build the bloom from that point forward; historical rowsets require an explicit backfill:

BUILD INDEX idx_name ON t_evt;

The bloom FPP is chosen conservatively to account for the union of multiple rowset blooms within a tablet. Using the same FPP independently for every rowset does not work well: for example, with a 1% FPP per bloom and 50 rowsets, the probability that at least one bloom produces a false positive approaches 40%.

The exact FPP budgeting policy can be discussed as part of the design; the important property is that false positives only reduce pruning effectiveness and never affect correctness.


Architecture

1. Write path

Every rowset writer (including load, compaction, and backfill) feeds indexed values into a separate bloom writer.

When the rowset is finalized:

  • the bloom is written to a self-validating index file;
  • a small descriptor is recorded in RowsetMeta;
  • the normal rowset data path is unchanged.

The additional write-side work is approximately one hash-and-insert per indexed value, with no additional read I/O.

2. Plan-time pruning

This is the main optimization.

The FE:

  1. extracts supported = / IN predicates;
  2. encodes probe values identically to the write path;
  3. pins the query snapshot version;
  4. sends one pruning RPC per BE, rather than one RPC per tablet.

Each BE probes the bloom filters of rowsets visible at that snapshot.

The FE then intersects the result with the tablets already selected by the normal planner. The global index can only remove tablets; it never adds tablets that the planner did not select.

The planning-time probe path is designed to avoid turning an index lookup into additional remote I/O. If the required bloom is unavailable locally or cannot be used safely, the result is treated as unknown and the tablet is kept.

3. Scan-time gate

As a second layer, the BE scanner can re-check the bloom at rowset granularity and skip a rowset on a definite miss.

This provides an additional guard when plan-time pruning is unavailable or degraded.

4. Background warm-up

An FE-side daemon detects newly available or restarted BEs and asks them to prefetch relevant bloom files into the local file-cache INDEX queue.

This allows the first query after a BE restart to probe blooms at local-storage latency rather than requiring query-time remote fetches.


Why one bloom per rowset?

A committed rowset is immutable, so its bloom can also be immutable:

build on write
append with new rowsets
never rewrite an existing bloom

This aligns naturally with Doris snapshot semantics.

At a given query snapshot, the visible data consists of a known set of rowsets. Pruning is based only on the blooms corresponding to those visible rowsets.

A mutable tablet-level aggregate bloom would introduce a much harder correctness problem: after a rowset is committed, the aggregate must immediately and reliably reflect that rowset before it can be used for pruning. If the aggregate lags behind the committed data, the system could incorrectly prune a tablet that actually contains a matching row.

That is the failure mode this design must avoid.

Per-rowset decomposition removes the need for a mutable membership structure and makes stale or missing state naturally fail open.

Blooms are also not merged across rowsets. Compaction instead rebuilds a new bloom from the compacted data, just as it produces a new immutable rowset.


Fail-open by construction

The core rule is:

Over-scan, never skip.

A probe has three logical outcomes:

hit
definite miss
unknown

Only a definite miss is allowed to prune data.

Any failure or uncertainty maps to unknown, which means the tablet or rowset is kept.

Examples include:

  • missing bloom descriptor;
  • historical rowset not yet backfilled;
  • unreadable or corrupt index file;
  • local cache miss;
  • RPC timeout;
  • unsupported predicate or type;
  • snapshot/version uncertainty.

A broken, incomplete, or temporarily unavailable index can therefore make a query slower, but it must not change the query result.

The prune view must also be at least as fresh as the query snapshot. If that cannot be guaranteed, pruning is disabled for the affected tablet rather than using an older view.

Finally, global pruning only intersects with the planner's existing tablet set:

final tablets = planner tablets ∩ global-index candidates

It never introduces tablets the planner itself did not select.

With no applicable index, or with the feature disabled, the behavior is equivalent to the existing scan path.


Observability

EXPLAIN shows an attribution line whenever global pruning is attempted:

tablets=1/64, tabletList=1785089744872
globalFilter: ev[bloom] -> 1/64 tablets (probes=1, degraded=0)

degraded counts tablets that were kept for safety reasons rather than because of a bloom hit.

This makes a partially unavailable or ineffective index visible as an explicit metric instead of appearing only as increased query latency.


Scope and limitations

  • Plan-time pruning is currently implemented only for storage-compute-separated mode. This is where tablet fan-out is most expensive and where the current metadata/cache design was developed. The write path and scan-time gate are not inherently cloud-specific, so non-cloud plan-time pruning could be added later.
  • Supported column types must have stable, exact-match encoding between the write and probe paths. The current implementation supports integers, strings, and date/datetime types.
  • DECIMAL is currently rejected because differing scale representations can break byte-identical probe encoding. Any type mismatch or unsupported encoding disables pruning rather than risking a false negative.
  • The index is not useful for low-cardinality values that occur in most tablets. A global membership filter cannot prune a value that is genuinely present almost everywhere.
  • This is not a row-level lookup index. It only decides which tablets or rowsets are worth opening. Row location inside the selected data continues to use the existing local indexes and scan path.
  • IS NULL and predicates involving NULL do not currently participate. Bloom filters contain only non-null values.

Verification

Test environment:

Cloud smoke environment:
FDB + Meta Service + FE/BE + MinIO

Table:
64 tablets
10,004 rows
GLOBAL_POINT indexes on ev and name

Existing value:

EXPLAIN SELECT * FROM gp_cloud.t_evt WHERE ev = 4242;

Result:

tablets=1/64, tabletList=1785089744872
globalFilter: ev[bloom] -> 1/64 tablets (probes=1, degraded=0)

The query is reduced from 64 tablets to 1.

Nonexistent value:

EXPLAIN
SELECT *
FROM gp_cloud.t_evt
WHERE ev = 999999999;

Result:

globalFilter: ev[bloom] -> 0/64 tablets (probes=1, degraded=0)

All 64 tablets are pruned.

As a correctness cross-check, pruning can be disabled and the same predicate evaluated through the normal path:

SET enable_global_point_index_prune = false;

SELECT COUNT(*)
FROM gp_cloud.t_evt
WHERE ev = 999999999;

Result:

0 rows

Warm-up was verified separately by clearing the local file cache and restarting the BE. The background daemon repopulated the bloom files within approximately 30 seconds, allowing subsequent pruning to use locally cached index data without requiring a user query to trigger the warm-up.

The 566-tablet, ~5B-row table described in the Problem section is the production workload that motivated the design.

The pruning numbers above are from the smaller reproducible smoke environment. We do not yet have a published end-to-end benchmark at the full production scale.


Alternatives considered

1. Exact value → tablet mapping in Meta Service

An exact mapping avoids false positives, but its minimum information footprint is already large.

For approximately 5B distinct values across 566 tablets, the theoretical lower bound for storing only the tablet identity is roughly:

N × log2(T) / 8 ≈ 5.3 GiB

before accounting for keys, versions, metadata, or database overhead.

In our sizing, a straightforward FDB KV representation would require roughly 200–300 GB once practical key/value overhead is included.

It would also make every point lookup depend on an FDB lookup rather than a reusable BE-local cache.

Rejected primarily on metadata capacity and lookup-path cost.

2. Mutable tablet-level aggregate bloom

This would reduce the number of bloom files to probe, but introduces a synchronization requirement between rowset commit and aggregate-bloom publication.

If committed data becomes visible before the aggregate is updated, the stale bloom can incorrectly prune a tablet containing matching data.

That is the one class of failure this design cannot permit.

Rejected on correctness complexity.

3. Per-segment blooms

A wide-table base rowset may contain tens of segments.

Keeping the same effective tablet-level FPP would require much tighter per-bloom FPP as the number of independent filters increases, while also multiplying the number of descriptors, files, and probes.

Per-rowset granularity provides a better tradeoff between bloom size, probe count, and pruning effectiveness.

4. Reusing the existing bloom-filter index

The existing bloom-filter index is designed for scan-time local pruning.

It is not currently exposed as a compact cross-tablet membership structure that can be probed before the FE commits to scanning a tablet.

The missing capability here is not bloom filtering itself; it is making membership information available safely and cheaply at planning time.


Questions for the community

  1. Should this be introduced as a distinct index type such as GLOBAL_POINT, or would it fit better as a global/tablet-pruning mode of the existing bloom-filter index?
  2. Should plan-time pruning support non-cloud mode in the first version, or is landing the storage-compute-separated implementation first a reasonable scope?
  3. Does this feature warrant a DSIP before submitting the implementation PRs?

A working implementation is currently running in our branch-4.0-based environment. We would appreciate feedback on the overall direction and interface before submitting the implementation.

Use case

No response

Related issues

No response

Are you willing to submit PR?

  • Yes I am willing to submit a PR!

Code of Conduct

Metadata

Metadata

Assignees

No one assigned

    Labels

    kind/featureCategorizes issue or PR as related to a new feature.

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions