Skip to content

Core, Spark: SCALAR Key Lookup Index — reference implementation of #16961 - #18187

Draft
osscm wants to merge 39 commits into
apache:mainfrom
osscm:scalar_sec_index
Draft

osscm wants to merge 39 commits into
apache:mainfrom
osscm:scalar_sec_index

Conversation

@osscm

@osscm osscm commented Sep 20, 2026 •

Copy link
Copy Markdown

Reference implementation of the SCALAR index type defined in #16961.

Supersedes #17426, which the stale bot auto-closed for 30 days of inactivity (explicitly not a
judgement on the merits) before this branch's Spark integration existed. Continuing that work here
with a full Spark integration added since.

What is included

Spec (format/index.md from #16961) + Core/API/Data/Spark layers:

Layer 0 — Data model: IndexMetadata, IndexSnapshot interfaces and JSON parsers
Layer 1 — Storage: IndexCatalog, IndexMetadataIO, InMemoryIndexCatalog
Layer 2 — Tracking file: TrackingFileWriter + TrackingFileReader (Avro, same as Iceberg manifests)
Layer 3 — Build: HashTransform, LeafFileMetadata, ScalarIndexCommitter, LeafFileEntry/Writer/Reader
Layer 4 — Spark integration:

  • CALL system.build_scalar_index(table, columns, transform, options) populates an index from a live table
  • Read-path: SparkScanBuilder resolves equality, IN, and range predicates (<, <=, >, >=,
    including a BETWEEN, which Spark decomposes into two range predicates on the same column)
    against an existing index -- IN and equality work against HASH or IDENTITY (each literal
    resolves to its own transform value, queried as a separate point); range resolution only for
    IDENTITY-transform indexes, since HASH scatters values across buckets and a contiguous range
    on the original column does not correspond to a contiguous range of transform values the way
    it does for IDENTITY
  • FileScanTaskFilteringScan enforces file-level pruning at scan-planning time (not just advisory) --
    self-verifying: falls back to the unfiltered file set if resolved paths don't match real
    candidates, so a bug here can only miss an optimization, never return a wrong result
  • Covered/uncovered-files staleness handling (Huaxin Gao's PK Index proposal, Section 8): a stale
    index still prunes files that existed at its own snapshot, and always scans files added since --
    rather than falling back to no pruning at all on any snapshot mismatch
  • Compaction detection: if a non-append snapshot (e.g. a compaction/rewrite) sits between the
    index's snapshot and the current one, falls back to no pruning rather than risk resolving
    matches against file paths compaction may have rewritten away
  • Incremental (append-only) build mode: options => map('mode', 'incremental') builds leaf files
    only for data added since the index's last snapshot, adopting the append-only option from
    Huaxin Gao's PK Index proposal (Section 7.2)
  • Durable index catalog: registrations now survive a JVM restart (DurableIndexCatalog persists a
    pointer file per index under the table's own location, scoped by table UUID)
  • Planning-cost bound (Open Question 2): a query resolving to more candidate leaf files than a
    configurable cap (default 100, scalar-index.max-candidate-leaf-files table property) falls
    back to normal planning rather than risk making planning slower than no index at all
  • Partition-column warning (Open Question 5): building an index on a column that's already in
    the table's partition spec logs a warning rather than rejecting -- redundant overhead, not wrong
  • Index discovery: the read path now discovers a covering index via IndexCatalog#listIndexes
    and matches by keyColumnIds(), matching the design doc, instead of a hardcoded
    <column>_idx name lookup

Test status

Core/API/Data: unit tests from earlier layers (index metadata round-trip, tracking/leaf file
read-write, in-memory catalog, end-to-end build+commit+lookup wiring), plus TestDurableIndexCatalog
for the new durable catalog.
Spark: TestBuildScalarIndexProcedure, TestScalarIndexScanPruning, and TestFileScanTaskFilteringScan
green, including new coverage for staleness handling, compaction fallback, incremental builds,
range-predicate resolution (two-sided, one-sided, and confirming a HASH index correctly declines
the range path), IN-predicate resolution (against both HASH and IDENTITY, mixed present/absent
values, and no values present at all), a forced-low planning-cost bound still returning correct
results via fallback, and building on a partition column succeeding (warns, does not reject). New
TestSparkScanBuilderCandidateLeafFiles verifies the IN-predicate leaf-file dedup mechanism itself
(not just correctness -- a black-box SQL result is a Set that would mask a missing dedup
by coincidence, so this asserts the deduped candidate list directly against a real tracking file).
Full spark-extensions module regression run clean (1675 tests, 106 skipped, 0 failures).

Known gaps / explicit non-goals for this pass

  • DurableIndexCatalog.listIndexes (used by the new discovery path above) only works when the
    table's FileIO supports prefix listing (SupportsPrefixOperations)
  • No snapshot-expiration-triggered cleanup: the design doc says an index snapshot becomes eligible
    for removal (and its leaf files deleted, if unreferenced) when the corresponding table snapshot
    expires; nothing implements this yet, so orphaned index data accumulates indefinitely
  • No schema-evolution handling: the design doc calls for disabling the index on a dropped key
    column, and forcing a rebuild on type promotion (e.g. int -> long, since the physical encoding
    changes and old entries would silently produce wrong lookups); neither is detected
  • Pruning is file-level only, not exact-row/position-level (tracked as a follow-up)
  • When the index confirms a key is absent, the scan is NOT pruned to zero files -- deliberately
    deferred since a bug there would silently return wrong (empty) results, not just miss an
    optimization
  • Scoped to plain SELECT batch scans; incremental-append/changelog/merge-on-read/copy-on-write
    scans are untouched
  • Trino integration not started (planned as a separate follow-up phase)
  • Marked Draft -- looking for early feedback on direction, not requesting a merge review yet

Relationship to the PK Index proposal

Thanks @huaxingao for putting together the PK design, it already has covered quite a few categories, specially delete.
I’ll further review it and find delta between the two design docs, and we can converge iI, appreciate your time!

Huaxin Gao's Primary Key Index for Apache Iceberg
proposes another concrete index type under the same Secondary Index framework (#16961), looks like both has map key -> (file, position) and both use scan-time file pruning as a core mechanism.

Builds on

osscm added 23 commits July 21, 2026 12:16
Applies the index spec document from apache/iceberg PR apache#16961 authored
by pvary. This defines the generic secondary index framework: SCALAR
index type, HASH/IDENTITY transforms, Index Metadata, Tracking File,
and Leaf File structure.
Implements IndexMetadata and IndexSnapshot POJOs matching the field
definitions in format/index.md (PR apache#16961). Includes JSON parsers
and a unit test suite covering round-trip serialization, snapshot
lookup helpers, builder validation, and spec field name verification.

6/6 unit tests passing.
Adds Layer 1 of the SCALAR index implementation:
- IndexIdentifier: identifies an index by table + name
- IndexCatalog: interface with create/load/update/drop/list operations
  with optimistic concurrency on updateIndex
- IndexMetadataIO: read/write index metadata JSON via FileIO, with
  metadata file location naming (00001-{uuid}.metadata.json)
- InMemoryIndexCatalog: thread-safe in-memory implementation for tests,
  includes CAS-based conflict detection on updateIndex

14/14 unit tests passing.
Implements the tracking file format for the secondary index spec:
- TrackingFileEntry: value class holding leaf file location, bounds,
  record count, size, and optional key metadata (field IDs from spec)
- TrackingFileWriter: writes entries as an Avro DataFile using Snappy
  compression, writing to Iceberg OutputFile (same pattern as manifests)
- TrackingFileReader: readAll() and readMatching(min, max) for
  planning-time pruning — returns only leaf files whose transform-value
  range overlaps the query range

18/18 unit tests passing.
Implements the core build logic for the SCALAR index:
- HashTransform: maps key values to hash buckets [0, numBuckets).
  Phase 1 uses Java hashCode; will align with Iceberg murmur3 bucket
  transform in a follow-up.
- LeafFileMetadata: captures path, bounds, record count, and size of
  one written leaf file; converts to TrackingFileEntry for the tracking
  file writer.
- ScalarIndexCommitter: given a list of leaf files from the Spark build
  job, writes the tracking file (Avro), creates an IndexSnapshot, writes
  the metadata JSON, and commits to the IndexCatalog with CAS semantics.

Also adds scripts/build_scalar_index_taxi.scala — an end-to-end Spark
script demonstrating the full build flow on NYC Yellow Taxi data using
the medallion column as the SCALAR HASH index key.

25/25 unit tests passing.
…hot, and commit validation

Fixes three errorprone ERROR-level findings that blocked compilation:
missing Locale.ROOT in String.format calls in ScalarIndexCommitter and
IndexMetadataIO, and a RuntimeException that should have been
UncheckedIOException in IndexSnapshotParser.

Adds test coverage that was previously missing entirely:
- IndexMetadataIO.read() round-trip against the actual bytes written by
  ScalarIndexCommitter, not just the in-memory catalog object.
- GenericIndexMetadata.removeSnapshot(), including currentSnapshotId
  reassignment when the removed snapshot was current.
- ScalarIndexCommitter validation for null/empty leafFiles.
- A deterministic characterization test (via a new RaceInjectingCatalog
  test double) pinning down that a race on the first commit for an
  identifier throws AlreadyExistsException, not the
  ConcurrentModificationException the class's javadoc says callers
  should retry on. This documents current behavior; it is not a fix.
…dex reads

Implements the read side of the SCALAR index, which previously did not
exist at all -- only the build/commit side (ScalarIndexCommitter,
TrackingFileWriter/Reader) was implemented before this.

Lives in iceberg-data, not iceberg-core, since leaf files are Parquet
and iceberg-core's main sourceset has no compile-time dependency on
iceberg-parquet (only testRuntimeOnly) -- that boundary looks
intentional (core is kept format-agnostic), so this follows the
project's existing convention of putting Parquet-specific I/O where the
dependency already exists, same as GenericParquetWriter/Readers.

LeafFileEntry.schema(NestedField keyField) is the single source of
truth for the leaf-file schema, used by both writer and reader, so they
cannot drift apart on field IDs. The three synthetic fields
(transform_value, file_path, position) use Integer.MAX_VALUE-(101-103)
as field IDs, placed just past Iceberg's own reserved metadata-column
range so they never collide with a source table's real field IDs or
Iceberg's own metadata columns.

LeafFileReader.readMatching pushes the predicate to Parquet via
ReadBuilder#filter for row-group-level statistics pruning, then applies
an Evaluator bound to the leaf schema for exact per-record matching.
This second step is necessary: the createReaderFunc-based read path
(org.apache.iceberg.parquet.ParquetReader) only skips whole row groups
by statistics and does not filter individual records within a row
group that wasn't skipped, unlike the older readSupport-based path.
Confirmed by two failing tests before this fix -- filtered reads were
returning every row in the file.

LeafFileWriter does not sort; it writes rows in whatever order the
caller provides. Correct pruning depends on the caller (the future
index build job) having sorted by (transform_value, key_value) first --
this is not yet enforced or verified anywhere.
LeafFileWriter now tracks the last written (transform_value, key_value)
tuple and throws IllegalStateException immediately if a new entry is
out of order, rather than silently writing an unsorted leaf file.

This matters even for today's single-column HASH/IDENTITY transforms,
not just a hypothetical multi-column future: unsorted leaf-file rows
mean Parquet's row-group statistics can't prune anything, since every
row group ends up spanning close to the full transform_value range.
Results stay correct either way (row-group skip is always a safe
exclusion, and the per-record Evaluator in LeafFileReader guarantees
exact matches), but pruning effectiveness collapses silently -- no
error, no warning, just a query that's slow for no visible reason.

Deliberately does not sort internally (that's the rejected alternative,
option 1): the real build job already has to shuffle/partition data by
transform_value bucket to write leaf files at all, and Spark's
sortWithinPartitions gets the sort essentially for free as part of that
same operation. Sorting again here would duplicate that cost on every
leaf file to guard against a bug a correct caller never triggers.
Equal consecutive keys are allowed, since a key value is not required
to be unique across rows.
Adds TestScalarIndexEndToEnd, exercising the full SCALAR index chain as
one flow for the first time: write sorted leaf files via LeafFileWriter,
commit them through ScalarIndexCommitter into an InMemoryIndexCatalog,
then look up a key exactly the way a planner would -- compute its
transform value via HashTransform, narrow to candidate leaf files via
TrackingFileReader, then resolve the exact row within those leaf files
via LeafFileReader. Covers both an exact-match lookup (verifying the
correct file_path/position comes back) and a missing-key lookup
(verifying an empty result).

Until now the build/commit side and the read side were two separately-
tested halves that had never actually been exercised together. This
doesn't wire the index into any real engine or source table -- it
proves the pieces we've built work together correctly, nothing more.

Uses a minimal local-filesystem FileIO test double (real files, since
Parquet needs seekable I/O) rather than the in-memory doubles used
elsewhere for Avro-only tests.
Implements "CALL system.build_scalar_index(table, columns, transform,
options)", populating a SCALAR index from an actual Spark job for the
first time -- everything before this was either a standalone demo
script or a JUnit test manually orchestrating the pieces.

Computes each row's position via row_number() over a window
partitioned by input_file_name(), before any shuffle, so it reflects
physical file-scan order. Computes transform_value by calling the real
HashTransform class through a typed UDF (HASH) or a direct cast
(IDENTITY, numeric keys only, since a string can't cast to the long
transform_value column) -- never reimplements the transform logic in
Spark-land. Shuffles by transform_value bucket and sorts within
partitions, then writes one leaf file per partition via LeafFileWriter,
collecting per-partition metadata back to the driver via
Encoders.javaSerialization (a Dataset<Row> encoder for a custom struct
isn't a stable enough API to rely on here). Commits through
ScalarIndexCommitter into a new session-scoped registry,
SparkIndexCatalogs, added alongside this and following the same
singleton pattern as the existing ScanTaskSetManager.

Only a single key column is supported (multi-column composite indexes
remain an explicit Non-Goal in the design proposal); the procedure
rejects multiple columns with a clear error rather than silently using
only the first one.

Important caveat: this Spark code has not been compiled or run.
Gradle cannot even start in this sandbox (blocked at the local
socket/daemon level), so every Iceberg-core API call was verified
against the actual source in this repo, but the Spark-specific pieces
(UDF registration, window functions, mapPartitions encoder usage) rely
on general Spark API knowledge plus in-repo precedent where found, not
execution. Needs a real compile and test run before being trusted.
Adds tryPruneUsingScalarIndex(), called from pushPredicates() after
normal predicate pushdown. On an equality predicate matching an
indexed column, resolves the exact source file via
TrackingFileReader + LeafFileReader (using SparkIndexCatalogs to find
the index) and, if it resolves to exactly one file, adds an
Expressions.equal(_file, ...) constraint on top of the existing
filters. File-level pruning only in this pass -- true row-position
pushdown into the scan tasks is a further, not-yet-attempted
refinement, noted in the method's javadoc.

Deliberately fail-open everywhere: no index registered, a stale index
snapshot relative to the table's current snapshot, an unsupported
predicate shape, zero or multiple leaf-file matches, or any I/O error
all fall back silently to normal planning. For a table with no index
(the common case) this method is a complete no-op, which is the whole
point -- existing tables and existing tests should see no behavior
change at all.

Caught and fixed a real bug before committing: the write side
(BuildScalarIndexProcedure) and this read side each need to construct
byte-identical TableIdentifiers to look up the same IndexIdentifier --
otherwise indexExists() here would silently and permanently return
false with no error, making the whole integration a no-op that never
activates. Both sides now derive it the same way, from
TableIdentifier.parse(table.name()) on the core Table object, not from
Spark's own catalog Identifier type.

Same caveat as the previous commit: unverified against a real Spark
runtime, since gradle cannot start in this sandbox at all (blocked at
the local socket/daemon level, not just the domain-allowlist issue
from earlier in this session). Every Iceberg-core API call
(UnboundPredicate, NamedReference, Schema.findField, MetadataColumns,
TableIdentifier.parse) was verified against the actual source in this
repo; needs a real compile and test run before being trusted.

Also adds TestScalarIndexScanPruning: functional-correctness tests
(equality query returns the right row / no rows for a missing key /
still correct with no index at all) plus one best-effort test using
Dataset#inputFiles() to check the indexed query reads no more files
than an equivalent un-indexed query -- flagged in that test's own
javadoc as unverified whether inputFiles() actually reflects Iceberg's
final pruned file set for this Scan implementation.
The previous version hand-rolled the whole build pipeline (compute
buckets, sort, write Parquet, commit) directly in the script, in
parallel to the actual production code -- meaning it could silently
drift from what BuildScalarIndexProcedure actually does. It also never
computed `position`, so it only ever achieved file-level pruning, not
the exact-row (file, position) pruning that's the actual differentiator
of a SCALAR index over the sibling Bloom filter index.

This version calls the real procedure via SQL instead, so the demo
exercises the real code path. It also picks a real medallion value out
of the table for the lookup rather than a hardcoded one that might not
exist in whatever dataset is loaded, and uses Dataset#inputFiles() to
show the file count before/after rather than hand-computing the hash
bucket and calling TrackingFileReader directly.

Same caveat as the last two commits: not run against a real Spark
session.
Table.uuid() returns java.util.UUID, not String -- ScalarIndexCommitter.
commit() needs the String form. First real compile error surfaced by
actually running the build, exactly as expected given this code was
entirely unverified before now.
transformValueColumn() referenced the key column by its original name,
but by the time it runs the earlier .select() has already renamed it
to "__key" -- the original name no longer exists in that DataFrame's
schema. Caught by actually running the test: UNRESOLVED_COLUMN on the
UDF's input reference.
Encoders.javaSerialization() requires the target class to be public,
not just Serializable -- caught by actually running the test:
SparkUnsupportedOperationException, "is not a public class."
BuildResult doesn't need the same fix since it's a plain driver-side
return value, never passed through any Spark Dataset/Encoder API.
Spark's ClosureCleaner requires every object captured inside a UDF
closure to be serializable, even for local (non-distributed)
execution -- caught by actually running the read-path test:
java.io.NotSerializableException: org.apache.iceberg.index.HashTransform,
surfaced through BuildScalarIndexProcedure's UDF-based transform_value
computation. HashTransform only wraps a single int field, so this has
no real downside -- the same pattern Iceberg's own core classes
(Schema, Types.NestedField) already follow for the same reason.
row_number() returns IntegerType, not LongType -- __position ended up
as an Integer-backed column, but writeLeafFilePartition read it back
via row.getAs("__position") with an inferred Long type, throwing
ClassCastException at runtime (not a compile error, since getAs()'s
type parameter is just a cast, not checked against the row's actual
schema). Fixed at the source by casting the column expression to
LongType explicitly, rather than only at the read site, so the
DataFrame's schema genuinely matches what downstream code expects.
…ning

tryPruneUsingScalarIndex() injected Expressions.equal("_file", ...) into
filterExpressions to constrain the scan to a resolved file. That list
also feeds SparkSchemaUtil.prune() (via pruneColumns) and the eventual
core Scan#filter() call, both of which bind expressions against the
real table schema via Binder -- but "_file" is a Spark-only metadata
column, not a real schema field, so binding it always threw
ValidationException: Cannot find field _file. This broke every query
that reached column pruning, not just ones an index could resolve.

Core Iceberg Scan/Expression/Binder has no concept of filtering by
metadata columns like _file at all -- enforcing this would need a Scan
decorator around planFiles()/planTasks() at the task level, a larger
follow-up. For now tryPruneUsingScalarIndex logs the resolved file but
does not enforce it; the underlying predicate is still pushed down and
applied normally, so functional correctness is unaffected.
The test built its own IndexIdentifier from the Spark catalog
Identifier's namespace()/name(), but BuildScalarIndexProcedure commits
under an IndexIdentifier derived from TableIdentifier.parse(table.name())
-- the core Iceberg Table's own name, chosen deliberately because
SparkScanBuilder on the read side only has the core Table, not the
original Spark Identifier, so both sides need a common, independently
computable source. The two derivations do not necessarily produce the
same TableIdentifier (table.name() can be catalog-qualified), so the
test's indexExists() lookup used a different key than what was
committed and always returned false. Fixed by deriving the test's
IndexIdentifier the same way production code does.
Still failing after aligning both sides to derive IndexIdentifier from
table.name() -- root cause was one level deeper: validationCatalog is
a separate Catalog handle (e.g. a freshly-constructed HadoopCatalog, or
a HiveCatalog/RestCatalog initialized under its own catalog-name
property), independent of the Spark-registered catalog that
BuildScalarIndexProcedure and SparkScanBuilder actually load the table
through. Several Iceberg catalog implementations embed the loading
catalog's own name into Table.name(), so the same physical table
loaded via validationCatalog vs. via Spark produced two different
name() strings, hence two non-equal TableIdentifiers, hence
indexExists() always returning false in the test -- consistent with
every catalog config in the parameterized suite failing identically.

Fixed by loading the table via Spark3Util.loadIcebergTable(spark,
tableName), the same pattern used throughout the rest of this test
suite, so the test's table.name() matches what production code
actually sees.
tryPruneUsingScalarIndex() previously only logged the resolved file
path -- core Iceberg Scan/Expression/Binder has no concept of
filtering by metadata columns, so there was no way to push the
constraint through the normal Expression pipeline. Enforcing it
requires intercepting task planning directly.

Added FileScanTaskFilteringScan, a thin BatchScan decorator that
overrides only planFiles() to filter to a resolved set of file paths,
delegating everything else. This is enough because
SparkPartitioningAwareScan (confirmed by reading it) only ever calls
planFiles(), never planTasks(), when materializing Spark input
partitions -- so no changes were needed to the scan-wrapper classes
that consume it. Wired in only for the plain SELECT batch-scan path in
SparkScanBuilder#buildBatchScan(); incremental-append, changelog,
merge-on-read, and copy-on-write scans are left alone on purpose,
since row-level operations have different correctness considerations
worth their own review.

Generalized tryPruneUsingScalarIndex() from acting only when exactly
one file resolves, to the full set of matching files, so the
optimization is not limited to single-row lookups. The zero-match case
(index confirms the key is absent) is deliberately not pruned to zero
files: that would be correctness-sensitive (a bug would silently
return wrong empty results) rather than just missing an optimization,
and is left as a documented follow-up.

FileScanTaskFilteringScan self-verifies before trusting the resolved
paths: it matches them against real candidate files from normal
planning, and falls back to the unfiltered file set if none match --
guarding against a possible path-format mismatch between how the
index recorded paths (Spark input_file_name() at build time) and how
Iceberg reports them at scan time (DataFile#path()), since that
agreement has not been verified across every FileIO implementation.
This keeps the design invariant intact: the index must never be
required for correctness, only used opportunistically for pruning.

Strengthened TestScalarIndexScanPruning to assert an exact file count
now that pruning is actually enforced, and added
testPrunesBeyondNativeMinMaxStats, which constructs files whose
min/max ranges all cover the query literal so Iceberg own manifest
stats pruning cannot exclude any of them, to demonstrate the index
adds real pruning value beyond native stats rather than a test that
would pass with the pruning logic disabled.
…n test

The full spark-extensions module run (2722 tests) came back with only
8 failures, all in TestScalarIndexScanPruning, all "expected: 1 but
was: 0" on Dataset#inputFiles().length -- no regressions anywhere else,
confirming the decorator itself is not the problem, the verification
mechanism is.

Dataset#inputFiles() returns empty for Iceberg's DataSourceV2 scans in
this Spark version regardless of pruning -- it is not wired up for
non-FileScan V2 sources, so both the indexed and unindexed queries
always reported 0 files, which is also why the original, pre-existing
"<=" comparison always passed trivially before this change strengthened
it to an exact count. Separately, the resultDataFiles SQL metric
(ScanMetricsUtil, called from ManifestGroup) is recorded inside the
wrapped scan's own planFiles() -- upstream of this decorator's filter
-- so it reflects native manifest-stats pruning only, never this
decorator's additional restriction. Neither is a usable signal for this
class specifically.

Removed both SQL-level file-count assertions and the now-redundant
testInputFilesReducedAfterIndexBuild test (its distinguishing signal
was unreliable, and its correctness check duplicated an existing test).
Added TestFileScanTaskFilteringScan, a focused unit test against a real
(if minimal) Iceberg table built via the existing TestTables test
helper, verifying planFiles() directly: filters to a single allowed
path, filters to multiple allowed paths, falls back to the unfiltered
set when no allowed path matches any real candidate (the self-verifying
safety net), and delegates other methods unchanged. This is a more
direct and deterministic proof of the decorator's own logic than any
Spark Dataset-level signal available here.

Also switched FileScanTaskFilteringScan's path comparison from the
deprecated ContentFile#path() to #location().
Closeable#close() declares throws IOException, so try-with-resources
on CloseableIterable requires the enclosing method to declare or catch
it -- missed across all three try-with-resources blocks.
@github-actions github-actions Bot added API spark core data Specification Issues that may introduce spec changes. labels Sep 20, 2026
tryPruneUsingScalarIndex previously required an exact match between the
index's snapshot and table.currentSnapshot(), falling back to zero
pruning on any mismatch -- one unrelated commit after the index was
built made it fully inert until rebuilt. Adopted the covered/uncovered
model from Huaxin Gao's Primary Key Index for Apache Iceberg proposal
(Section 8, Staleness Semantics): files present at the index's own
snapshot ("covered") are still pruned via the index; files added since
("uncovered") are always included in the resolved set unconditionally,
since the index has no information about them. A stale index now still
helps instead of becoming a complete no-op.

Also relaxed the zero-covered-matches fallback: when there are
uncovered files, resolving to (0 covered matches + all uncovered files)
is still sound and still narrower than an unfiltered scan, so it no
longer falls all the way back. The zero-files-total case (no covered
matches, no uncovered files -- i.e. a fully fresh index) is unchanged
and still deliberately not pruned to zero, since that remains
correctness-sensitive in a way this isn't.

Fixed a real bug found while implementing this: the initial version
used IncrementalAppendScan to compute the uncovered file set, but it
silently filters snapshots down to appends-only rather than throwing
when a non-append snapshot (e.g. a compaction/rewrite) sits in the
range -- which would have made this method return an incomplete
uncovered set instead of failing. That's a correctness risk, not just
a missed optimization: a row physically moved by compaction into a
file outside both the index's covered matches and the (incomplete)
uncovered set would never be scanned again. Replaced with an explicit
walk of the snapshot ancestry that requires every snapshot in range to
be a pure append, throwing otherwise so the caller falls back to no
pruning at all -- verified by a new test that runs rewrite_data_files
between index build and query and confirms the row is still found.
build_scalar_index previously only supported a full rebuild -- every
call re-scanned the entire source table and rewrote every leaf file,
even to index a single newly-added row. Adopted the append-only option
from Huaxin Gao Primary Key Index for Apache Iceberg proposal (Section
7.2, her own Phase-1 recommendation): options => map(mode,
incremental) now builds leaf files only for data files added since
the existing index last snapshot, and appends them to the existing
leaf files rather than rewriting everything. Stale entries from
updated or deleted rows are not cleaned up incrementally -- a full
rebuild is what does that, matching the same recommendation.

Extracted IndexSnapshotUtil.addedFilePathsSince (previously a private
method on SparkScanBuilder, added for the covered/uncovered staleness
fix) into a shared org.apache.iceberg.spark utility, since the write
side needs the exact same "files added between these two snapshots,
requiring pure appends" logic to find what to incrementally index.
Keeping this correctness-sensitive logic in one place means a future
fix cannot apply to one call site and not the other.

Falls back to a full rebuild if there is no existing index yet, or if
anything about determining the added files fails (for example a
compaction ran since the index was last built) -- incremental is
requested behavior, not a guarantee, consistent with the index never
being required for correctness. When there are no new data files at
all, either refreshes just the snapshot pointer (no new leaf files
needed) or, if already exactly fresh, skips committing entirely rather
than creating a redundant index snapshot.

BuildScalarIndexProcedure leaf-file-writing pipeline (position to
transform value to sort to write) is now shared between the full and
incremental paths via a new buildLeafFiles helper, rather than being
duplicated for the incremental case.
InMemoryIndexCatalog (still used directly by SparkIndexCatalogs until
this change) loses track of which metadata file is current for an
index as soon as the JVM restarts, even though the metadata, tracking,
and leaf files themselves are already on durable storage -- only the
pointer to the current one lived in memory. IndexCatalog own javadoc
already describes the intended design ("stores a pointer to the
current index metadata file"), so a file-backed implementation of
that pointer is what was actually missing, not a new concept.

Added DurableIndexCatalog: one small pointer file per index name, at
<tableLocation>/metadata/scalar-indexes/<indexName>.pointer, containing
just the current metadata file location as plain text. createIndex,
updateIndex (optimistic-concurrency, same read-then-conditional-write
model InMemoryIndexCatalog already documents), dropIndex, and
indexExists all read or write that one file via FileIO -- no new
concept beyond what IndexMetadataIO already does for the metadata
files themselves. listIndexes needs directory listing, which not
every FileIO guarantees; it uses SupportsPrefixOperations when
available (true for most durable backends, including S3FileIO) and
throws UnsupportedOperationException otherwise rather than guessing.

Wired into SparkIndexCatalogs.catalogFor, replacing InMemoryIndexCatalog
as the default there. SparkIndexCatalogs own in-memory map now caches
only the IndexCatalog object for reuse within one process, not the
index metadata itself -- losing that cache after a restart costs
nothing beyond re-constructing a cheap wrapper object, since a fresh
DurableIndexCatalog re-derives the same state from its pointer files.
InMemoryIndexCatalog is unchanged and still used directly by core/data
module tests that do not need real durability.

New TestDurableIndexCatalog (core) exercises this against a real
HadoopFileIO and a real temp directory rather than opaque s3:// strings
like TestInMemoryIndexCatalog uses, since durability across a fresh
catalog instance (standing in for a process restart) is the entire
point being tested and cannot be verified with in-memory fakes.
Test run surfaced runaway snapshot counts across TestBuildScalarIndexProcedure
tests (expected 1 or 2 snapshots, got 3-12) caused by DurableIndexCatalog
keying its on-disk pointer files by table.location() alone. A table
location string can be reused after a drop and recreate (most catalogs
do not purge files this class does not know about on drop, since they
live outside Iceberg own metadata and manifest tracking entirely), so a
fresh table at a reused location was silently inheriting a previous,
logically unrelated table index registrations -- each new full-rebuild
test found an old pointer file still there and treated it as an update
of an existing index, carrying its entire snapshot history forward.

SparkIndexCatalogs already keys its own in-memory cache of IndexCatalog
instances by table UUID for exactly this reason; DurableIndexCatalog now
does the same for the durable registry path it manages, folding tableUuid
into the constructor and the on-disk directory structure.

Added a regression test constructing two catalogs at the same location
with different table UUIDs and asserting one cannot see the other
registrations.
…exes

Our own design doc lists range scans (WHERE col BETWEEN a AND b, via
IDENTITY) as a primary Goal alongside point lookups, but
tryPruneUsingScalarIndex only ever recognized EQ predicates -- IDENTITY
was fully implemented on the build side and completely unused for
pruning on the read side.

Generalized the read path: predicates are now grouped by column first
(a BETWEEN decomposes into two separate range predicates on the same
column, which need to be combined into one bound), then resolved
either as an equality lookup (HASH or IDENTITY, unchanged behavior) or
a range lookup (IDENTITY only -- HASH scatters values across buckets,
so a contiguous range on the original column does not correspond to a
contiguous range of transform values the way it does for IDENTITY).

The coarse tracking-file-level range bound is deliberately loose at
the GT/LT boundary (uses the literal value directly rather than value
+/- 1): LeafFileReader re-applies the exact original predicate
expression at the leaf-file level regardless, so looseness here can
only cost scanning one extra boundary leaf file, never exclude a real
match. Both TrackingFileReader and LeafFileReader already supported
range queries and range expressions respectively (their own javadocs
already described this case) -- the gap was entirely in
SparkScanBuilder never constructing one.

IN predicates are still not resolved against the index -- noted as a
follow-up in the updated javadoc, not attempted in this change.

Added tests for a two-sided range (BETWEEN-equivalent), a one-sided
range, a range matching nothing, and confirmed a HASH-transform index
correctly declines the range path and still returns correct results
through the normal residual-filter path instead.
Our design doc's own "Missing Building Blocks" table calls for Spark
integration covering "= and IN predicates"; only = (and, since the
previous commit, range comparisons) were resolved. tryPruneUsingScalarIndex
now also recognizes IN.

Each literal in an IN list resolves to its own transform value, queried
as a separate point rather than folded into one [min, max] range: HASH
in particular can scatter an IN-list's values across unrelated,
non-contiguous buckets, so there often is no single range that is both
correct and useful. IDENTITY works the same way for consistency, even
though its values happen to be orderable. Candidate leaf files from all
of an IN predicate's target points are deduplicated by location before
being read, since two IN values landing in the same bucket would
otherwise be read twice for no benefit. The leaf-file-level exact check
uses Expressions.in(...) directly, so a single leaf file scan picks out
every matching value in one pass rather than one equality check per
value.

Extracted the shared HASH/IDENTITY transform-value computation
(previously duplicated between the equality and, after the last commit,
IN branches) into a small transformValue helper, and introduced
TransformValueRange as a small holder so equality, IN, and range
predicates all funnel through the same "list of [min, max] sub-ranges to
query the tracking file for" shape rather than three divergent code
paths.

Added tests for IN against both a HASH and an IDENTITY index, including
a mix of present/absent values and no values present at all.
…rectness

The existing SQL-level IN-predicate tests (testInPredicateResolvesViaHashIndex,
testInPredicateResolvesViaIdentityIndex) prove correctness, but cannot actually
prove the dedup logic added alongside IN support is doing anything: the final
resolved path set is a Set<String>, which would naturally collapse duplicates
downstream even if the same leaf file were read twice for two colliding IN
values. A black-box SQL test genuinely cannot distinguish "dedup worked" from
"dedup is missing but the result is still correct by coincidence."

Extracted the dedup loop out of tryPruneUsingScalarIndex into a package-private
static method, collectCandidateLeafFiles(FileIO, String, List<TransformValueRange>),
and made TransformValueRange package-private too, so both can be exercised
directly against a real tracking file (via TrackingFileWriter) without a full
Spark session. New TestSparkScanBuilderCandidateLeafFiles asserts the actual
return value: two target ranges matching the exact same leaf file collapse to
one entry, two ranges overlapping the same wider leaf file also collapse to
one, and two ranges landing in genuinely different leaf files are not
incorrectly merged into one.
@osscm

osscm commented Sep 23, 2026 •

Copy link
Copy Markdown
Author

cc @huaxingao
cc @pvary

…uestions 5, 2)

Design doc's Open Question 5: should the spec discourage or reject a
SCALAR index on an already-partitioned column, since partitioning
already prunes on it and the index is redundant overhead there.
BuildScalarIndexProcedure now checks the key column against the
table's current partition spec's source ids and logs a warning if it
matches -- warns, does not reject, matching the "always advisory,
never required" design and the fact that the doc frames this as an
open question, not a settled rejection. Redundant is not wrong.

Design doc's Open Question 2: without a bound, a misconfigured index
(a very large hash.num-buckets, or a wide IN list) could resolve to
more candidate leaf files than would ever be worth opening, making
planning slower than no index at all. tryPruneUsingScalarIndex now
checks the deduped candidate leaf file count against a cap (default
100, overridable per table via the scalar-index.max-candidate-leaf-files
property) and falls back to normal planning if exceeded -- the original
predicate is still applied downstream regardless, so this can only cost
a missed optimization, never a wrong result.

Added tests for both: building on a partition column still succeeds
(not rejected), and a query forced past an aggressively low bound
(set via ALTER TABLE ... SET TBLPROPERTIES) still returns correct
results through the fallback path.
tryPruneUsingScalarIndex previously derived a fixed IndexIdentifier of
"<column>_idx" and called indexExists/loadIndex directly, which only
worked because BuildScalarIndexProcedure happens to name indexes that
way today. The design doc describes the planner discovering indexes via
IndexCatalog#listIndexes(tableIdentifier) and matching by which columns
an index actually covers -- not by guessing a name.

Now lists all indexes for the table once per predicate-pushdown call and
picks the first SCALAR index whose keyColumnIds() contains the queried
column's field ID, same correctness behavior as before but no longer
coupled to the build procedure's naming choice.
The secondary-index sources had not been run through spotlessApply, so
spotlessJavaCheck failed in build-checks and in every module's check task.
Reformat them to the project style.

DurableIndexCatalog's class javadoc linked org.apache.iceberg.spark.
SparkIndexCatalogs, which core cannot resolve because core does not depend
on spark; this failed build-javadoc. Reference the class as plain text.

Add the missing Apache license header to the build_scalar_index_taxi.scala
demo script, which failed the rat license check.
If the key column is promoted (e.g. int to long) after a SCALAR index is
built, a HASH transform value recomputed from the current type no longer
matches the value stored in the index, so pruning would drop files that
actually contain matching rows and return a wrong (empty) result.

Add LeafFileReader.storedKeyType to read the key type physically stored in
a leaf file. SparkScanBuilder compares it against the current key type and,
when they differ, skips index pruning and falls back to normal planning,
logging a warning to rebuild the index.
The SCALAR index code lives only in spark/v3.5, but the default
spotlessApply configures the default Spark version (4.1), so these files
were never formatted and failed spotlessJavaCheck in build-checks and the
3.5 spark-tests job (which fail-fast-cancelled the other Spark versions).
Run spotlessApply for the 3.5 modules.
The index code had never been through checkstyle (only build-checks runs it,
and it was failing earlier on spotless/javadoc). Resolve the violations:
relocated Guava imports, Guava factory methods over new HashMap/ConcurrentHashMap,
descriptive local/parameter/member names, braces on single-line if blocks,
before() test setup naming, assertThatThrownBy message assertions, and an
unclosed HTML tag in an IndexIdentifier javadoc.
Decompose tryPruneUsingScalarIndex into focused helpers (predicate
resolution per operation, the type-promotion guard, and file-path
resolution) so each stays within the complexity and length limits, and
read the tracking file once per query, reusing the entries for both the
type-promotion check and candidate collection. Previously the tracking
file was read once for the type check plus once per target range during
candidate collection -- four reads for a three-value IN predicate.
Leaf files were written with Parquet's default 128 MB row-group size, so a
multi-million-entry leaf was only a few huge row groups and the sorted-key
statistics filter in LeafFileReader could not skip -- a point lookup
generically decoded roughly a whole row group (~1M+ entries), dominating
query planning at hundreds of ms per candidate leaf.

Default the leaf-file row-group size to 1 MB so a lookup prunes to one small
row group, and expose a leaf.row-group-size-bytes build option to tune it.
Measured on a 10M-row table, this cut per-leaf resolution from ~600 ms to
~7 ms and made indexed point/IN lookups faster than a full scan.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

API core data spark Specification Issues that may introduce spec changes.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant