Architecture
How Ariadne is organized internally — the trait stack, what each layer is responsible for, on-disk storage, the data flows behind each public operation, and the design decisions worth knowing if you're contributing code.
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.
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.
| Path | Kind | Purpose |
|---|---|---|
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
- Pre-flight scan.
analyzeFilesreads distinct counts so the batcher can avoid producing a batch that pushes a column overlargeIndexLimit. - Optimal batching.
createOptimalBatchesgroups files into batches sized to keep memory bounded and large-index splits predictable. - Per-batch build.
Read files → apply computed indexes →
classifyLargeFilesmarks the(file, column)pairs that skip inline array storage → build regular, exploded, bloom, temporal, range, and auto-bloom DataFrame columns → append tostaging/and stream large pairs into eachlarge_indexes/{column}. - Periodic consolidation.
Every
stagingConsolidationThresholdbatches,consolidateStagingmerges staging into the main index via DeltaMERGE. IfautoCompactThresholdis set and reached, anOPTIMIZEruns too. - Final consolidation.
At the end, any remaining staged data is merged in and the staging table is deleted.
join
- Locate candidate files.
locateFilesFromDataFrameintersects per-index-type candidate sets: regular array matches, bloom probes, temporal matches, range overlap, and auto-bloom pre-filter for any large-index column. - Read the filtered files.
readFiles(candidateSet)goes through the full read pipeline (format → computed → exploded → select). - Temporal deduplication.
For any join column that is temporally indexed, keep only the latest row per key.
- Standard Spark join.
The result joins against the user's DataFrame on the requested keys with the requested join type.
deleteFiles
- Acquire update lock.
Single-writer guarantee while we touch Delta state.
- Read file sizes from the main index.
So we can correctly decrement
total_indexed_file_size. - Merge-delete from the main index.
Removes rows matching the filenames.
- Merge-delete from each large-index table.
One per column.
- Merge-delete from staging (if present).
In case a previous update was interrupted mid-flight.
- 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:
| Version | Adds |
|---|---|
| v1 → v2 | computed_indexes |
| v2 → v3 | exploded_field_indexes |
| v3 → v4 | read_options |
| v4 → v5 | bloom_indexes |
| v5 → v6 | temporal_indexes |
| v6 → v7 | range_indexes |
| v7 → v8 | auto_bloom_indexes |
| v8 → v9 | total_indexed_file_size (per-file size tracking, pruning metrics) |
| v9 → v10 | batches_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 version | Invariant |
|---|---|
| v1 | Alpha37 compatibility baseline; exploded columns may use array_column and file_size is not required. |
| v2 | Main and staging Delta tables contain non-null file_size BIGINT. |
| v3 | Exploded 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 by | Support | Operational preflight |
|---|---|---|
0.0.1-alpha-37 through 0.0.1-alpha44 | Supported | Validates 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-beta | Supported | Structurally 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 v2 | Supported | Resumes the ordered, idempotent migration from the declared checkpoint. |
| Explicit storage v4 | Current | Uses the declared current layout without repeating historical schema inspection. |
Earlier than 0.0.1-alpha-37 | Unsupported | Rebuild 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 v4 | Requires a newer library | Fails 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.
select()mutates internalselectedColumnsstate and returnsthis, not a copy._metadatais a mutablevarread and written from multiple call sites without synchronization.IndexMetadatauses mutable Java collections that are not thread-safe.
Safe patterns
- One thread per
Indexinstance. - Separate
Indexinstances per thread — they coordinate via the file-based locks. - Read-only queries (
locateFiles,stats,printIndex) never modify index contents, so they are safe across threads as long as nothing is mutating concurrently. They are not, however, purely read-only against storage: each runs the migration preflight first, which can acquire the update lock and write migrated data and metadata on the first call against a non-current index.
Lock file concurrency
File-based locks coordinate across separate Spark jobs and applications:
- Atomic acquisition via Hadoop
create(overwrite=false). - Stale-lock auto-healing has a known TOCTOU window, mitigated by a retry-depth guard.
- Lock refresh during long
updatecalls prevents false stale detection.
Driver memory considerations
A few operations collect data to the driver. They're documented in scaladoc with @note or inline comments:
- Unbounded:
IndexQueryOperations.collectFilenamesViaStaging()collects the full distinct filename set. This is the one genuinely unbounded driver-side collect on the query path — an index with millions of files can OOM the driver. - Bounded: bloom and auto-bloom probes collect only query-side values, capped by the derived probe bound (
BloomFilterOperations.maxProbeValues). Bloom filter binary data is never collected —probeBloomFiltersbroadcasts the values and evaluates the filters on the executors, returning only matching filenames. - Bounded:
IndexQueryOperations.getAutoBloomCandidates()probes distributed; its strategy-selection collect is capped one value past the threshold. - Bounded:
locateFilesWithRangeFromDataFrame()truncates to 10,000 distinct values. - Bounded by staging size:
Index.stats()collects the staging table's contents when one exists, so it scales with unconsolidated staging rather than with the index.
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:
- Metadata operations — extend
IndexMetadataOperations. - File reading — extend
IndexFileOperations. - Bloom-filter work — extend
BloomFilterOperations. - Index building — extend
IndexBuildOperations. - Join behavior — extend
IndexJoinOperations. - Query / introspection — extend
IndexQueryOperations. - Path or naming utilities — add to the
IndexPathUtilsobject.
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.