Repository navigation
Conversation
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.
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.
Author
|
cc @huaxingao |
…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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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 tableIN, and range predicates (<,<=,>,>=,including a
BETWEEN, which Spark decomposes into two range predicates on the same column)against an existing index --
INand equality work against HASH or IDENTITY (each literalresolves 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
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
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
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
options => map('mode', 'incremental')builds leaf filesonly for data added since the index's last snapshot, adopting the append-only option from
Huaxin Gao's PK Index proposal (Section 7.2)
pointer file per index under the table's own location, scoped by table UUID)
configurable cap (default 100,
scalar-index.max-candidate-leaf-filestable property) fallsback to normal planning rather than risk making planning slower than no index at all
the table's partition spec logs a warning rather than rejecting -- redundant overhead, not wrong
IndexCatalog#listIndexesand matches by
keyColumnIds(), matching the design doc, instead of a hardcoded<column>_idxname lookupTest 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 thetable's FileIO supports prefix listing (
SupportsPrefixOperations)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
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
deferred since a bug there would silently return wrong (empty) results, not just miss an
optimization
scans are untouched
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