Overview

The Ariadne index system is a modular Scala layer over Apache Spark and Delta Lake. The single public type, Index, is a case class that extends a stack of traits — each adding one tightly-scoped layer of behavior to the one below it. There's no inheritance branching; the chain is strictly linear, which keeps testing focused and dependencies obvious.

Around the trait stack sit a few utility objects (IndexPathUtils, IndexCatalog, SchemaHelper), supporting types for tracked files and locks (FileList, IndexLock), and the metadata container (IndexMetadata) that round out the system.

Trait stack

Each trait below builds on the one above using self: Index =>, so every trait can call into Index's case-class members but only the bottom layer (AriadneContextUser) directly touches Spark.

Index (case class · public API)
IndexQueryOperations
IndexJoinOperations
IndexBuildOperations
BloomFilterOperations
IndexFileOperations
IndexMetadataOperations
AriadneContextUser
IndexPathUtilsutility object
FileListfile tracking
IndexLockconcurrency
IndexCatalogglobal directory

Storage layout

Everything Ariadne owns lives under spark.ariadne.storagePath/{indexName}/. The layout is deliberately flat: metadata is one JSON file, hot index data is one Delta table, anything that explodes too large gets its own Delta table.

PathKindPurpose
metadata.json JSON file Index config: schema, format, index types, read options
index/ Delta table Main index — value arrays, range structs, bloom bytes
staging/ Delta table Transient — holds batches during update before merge
large_indexes/{column}/ Delta table (per column) One per column with at least one file at or above largeIndexLimit distinct values — (filename, value) rows
.filelist.lock Lock file Held during addFile
.update.lock Lock file Held during update, deleteFiles, compact, vacuum

The list of tracked files is kept separately by FileList (a tiny two-column Delta table: filename, addedAt) so adding files doesn't require touching the main index.

The traits in detail

AriadneContextUser

The base. Provides the SparkSession, Hadoop FileSystem handle, the resolved storagePath, and a small set of filesystem and Delta convenience methods (exists, delete, open, delta(path)). Every other trait sits on top of this.

IndexMetadataOperations

Reads and writes metadata.json, validates the format string against the supported set, caches the parsed metadata, and exposes the public format and refreshMetadata() accessors. Schema serialization (round-tripping StructType through JSON) lives here.

IndexFileOperations

Owns the read pipeline. Given a Set[String] of file paths, it produces a DataFrame that has had: format-specific read options applied, computed-index columns added, array fields exploded into their virtual columns, and (if the user called select) been projected down to the requested columns. It also adds a stable filename column via input_file_name() when needed for downstream joins.

BloomFilterOperations

Build and query layer for Guava-backed bloom filters. Builds binary bloom_* columns during update, deserializes and tests them at query time, and exposes the set of columns that have bloom indexes. The implementation distinguishes explicit bloom indexes (added via addBloomIndex) from auto-bloom filters added automatically to large regular, computed, exploded-field, and temporal index columns (range indexes and binary-typed columns are excluded) — both are handled here.

IndexBuildOperations

The most substantial layer. Handles the full update pipeline: pre-flight file analysis, optimal batching, building each index type into DataFrame columns, splitting columns into the large-index pathway when they exceed thresholds, appending to staging and large-index tables, and consolidating staging into the main index via Delta MERGE.

IndexJoinOperations

Provides the optimized join entry point and the internal joinDf that consumers like the SQL catalog use. Maps the join columns to their storage-side names (e.g. user_id → bloom_user_id for bloom indexes; exploded-field aliases back to their array column), applies temporal deduplication for any temporal-indexed join column, and produces the joined DataFrame.

IndexQueryOperations

File location and statistics. locateFiles(map) and locateFilesFromDataFrame intersect the candidate-file sets across every index type involved in a query (regular, bloom, temporal, range, auto-bloom). stats() assembles a per-column statistics DataFrame. Internal helpers like collectFilenamesViaStaging materialize the distinct filename set to a temporary CSV so the driver collects from a trivial scan instead of re-running the full candidate plan. The filename set itself is still collected to the driver, so it remains the one genuinely unbounded driver-side collect in the query path.

Index (the case class)

The public API surface. Combines every trait above, plus the user-facing methods that don't belong to any one layer: addIndex, addBloomIndex, addComputedIndex, addExplodedFieldIndex, addTemporalIndex, addRangeIndex, addFile, update, deleteFiles, compact, vacuum, select, indexes, hasFile.

Factory methods live in the companion object:

Index(name, schema, format)
Index(name, schema, format, allowSchemaMismatch)
Index(name, schema, format, readOptions)
Index(name, schema, format, allowSchemaMismatch, readOptions)
Index(name)  // reconnect to an existing index — metadata must exist

Supporting types

IndexPathUtils

Stateless object for path-shaped utilities: sanitizing file names for storage, resolving the storage path, checking if an index exists, removing an index. Used everywhere paths are constructed; centralizing it keeps the path conventions in one place.

IndexCatalog

Global, stateless directory of every index under storagePath. Every call rescans the filesystem — there's no in-memory cache to invalidate. Exposes list, exists, describe, describeAll, get, toDF, remove, and findIndexes(file). The companion IndexSummary case class is the projection used by describe.

FileList

The tracked-file ledger for a single index. Backed by a two-column Delta table (filename: String, addedAt: Timestamp). Add and remove operations are merge-based; hasFile reads the table. Held under its own lock so updates to the file list don't block the main update lock.

IndexLock

File-based mutex implemented on top of Hadoop create(overwrite=false) for atomic acquisition. Lock files carry a JSON payload — correlation ID, acquisition time, last-refresh time, owner — used to detect stale locks and auto-heal them after lockTimeout seconds. Long-running operations call refresh() periodically so an honest in-progress job isn't mistaken for a crash.

{
  "correlationId": "uuid",
  "acquiredAt": "ISO-8601",
  "lastRefreshedAt": "ISO-8601",
  "owner": "spark-app-id"
}

IndexMetadata

The case class persisted to metadata.json. Uses Java collections (util.List, util.Map) for Gson serialization compatibility — don't switch these to Scala collections without rewriting the serializer. Supports forward migration from older metadata versions via null-checks in IndexMetadata.apply(jsonString).

SchemaHelper

Small utility object for schema introspection: fieldExists and fieldType, both with nested-field-path support.

Data flows

update

  1. Pre-flight scan.

    analyzeFiles reads distinct counts so the batcher can avoid producing a batch that pushes a column over largeIndexLimit.

  2. Optimal batching.

    createOptimalBatches groups files into batches sized to keep memory bounded and large-index splits predictable.

  3. Per-batch build.

    Read files → apply computed indexes → classifyLargeFiles marks the (file, column) pairs that skip inline array storage → build regular, exploded, bloom, temporal, range, and auto-bloom DataFrame columns → append to staging/ and stream large pairs into each large_indexes/{column}.

  4. Periodic consolidation.

    Every stagingConsolidationThreshold batches, consolidateStaging merges staging into the main index via Delta MERGE. If autoCompactThreshold is set and reached, an OPTIMIZE runs too.

  5. Final consolidation.

    At the end, any remaining staged data is merged in and the staging table is deleted.

join

  1. Locate candidate files.

    locateFilesFromDataFrame intersects per-index-type candidate sets: regular array matches, bloom probes, temporal matches, range overlap, and auto-bloom pre-filter for any large-index column.

  2. Read the filtered files.

    readFiles(candidateSet) goes through the full read pipeline (format → computed → exploded → select).

  3. Temporal deduplication.

    For any join column that is temporally indexed, keep only the latest row per key.

  4. Standard Spark join.

    The result joins against the user's DataFrame on the requested keys with the requested join type.

deleteFiles

  1. Acquire update lock.

    Single-writer guarantee while we touch Delta state.

  2. Read file sizes from the main index.

    So we can correctly decrement total_indexed_file_size.

  3. Merge-delete from the main index.

    Removes rows matching the filenames.

  4. Merge-delete from each large-index table.

    One per column.

  5. Merge-delete from staging (if present).

    In case a previous update was interrupted mid-flight.

  6. Update metadata, remove from FileList, release lock.

compact & vacuum

compact() runs Delta OPTIMIZE on the main index table and every large-index table under the update lock. vacuum(retentionHours) does the same with Delta VACUUM. Auto-compaction is driven by autoCompactThreshold and a batches_since_compact counter persisted in metadata, so the counter survives across separate Spark jobs.

Key features

Large index handling

When a file contributes largeIndexLimit or more distinct values to a column, that pair can no longer live as a Spark array column in the main index — the array would blow memory at both build and query time. Instead, the column gets its own consolidated Delta table at large_indexes/{column}/ containing (filename, value) rows, and the inline array is left null. Delta optimization works cleanly on this shape; compact() runs plain OPTIMIZE bin-packing over it.

Classification is per (file, column) pair and is driven by the distinct counts already gathered by analyzeFiles (see classifyLargeFiles). It deliberately keys on the distinct value count rather than the row count, so a file with many duplicate rows but few distinct values stays inline. Because largeness is known before aggregation, the values for a large pair are streamed straight into large_indexes/{column}/ from the source rows — the giant per-file array is never materialized.

Staged append + consolidation

Merging into the main index on every batch is expensive: each MERGE rewrites Delta files that grow as the index grows, so per-batch merge cost climbs as the update proceeds. Ariadne instead appends every batch to a staging/ Delta table and only merges into the main index once per stagingConsolidationThreshold batches (plus a final consolidation at the end). For a 180-batch update with the default threshold of 50, that's 4 merges instead of 180.

Staging also acts as a checkpoint: a crash mid-update leaves consolidated work in the main index and any in-flight work in staging, which the next update picks up.

Bloom filter indexes

Stored as binary columns in the main index Delta table with a bloom_ prefix. The serialized payload is the raw bytes of a Guava BloomFilter[CharSequence]. Lower fpr = larger filters but fewer wasted file reads.

Build. Filters are folded with BloomFilterAggregator, a Spark Aggregator whose buffer is the filter: values are hashed in and discarded, so executor memory tracks the filter size (~1.2 bytes per distinct value at 1% FPR) rather than the number of values. Because a bloom column stores only the filter — never the values — this removes the only cardinality ceiling a bloom index had. Each file's filter is sized from that file's own distinct count, which is carried on every input row since an Aggregator cannot see its grouping key.

Query. The probe values are broadcast and the filters are tested on the executors, so only matching filenames come back to the driver. Explicit bloom indexes and auto-bloom pre-filters share one bounded-collect helper and differ only in how a null filter is read: a missing auto-bloom means the file was never large enough to get one and stays a candidate, while a missing explicit bloom means the file held no values and cannot match. A filter prunes (1 − fpr)n of non-matching files, so pruning power collapses as the probe set grows; past the point where it would prune less than 1% the pre-filter is skipped rather than truncated. Skipping only widens the candidate set, which the subsequent read and join narrow correctly, whereas truncating would prune away files holding the dropped values. The bound is derived from the filter’s FPR, not configured.

String representation. Build and query must agree byte-for-byte on how a value is stringified, or the no-false-negatives guarantee breaks. The query side stringifies driver-side JVM objects with toString, so the build side does the same via canonicalStringColumn. A plain cast(StringType) would not be safe: Spark renders a timestamp as 2024-01-01 00:00:00 while java.sql.Timestamp.toString renders 2024-01-01 00:00:00.0.

Range indexes

Per-file min and max stored as a struct column in the main index with a range_ prefix. Range columns are intentionally kept out of the regular storageColumns set — they have their own rangeStorageColumns handling because they're struct-typed, not array-typed. At query time, files whose range doesn't overlap any query value are dropped before reads start.

Auto-bloom for large indexes

Querying a large_indexes/{column} table is more expensive than reading an array column, because each row is one filename/value pair. Auto-bloom adds a bloom filter (column-prefixed auto_bloom_) to the main index the first time any file crosses largeIndexLimit for that column; from then on every file gets a filter. At query time, the bloom pre-filters candidate files before the more expensive large-index scan.

Auto-bloom filters are built from the same source (filename, value) rows that feed large_indexes/, using the streaming aggregator described above. A large file has no inline array to read from, and building from rows means auto-bloom inherits the same absence of a cardinality ceiling.

Build flow is integrated into the staged-append path: bloom filters are added as a binary column to the staging DataFrame, and metadata is only updated after data is safely staged.

The surviving candidates are applied to the large-index scan as filename IN (…) (see pruneLargeIndexRows). Delta evaluates an IN list against file statistics as a range rather than as set membership, so file-level skipping keeps every data file whose filename min/max overlaps the span between the lowest and highest candidate — a scattered candidate set therefore prunes far fewer files than its size suggests. The saving that does materialize comes mostly from row-group filtering within those files, which is effective because each source file's rows stay contiguous through bin-packing. An empty candidate set is the case that pays off fully: the large index is then skipped outright.

Index repartitioning for very large joins

The explode step on index array columns can produce massive intermediate DataFrames that cause FetchFailedException on shuffles. The optional indexRepartitionCount setting repartitions the index DataFrame both before explode in locateFilesFromDataFrame and after readFiles in joinDf. It's not enabled by default because the extra shuffle hurts small joins.

Deterministic staging recovery

Each staged batch records an internal timestamp and batch ID. Consolidation ranks duplicate filenames by newest batch, row completeness, batch ID, and a stable payload hash, then removes the internal columns before merging into the main index. Legacy staging tables without ordering columns use completeness and payload hash only. update always consolidates stale staging under the update lock, even when no new files or column backfills are pending.

Metadata versioning

IndexMetadata.apply(jsonString) migrates old metadata files forward via sequential null-checks. When adding a new field, add a corresponding null-check migration step. Current migration history:

VersionAdds
v1 → v2computed_indexes
v2 → v3exploded_field_indexes
v3 → v4read_options
v4 → v5bloom_indexes
v5 → v6temporal_indexes
v6 → v7range_indexes
v7 → v8auto_bloom_indexes
v8 → v9total_indexed_file_size (per-file size tracking, pruning metrics)
v9 → v10batches_since_compact (cross-job auto-compaction counter)

Physical storage versions

storage_format_version versions the complete persisted index layout independently from additive metadata fields. Indexes created before explicit versioning are treated as storage v1 after baseline validation. Migration preflight acquires the update lock, applies each step idempotently, verifies the physical result, and writes the new version only after all steps succeed.

Storage versionInvariant
v1Alpha37 compatibility baseline; exploded columns may use array_column and file_size is not required.
v2Main and staging Delta tables contain non-null file_size BIGINT.
v3Exploded columns and large-index paths use canonical as_column names.
v4 (current)Every column with a large_indexes/ table is registered in auto_bloom_indexes, stores its values under the canonical column name, and carries a non-null auto_bloom_{column} filter in the main table for every file holding at least one non-null value. A file holding no value for the column keeps a null filter and is treated as an unconditional match candidate.

Plain Index(name) construction remains read-only. Operational entry points run preflight; versions newer than the current build raise a typed compatibility exception instead of returning incomplete results.

The compatibility floor is covered by a physical fixture under src/test/resources/fixtures/alpha37/, generated from tag 0.0.1-alpha-37 with dev/scripts/generate-alpha37-fixture.sh. Normal tests use only the committed fixture and require no network access.

Supported release cohorts

Persisted bySupportOperational preflight
0.0.1-alpha-37 through 0.0.1-alpha44SupportedValidates the unversioned v1 baseline, backfills file sizes, normalizes exploded aliases, verifies the result, then records v4. Auto-bloom filters are backfilled for any column that already has a large_indexes/ table.
0.1.0-beta through 0.1.10-betaSupportedStructurally validates the unversioned layout, applies only missing steps, verifies the result, then records v4. Auto-bloom filters are backfilled for any column that already has a large_indexes/ table.
Explicit storage v1 or v2SupportedResumes the ordered, idempotent migration from the declared checkpoint.
Explicit storage v4CurrentUses the declared current layout without repeating historical schema inspection.
Earlier than 0.0.1-alpha-37UnsupportedRebuild from source data with a current Ariadne version. Note that Ariadne cannot detect this case: unversioned metadata is assumed to be the v1 baseline, so migration is attempted and may fail during a verification step rather than being rejected up front.
Version newer than v4Requires a newer libraryFails before physical or metadata writes; never downgrades an index in place.

This build supports persisted indexes from 0.0.1-alpha-37 through 0.1.10-beta. Storage v2, v3 and v4 are migration checkpoints, not promises that every historical release wrote an explicit version. All releases before the version framework wrote unversioned metadata and are classified by the validated physical layout. Future releases may introduce newer formats that require a newer Ariadne build.

Thread safety

Index instance thread safety

Index instances are not safe for concurrent mutation.

Safe patterns

Lock file concurrency

File-based locks coordinate across separate Spark jobs and applications:

Driver memory considerations

A few operations collect data to the driver. They're documented in scaladoc with @note or inline comments:

FileList.addFile() does not collect — it deduplicates with a Delta MERGE against the tracked file list.

Extending Ariadne

The trait stack is intentionally easy to extend in place:

Each trait has a focused test suite. For example, join changes live in IndexJoinOperationsTests:

class IndexJoinOperationsTests extends SparkTests {
  test("join should cache results for repeated operations") {
    val index = Index("test", schema, "parquet")
    index.addIndex("id")
    // ...
  }
}

New error conditions get a typed exception extending AriadneException — see the table on the Troubleshooting page for the existing set and conventions.

Logging convention: use logger.warn(...) for normal operational messages. This is intentional — Spark surfaces warn-level output in cluster environments. Reserve debug for verbose entries, and never silently swallow exceptions: log them, then rethrow or handle.