case class Index extends IndexQueryOperations with Product with Serializable
Represents an Index for managing metadata and file-based indexes in Apache Spark.
This class provides methods to add, locate, and manage file-based indexing in Spark using Delta Lake. It supports schema enforcement, metadata persistence, and file tracking.
- Note
Index instances are NOT safe for concurrent use from multiple threads. Each thread should use its own Index instance.
- Alphabetic
- By Inheritance
- Index
- Serializable
- Serializable
- Product
- Equals
- IndexQueryOperations
- IndexJoinOperations
- IndexBuildOperations
- BloomFilterOperations
- IndexFileOperations
- IndexMetadataOperations
- AriadneContextUser
- AnyRef
- Any
- Hide All
- Show All
- Public
- All
Type Members
-
case class
FileAnalysis(filename: String, distinctCounts: Map[String, Long], maxDistinctCount: Long) extends Product with Serializable
Case class to hold file analysis results for batching decisions.
Case class to hold file analysis results for batching decisions.
- filename
The name of the file
- distinctCounts
Map of column name to distinct value count for that file
- maxDistinctCount
Maximum distinct count across all indexed columns
- Definition Classes
- IndexBuildOperations
-
case class
LargeFileClassification(byColumn: Map[String, Set[String]]) extends Product with Serializable
Identifies which files contribute too many distinct values to a column to be stored inline.
Identifies which files contribute too many distinct values to a column to be stored inline.
A
(file, column)pair is "large" when the file contributes at leastlargeIndexLimitdistinct values to that column. Those values are streamed intolarge_indexes/{column}as individual rows instead of being collected into a per-file array, and the inline array column is leftnull.Classification is deliberately based on the distinct value count rather than the row count, so a file with many duplicate rows but few distinct values stays inline.
- byColumn
map of column name to the set of filenames that are large for that column; columns with no large files are omitted
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
Value Members
-
final
def
!=(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
-
final
def
##(): Int
- Definition Classes
- AnyRef → Any
-
final
def
==(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
-
def
addBloomIndex(column: String, fpr: Double = 0.01): Unit
Adds a bloom filter index for the specified column.
Adds a bloom filter index for the specified column.
Bloom filters are probabilistic data structures that provide:
- Guaranteed NO false negatives (if filter says "no", value definitely absent)
- Configurable false positive rate (if filter says "yes", value MIGHT be present)
- Space-efficient storage (approximately 10 bits per element at 1% FPR)
Binary columns cannot be bloom-indexed. Values are canonicalized to a string before hashing, and Spark materializes
BinaryTypeas a JVMArray[Byte]whosetoStringis identity-based ([B@1b6d3586), so the same bytes hash to a different token on every occurrence. Such a filter would match nothing, ever. Columns that contain a binary field at any depth are rejected for the same reason. Use addIndex on the column instead — regular indexes compare values, not their string form, and work correctly on binary data.- column
The column name to index with a bloom filter
- fpr
False positive rate between 0.0 and 1.0 (default 0.01 = 1%)
index.addBloomIndex("sessionId") index.addBloomIndex("ipAddress", 0.001)
- Exceptions thrown
ColumnNotFoundExceptionif column doesn't exist in schemaIllegalArgumentExceptionif column is null/blank, FPR out of range, column is a nested path, column is of a binary type, or column is already indexed by any other type
Example: -
def
addComputedIndex(name: String, sql_expression: String): Unit
Adds a computed index derived from a SQL expression.
Adds a computed index derived from a SQL expression.
The expression is evaluated during update to produce a virtual column whose distinct values are stored in the index. Idempotent: calling again with the same name is a no-op.
- name
The alias name for the computed column
- sql_expression
A Spark SQL expression evaluated against each data file
index.addComputedIndex("yearMonth", "date_format(event_date, 'yyyy-MM')")
- Exceptions thrown
IllegalArgumentExceptionif name or sql_expression is null/blank, or name conflicts with another index type
Example: -
def
addExplodedFieldIndex(arrayColumn: String, fieldPath: String, asColumn: String): Unit
Adds an exploded field index for a nested field inside an array column.
Adds an exploded field index for a nested field inside an array column.
During update, each element of
arrayColumnis exploded, thefieldPathis extracted, and the distinct values are stored underasColumnin the index. Joins onasColumnwill locate files containing any matching array element. Idempotent: calling again with the sameasColumnis a no-op.- arrayColumn
The array column to explode.
- fieldPath
The dot-separated field path to extract from each array element (e.g., "id" or "profile.user_id").
- asColumn
The virtual column name exposed for joins.
// Index the "id" field from each element of the "items" array column index.addExplodedFieldIndex("items", "id", "item_id")
- Exceptions thrown
IllegalArgumentExceptionif any parameter is null/blank, orasColumnconflicts with another index type
Example: -
def
addFile(fileNames: String*): Unit
Adds files to the index's file list for future indexing.
Adds files to the index's file list for future indexing. Acquires a file list lock to prevent concurrent modifications.
- fileNames
One or more file paths to register
index.addFile("/data/events/2024-01-01.parquet") index.addFile("/data/events/2024-01-02.parquet", "/data/events/2024-01-03.parquet")
- Exceptions thrown
IllegalArgumentExceptionif fileNames is null/empty or any fileName is null/blank
Example: -
def
addFilenameColumn(df: DataFrame, files: Set[String]): DataFrame
Adds a stable filename column to the DataFrame.
Adds a stable filename column to the DataFrame.
Uses
input_file_name()to capture the source file path. For single-file reads where Spark may report an empty filename, falls back to the known file path.- df
The DataFrame to add the filename column to
- files
The set of file paths being read
- returns
DataFrame with a
filenamecolumn added
- Attributes
- protected
- Definition Classes
- IndexFileOperations
-
def
addIndex(index: String): Unit
Adds a regular (array-of-distinct-values) index for the specified column.
Adds a regular (array-of-distinct-values) index for the specified column.
Idempotent: calling again with the same column is a no-op.
- index
The column name to index.
val index = Index("myIndex", schema, "parquet") index.addIndex("userId")
- Exceptions thrown
ColumnNotFoundExceptionif the column doesn't exist in the schemaIllegalArgumentExceptionif the column name is null/blank, is a nested path, or is already indexed by any other type
Example: -
def
addRangeIndex(column: String): Unit
Adds a range index for the specified column.
Adds a range index for the specified column.
Range indexes store min/max values per file, enabling file pruning at query time. Files whose [min, max] range does not overlap with the queried values are skipped.
- column
The column to index with min/max range
index.addRangeIndex("timestamp")- Exceptions thrown
ColumnNotFoundExceptionif column doesn't exist in schemaIllegalArgumentExceptionif column is null/blank or already indexed by another type
Example: -
def
addTemporalIndex(column: String, timestampColumn: String): Unit
Adds a temporal index for the specified column using a timestamp for versioning.
Adds a temporal index for the specified column using a timestamp for versioning.
When joining on a temporal index column, only the latest version (by timestamp) of each value is returned. This is useful when multiple files contain the same entity at different points in time.
The timestamp may be a nested path such as
meta.updatedAt, since it is only read during index construction and deduplication. The value column must be top-level: it is persisted as a column of the index table under its own name, so a dotted path could not be read back.- column
The top-level value column to index on (e.g., "user_id")
- timestampColumn
The timestamp column for ordering versions (e.g., "updated_at"), which may be a nested path
index.addTemporalIndex("userId", "updated_at") // the timestamp may also be a nested path index.addTemporalIndex("orderId", "meta.updatedAt")
- Exceptions thrown
ColumnNotFoundExceptionif either column doesn't exist in schemaIllegalArgumentExceptionif column or timestampColumn is null/blank, if column is a nested path, or if column is already indexed by another type
Example: -
def
analyzeFiles(files: Set[String]): Seq[FileAnalysis]
Performs pre-flight analysis on unindexed files to determine optimal batching strategy.
Performs pre-flight analysis on unindexed files to determine optimal batching strategy.
Reads the source files, applies computed indexes, and counts distinct values per indexed column per file. The resulting FileAnalysis objects are used by createOptimalBatches to group files into batches that stay under the
largeIndexLimit.If no storage columns are configured (e.g., only bloom or range indexes), returns trivial analyses with zero distinct counts.
Each column is counted on the same basis its builder uses, because these counts also drive classifyLargeFiles. Regular, computed, and temporal columns are counted before exploded fields are applied:
applyExplodedFieldsuses an innerexplode, which drops rows whose array is null or empty and would undercount every other column. Temporal columns add one for a present null value, matching thecollect_setofstruct(value, max_ts)in buildTemporalIndexes, which retains a struct for the null group.- files
set of file paths to analyze
- returns
sequence of FileAnalysis objects with per-column distinct counts; empty if
filesis empty
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- Note
This method calls
.collect()to bring per-file distinct counts to the driver. For indexes covering millions of files with many indexed columns, the collected result set can be large enough to cause driver OOM. Consider limiting the number of files analyzed per call or increasing driver memory.
-
def
appendToLargeIndex(sourceDf: DataFrame, large: LargeFileClassification, columns: Set[String]): Unit
Appends large index data to per-column Delta tables under
large_indexes/.Appends large index data to per-column Delta tables under
large_indexes/.For each column with at least one large file, the values of those files are read straight from the source DataFrame as distinct
(filename, value)rows and written tolarge_indexes/{column}/. Values are never collected into a per-file array first, so a file's cardinality is bounded by what Delta can store rather than by what fits in an executor-side array. Before appending, any existing rows for the same filenames are removed via Delta MERGE to prevent duplicates (important during column backfill or re-indexing).- sourceDf
the base DataFrame with
filenameand all source columns (computed indexes already applied)- large
the classification identifying which files are large for which columns
- columns
the columns to process: those with large files in this batch, plus any with an existing table that may hold stale rows for the files being re-indexed
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- Note
The
,count()call before the write is intentional: it materializes the DataFrame to determine whether any large-index rows exist for this column, avoiding the overhead of a Delta MERGE + write when no rows exist. This results in a double computation (count + write), but the alternative—writing unconditionally—would create empty Delta commits and unnecessary MERGE operations.Crash consistency: The MERGE that clears stale rows and the append that writes replacements are separate Delta commits, and this method runs before the matching main-index rows reach staging. A crash between any two of those commits therefore breaks the "inline or large, never both" invariant that loadLargeIndex readers rely on. For a file not yet in the main index this is self-correcting — it stays unindexed and the next update reprocesses it — but for a backfill over an already-indexed file the deletion can survive without its replacement, silently removing that column's values for that file. See Index.update for the user-facing statement of this limitation.
-
def
appendToStaging(df: DataFrame): Unit
Appends the processed DataFrame to the staging Delta table.
Appends the processed DataFrame to the staging Delta table.
Columns belonging to a large
(file, column)pair are alreadynullat this point: the index builders skip them and their values were streamed tolarge_indexes/instead. Staging deliberately does not re-derive largeness from array size. Doing so would be a second, independent decision, and any disagreement with classifyLargeFiles would null out an array whose values were never written anywhere else, silently losing them.- df
the combined index DataFrame to stage
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
applyColumnSelection(df: DataFrame, additionalColumns: Seq[String] = Seq.empty): DataFrame
Applies column selection if columns have been specified via select().
Applies column selection if columns have been specified via select(). Only reads the explicitly selected columns.
- df
The DataFrame to apply column selection to
- additionalColumns
Columns to retain alongside the user's selection. Entries absent from
dfare ignored.- returns
DataFrame with selected columns or original DataFrame if no selection
- Attributes
- protected
- Definition Classes
- IndexFileOperations
-
def
applyComputedIndexes(df: DataFrame): DataFrame
Applies computed indexes to a DataFrame by adding computed columns.
Applies computed indexes to a DataFrame by adding computed columns.
Each computed index is a Spark SQL expression stored in metadata. The expression is evaluated via
expr()and added as a new column. If no computed indexes are configured, the DataFrame is returned unchanged.- df
The base DataFrame
- returns
DataFrame with computed index columns added
- Attributes
- protected
- Definition Classes
- IndexFileOperations
- Exceptions thrown
org.apache.spark.sql.AnalysisExceptionif a computed index expression is invalid or references non-existent columns
-
def
applyExplodedFields(df: DataFrame): DataFrame
Applies exploded field transformations to a DataFrame.
Applies exploded field transformations to a DataFrame.
For each configured exploded field, extracts the nested field path from the array column and explodes it into a top-level column. This is used on the read/join path to prepare data for joining on exploded field values.
Note that
explode()(notexplode_outer()) is used intentionally — rows with null or empty arrays are dropped because they contain no matchable values for the join. This is consistent with the index build path which also excludes null/empty arrays.- df
The DataFrame to transform
- returns
DataFrame with exploded field columns added; rows with null/empty arrays are excluded
- Attributes
- protected
- Definition Classes
- IndexFileOperations
-
final
def
asInstanceOf[T0]: T0
- Definition Classes
- Any
-
val
autoBloomColumnPrefix: String
Column name prefix for auto-bloom filter storage in the main index table.
Column name prefix for auto-bloom filter storage in the main index table.
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
lazy val
autoBloomFpr: Double
False positive rate for auto-bloom filters on large index columns.
False positive rate for auto-bloom filters on large index columns. When a column exceeds largeIndexLimit, an auto-bloom filter is built with this FPR to enable file pruning at query time. Reads from spark.ariadne.autoBloomFpr configuration (default: 0.01).
- Definition Classes
- AriadneContextUser
-
def
autoBloomStorageColumns: Set[String]
Returns the set of auto-bloom storage column names (each prefixed with
auto_bloom_).Returns the set of auto-bloom storage column names (each prefixed with
auto_bloom_).- returns
set of prefixed column names (e.g.,
auto_bloom_user_id)
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
autoBloomValueColumn(column: String): Column
Returns the expression selecting the bloom-hashable scalar out of a columnValueRows row.
Returns the expression selecting the bloom-hashable scalar out of a columnValueRows row.
Every column type but temporal yields the value directly. Temporal rows carry a
(value, max_ts)struct, so the filter is folded over thevaluefield alone to match the bare scalars that queries probe with.- column
the storage column name, as produced by columnValueRows
- returns
the column expression to hash into the filter
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
lazy val
autoCompactThreshold: Option[Int]
Optional threshold for auto-compaction.
Optional threshold for auto-compaction. When set, every Delta table backing the index (main index and each large index) is compacted during an update once the persisted batches-since-compact counter meets or exceeds this value. The counter tracks update batches, not Delta log files. Reads from spark.ariadne.autoCompactThreshold configuration (default: not set).
- Definition Classes
- AriadneContextUser
-
val
batchesSinceCompact: Int
Counter tracking batches processed since the last auto-compaction.
Counter tracking batches processed since the last auto-compaction.
Incremented after each batch in
updateBatchedand reset to zero after each compaction cycle or at the start of each update call to prevent stale counts from a previousupdatefrom triggering premature compaction.- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
val
bloomColumnPrefix: String
Column name prefix used when storing bloom filter binary data in the index Delta table.
Column name prefix used when storing bloom filter binary data in the index Delta table.
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
bloomColumns: Set[String]
Returns the set of source column names that have bloom indexes configured.
Returns the set of source column names that have bloom indexes configured.
- returns
set of column names (without the
bloom_prefix)
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
bloomIndexConfigs: Seq[BloomIndexConfig]
Returns the bloom index configurations from the current index metadata.
Returns the bloom index configurations from the current index metadata.
- returns
sequence of BloomIndexConfig objects, one per bloom-indexed column
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
bloomStorageColumns: Set[String]
Returns the set of storage column names used for bloom filter binary data.
Returns the set of storage column names used for bloom filter binary data.
Each name is the source column name prefixed with bloomColumnPrefix.
- returns
set of prefixed column names (e.g.,
bloom_user_id)
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
buildAutoBloomIndexes(combinedDf: DataFrame, sourceDf: DataFrame, large: LargeFileClassification): DataFrame
Builds auto-bloom filters for columns that have at least one large file.
Builds auto-bloom filters for columns that have at least one large file.
A column becomes auto-bloom the first time any file contributes at least
largeIndexLimitdistinct values to it. From then on every file gets a filter, stored in the main index under theauto_bloom_prefix, so queries can cheaply skip files before touchinglarge_indexes/.Filters are folded directly from the source
(filename, value)rows with a streaming aggregator rather than from a collected array. Large files never have an array to read from, and building from the rows means a column's cardinality is limited only by what the filter itself costs (~1.2 bytes per distinct value at 1% FPR).- combinedDf
the combined index DataFrame, keyed by
filename- sourceDf
the base DataFrame with
filenameand all source columns (computed indexes already applied)- large
the classification identifying which files are large for which columns
- returns
combinedDfwith anauto_bloom_{column}binary column per auto-bloom column; unchanged if none qualify
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
buildBloomFilterIndexes(df: DataFrame): DataFrame
Builds bloom filter indexes for all configured bloom columns.
Builds bloom filter indexes for all configured bloom columns.
For each BloomIndexConfig, this method:
- Reduces the input to distinct (filename, value) pairs
- Computes each file's distinct value count so its filter can be sized individually
- Folds the values into a bloom filter with
BloomFilterAggregator - Joins the resulting bloom column back onto the accumulating DataFrame
The aggregation is streaming: values are hashed into the filter and discarded, so executor memory scales with the size of the filter (~1.2 bytes per distinct value at 1% FPR) rather than with the number of values held as boxed JVM objects. A bloom-indexed column is therefore not bounded by what fits in an executor-side array.
If no bloom indexes are configured, returns a distinct filename-only DataFrame.
- df
DataFrame containing a
filenamecolumn and all bloom-indexed source columns- returns
DataFrame with
filenameplus onebloom_{column}binary column per configured bloom index
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
buildExplodedFieldIndexes(baseData: DataFrame, resultDf: DataFrame, large: LargeFileClassification): DataFrame
Builds exploded field indexes and joins them onto the result DataFrame.
Builds exploded field indexes and joins them onto the result DataFrame.
For each configured exploded field index, extracts the nested field path from the array column, explodes it, collects distinct values back into an array per file, and joins the result onto
resultDf.- baseData
the full base DataFrame with all source columns
- resultDf
the accumulating result DataFrame to join with (must have
filename)- returns
DataFrame with exploded field index columns joined via
full_outer
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
buildRangeIndexes(df: DataFrame): DataFrame
Builds range indexes storing
Struct(min, max)per file.Builds range indexes storing
Struct(min, max)per file.For each range index configuration, groups by
filenameand computes theminandmaxof the column per file. The result is a struct column namedrange_{column}.- df
the base DataFrame with
filenamecolumn and range-indexed source columns- returns
DataFrame with
filenameplus one range struct column per configured range index; filename-only DataFrame if no range indexes are configured
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
buildRegularIndexes(df: DataFrame, large: LargeFileClassification): DataFrame
Builds regular (array-aggregated) indexes for all configured regular and computed columns.
Builds regular (array-aggregated) indexes for all configured regular and computed columns.
Groups the input data by filename and collects distinct values into arrays via
collect_set. If no regular or computed indexes are configured, returns a distinct filename-only DataFrame.- df
the base DataFrame with
filenamecolumn and all indexed source columns- returns
DataFrame with
filenameplus one array column per regular/computed index
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
buildStreamingBloomColumn(df: DataFrame, valueColumn: Column, bloomColumn: String, fpr: Double): DataFrame
Folds a column of values into one serialized bloom filter per file.
Folds a column of values into one serialized bloom filter per file.
The aggregation is streaming: values are hashed into the filter and discarded, so executor memory scales with the size of the filter (~1.2 bytes per distinct value at 1% FPR) rather than with the number of values held as boxed JVM objects. A bloom-backed column is therefore not bounded by what fits in an executor-side array.
Values are reduced to distinct
(filename, value)pairs first, so each value is inserted exactly once.- df
DataFrame containing a
filenamecolumn and the source values- valueColumn
the column expression holding the values to insert
- bloomColumn
name to give the resulting binary filter column
- fpr
desired false positive rate, strictly between 0 and 1
- returns
DataFrame of
filenameplus the serialized filter column; files contributing no non-null values are absent
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
buildTemporalIndexes(df: DataFrame, large: LargeFileClassification): DataFrame
Builds temporal indexes storing
Array[Struct(value, max_ts)]per file.Builds temporal indexes storing
Array[Struct(value, max_ts)]per file.For each temporal index configuration, groups by
(filename, value_column)to find the maximum timestamp per value per file, then collects the(value, max_ts)structs into an array per file. This enables temporal deduplication at query time.- df
the base DataFrame with
filename, value, and timestamp columns- returns
DataFrame with
filenameplus one struct-array column per temporal index; filename-only DataFrame if no temporal indexes are configured
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
canonicalStringColumn(column: Column): Column
Converts a column of any type to the canonical string form used for bloom hashing.
Converts a column of any type to the canonical string form used for bloom hashing.
Build and query sides must agree exactly on this representation: a bloom filter guarantees no false negatives, but only if the same value hashes to the same string in both places. The query side stringifies driver-side JVM objects with
toString, so this wraps the value in a single-element array and appliestoStringinside a UDF. Spark converts array elements to the same external JVM types (java.sql.Timestamp,java.math.BigDecimal, and so on) that the driver sees, making the two paths byte-identical for scalar types.This does not hold for types whose
toStringis not value-based.Array[Byte].toStringis an identity hash, so two equal binary values stringify differently on the build and query sides and every probe misses. Binary columns are therefore rejected by Index.addBloomIndex and excluded from auto-bloom. Nested array and map types are untested and carry the same risk.A plain
cast(StringType)would not be safe here — Spark renders a timestamp as2024-01-01 00:00:00whilejava.sql.Timestamp.toStringrenders2024-01-01 00:00:00.0, which would silently produce false negatives.- column
the source column, of any type
- returns
a string column holding the canonical form,
nullwhere the source isnull
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
classifyLargeFiles(analyses: Seq[FileAnalysis]): LargeFileClassification
Classifies each
(file, column)pair as large or inline using pre-flight distinct counts.Classifies each
(file, column)pair as large or inline using pre-flight distinct counts.- analyses
per-file analyses produced by analyzeFiles
- returns
the resulting LargeFileClassification; empty when no pair reaches
largeIndexLimit
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
clone(): AnyRef
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws( ... ) @native() @HotSpotIntrinsicCandidate()
-
def
columnValueRows(df: DataFrame, column: String): Option[DataFrame]
Produces the
(filename, value)rows for an index column, one row per occurrence.Produces the
(filename, value)rows for an index column, one row per occurrence.Rows are produced directly from the source data without aggregating into arrays. This is the value source for both
large_indexes/{column}and auto-bloom filters. Callers are responsible for anydistinctor filename filtering they need.Regular and computed indexes project the column directly; exploded-field indexes explode the configured nested path; temporal indexes reduce to one
(value, max_ts)struct per distinct value.- df
the base DataFrame with a
filenamecolumn, computed indexes applied, and all source columns present- column
the storage column name
- returns
SomeDataFrame offilenamepluscolumn, orNonewhencolumnis not a known storage column
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
compact(): Unit
Compacts all Delta tables belonging to this index using OPTIMIZE.
Compacts all Delta tables belonging to this index using OPTIMIZE. Acquires the update lock to prevent concurrent modifications.
index.compact()
- Exceptions thrown
IndexLockExceptionif the update lock cannot be acquired within the configured timeout
Example: -
def
compactDeltaTables(): Unit
Compacts all Delta tables (main index and large index tables) using Delta OPTIMIZE.
Compacts all Delta tables (main index and large index tables) using Delta OPTIMIZE.
Runs Delta Lake's OPTIMIZE command to consolidate small files into larger ones, improving read performance. Processes the main index table first, then each large index column table.
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
consolidateStaging(): Unit
Consolidates staged data into the main index table.
Consolidates staged data into the main index table.
Delegates to consolidateMainStaging which performs a Delta MERGE (upsert) of all staged rows into the main index, then deletes the staging table. Logs the total time taken.
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
createBaseDataFrame(files: Set[String]): DataFrame
Creates a base DataFrame from the provided files using the stored format and schema.
Creates a base DataFrame from the provided files using the stored format and schema. Applies any read options configured in the index metadata.
Returns an empty DataFrame (with the stored schema) if all file paths are blank after normalization.
- files
Set of file paths to read
- returns
DataFrame with data from the specified files
- Attributes
- protected
- Definition Classes
- IndexFileOperations
- Exceptions thrown
IllegalArgumentExceptionif the stored format is not csv, parquet, or json
-
def
createOptimalBatches(fileAnalyses: Seq[FileAnalysis]): Seq[Set[String]]
Groups files into batches whose aggregate
maxDistinctCountstays at or belowlargeIndexLimit.Groups files into batches whose aggregate
maxDistinctCountstays at or belowlargeIndexLimit.The algorithm sorts files by
maxDistinctCount(largest first) and packs them sequentially into batches. When adding a file would push the batch total past the limit, a new batch is started. Files that individually exceed the limit are placed into single-file batches to guarantee progress.- fileAnalyses
sequence of FileAnalysis objects from analyzeFiles
- returns
sequence of file batches (each a
Set[String]of filenames); empty iffileAnalysesis empty
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
lazy val
debugEnabled: Boolean
When true, logs detailed diagnostics during join operations including partition counts, physical plans, and per-phase timing.
When true, logs detailed diagnostics during join operations including partition counts, physical plans, and per-phase timing. Reads from spark.ariadne.debug configuration (default: false).
- Definition Classes
- AriadneContextUser
-
def
delete(path: Path): Boolean
Deletes a path recursively from the filesystem.
Deletes a path recursively from the filesystem.
- path
The Hadoop Path to delete
- returns
true if the path was successfully deleted
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulljava.io.IOExceptionif the filesystem operation fails
-
def
deleteFiles(filenames: String*): Unit
Deletes the specified files from the index, large index tables, and file list.
Deletes the specified files from the index, large index tables, and file list.
Acquires the update lock, removes matching rows from the main index Delta table, all large index Delta tables, and the FileList. If a filename doesn't exist in the index, it is silently ignored.
- filenames
One or more filenames to remove from the index.
index.deleteFiles("/data/events/2024-01-01.parquet") index.deleteFiles("/data/old1.parquet", "/data/old2.parquet")
- Exceptions thrown
IllegalArgumentExceptionif filenames is null/empty or contains null/blank entriesIndexLockExceptionif the update lock cannot be acquired within the configured timeout
Example: -
def
delta(path: Path): Option[DeltaTable]
Returns a io.delta.tables.DeltaTable if the path exists and contains valid Delta metadata.
Returns a io.delta.tables.DeltaTable if the path exists and contains valid Delta metadata.
Three outcomes are possible:
- the path holds a valid Delta table —
Some(DeltaTable); - the path is absent, or exists but is an empty directory —
None, so the caller may create the table there; - the path exists, is non-empty, and holds no readable Delta metadata — dev.cjfravel.ariadne.exceptions.InvalidDeltaTableException.
The third case is reported rather than repaired. Ariadne cannot distinguish an abandoned partial write of its own from data written by something else, so it never deletes the directory; cleanup is left to the caller, who can inspect the contents first. An empty directory is treated as absent because there is nothing to inspect and nothing to lose.
- path
The Hadoop Path to check for a Delta table
- returns
Some(DeltaTable)if a valid Delta table exists,Noneif the path is absent or an empty directory
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulldev.cjfravel.ariadne.exceptions.InvalidDeltaTableExceptionif the path exists and is non-empty but is not a readable Delta table
- the path holds a valid Delta table —
-
def
ensureStorageReadyUnderLock(lock: IndexLock, correlationId: String, refresh: Boolean): Unit
Runs migration preflight while the caller holds the update lock.
Runs migration preflight while the caller holds the update lock.
- refresh
whether to reload metadata after lock acquisition
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
final
def
eq(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
-
def
exists(path: Path): Boolean
Checks if a path exists on the filesystem.
Checks if a path exists on the filesystem.
- path
The Hadoop Path to check
- returns
true if the path exists
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulljava.io.IOExceptionif the filesystem operation fails
-
def
format: String
Returns the file format of the indexed data (e.g., "parquet", "csv", "json").
Returns the file format of the indexed data (e.g., "parquet", "csv", "json").
- returns
the format string from the stored metadata
- Definition Classes
- IndexMetadataOperations
-
lazy val
fs: FileSystem
Hadoop
FileSysteminstance resolved from the storagePath URI and the Spark Hadoop configuration.Hadoop
FileSysteminstance resolved from the storagePath URI and the Spark Hadoop configuration. Used for all filesystem operations (existence checks, reads, deletes, lock files).- Definition Classes
- AriadneContextUser
-
def
getAutoBloomCandidatesFromDf(storageColumn: String, joinColumn: String, valuesDf: DataFrame, indexDf: DataFrame): Option[Set[String]]
Gets candidate files from the auto-bloom filter using values extracted from a DataFrame.
Gets candidate files from the auto-bloom filter using values extracted from a DataFrame.
Delegates the bounded query-side collect to BloomFilterOperations.collectProbeValues, the same helper the explicit bloom path uses, then probes via getAutoBloomCandidates. The two bloom kinds differ only in how a
nullfilter is interpreted: a missing auto-bloom means the file was never large enough to get one and must stay a candidate, whereas a missing explicit bloom means the file held no values for the column and cannot match.- storageColumn
The storage column name in the index (used to look up the auto-bloom)
- joinColumn
The column name in
valuesDfcontaining the query values- valuesDf
DataFrame containing the values to probe against bloom filters
- indexDf
The main index DataFrame containing auto-bloom binary columns
- returns
Some(set of candidate filenames) if an auto-bloom index exists and the value set is within bounds, None otherwise
- Attributes
- protected
- Definition Classes
- IndexQueryOperations
- Note
The bound applies to query-side cardinality only; it does not limit how many values an index may hold. Auto-bloom is only a pre-filter, so past the bound this returns
Noneand the caller proceeds without pruning. Results are unaffected — the large index is simply scanned without bloom pruning. The value set is never truncated, since dropping probe values would prune away the files holding them and turn a pre-filter into a source of missing rows.
-
final
def
getClass(): Class[_]
- Definition Classes
- AnyRef → Any
- Annotations
- @native() @HotSpotIntrinsicCandidate()
-
def
getFileSizes(files: Set[String]): Map[String, Long]
Computes file sizes in bytes for the given files using the Hadoop FileSystem.
Computes file sizes in bytes for the given files using the Hadoop FileSystem.
Files that cannot be found or read are silently skipped with a warning log.
- files
set of fully-qualified file paths to measure
- returns
map from file path to size in bytes; missing/unreadable files are omitted
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
handleLargeIndexes(sourceDf: DataFrame, large: LargeFileClassification): Unit
Writes the values of every large
(file, column)pair into its per-column Delta table.Writes the values of every large
(file, column)pair into its per-column Delta table.- sourceDf
the base DataFrame with
filenameand all source columns (computed indexes already applied)- large
the classification identifying which files are large for which columns
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
hasFile(fileName: String): Boolean
Checks if a file is tracked by this index's file list.
Checks if a file is tracked by this index's file list.
- fileName
The file path to check
- returns
true if the file is in the file list
val tracked = index.hasFile("s3a://bucket/data/file1.parquet")
Example: -
def
index: Option[DataFrame]
Helper function to load the index.
Helper function to load the index.
On first call, checks for and performs any needed exploded field column migration (pre-0.1.1 indexes stored columns under
array_columnnames).- returns
Some(DataFrame)containing the latest version of the index, orNoneif the index Delta table does not yet exist
- Attributes
- protected
- Definition Classes
- IndexQueryOperations
-
def
indexFilePath: Path
Hadoop path for the main index Delta table (
{storagePath}/index/).Hadoop path for the main index Delta table (
{storagePath}/index/).- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
lazy val
indexRepartitionCount: Option[Int]
Optional number of partitions to use when repartitioning the index DataFrame during joins.
Optional number of partitions to use when repartitioning the index DataFrame during joins. When set, the index DataFrame is repartitioned before expensive operations like explode to reduce per-executor memory pressure and avoid FetchFailedExceptions on large indexes. Reads from spark.ariadne.indexRepartitionCount configuration (default: not set).
- Definition Classes
- AriadneContextUser
-
def
indexes: Set[String]
Returns all column names that can be used in joins across all index types.
Returns all column names that can be used in joins across all index types.
Includes regular, computed, exploded field, bloom, temporal, and range index columns.
- returns
the union of all indexed column names across every index type
val allIndexedColumns: Set[String] = index.indexes
Example: -
final
def
isInstanceOf[T0]: Boolean
- Definition Classes
- Any
-
def
join(df: DataFrame, usingColumns: Seq[String], joinType: String = "inner"): DataFrame
Joins the indexed data with the provided DataFrame.
Joins the indexed data with the provided DataFrame.
Locates relevant data files via the index, reads them, applies temporal deduplication if configured, and joins the result with the provided DataFrame.
What this returns: the minimal complete dataset for the join. Every indexed value that matches
dfis found, and index-side rows absent fromdfare discarded. Because pruning selects whole files, unmatched index-side rows that share a file with a match are still read, so the result can include incidental rows and can change after Index.compact. Join types whose result is defined by those rows —left,fullandleft_anti, since the index is on the left here — are rejected rather than returning a file-layout-dependent answer.- df
The DataFrame to join against indexed data
- usingColumns
The column names to join on (must be indexed columns)
- joinType
The Spark join type. Supported here: "inner", "left_semi", "right", "right_outer" (default: "inner")
- returns
The joined DataFrame
- Definition Classes
- Index → IndexJoinOperations
val lookupDf = spark.read.parquet("/data/lookups") val result = index.join(lookupDf, Seq("userId")) // Keep every row of lookupDf, matched or not: val rightResult = index.join(lookupDf, Seq("userId"), "right_outer")
- Exceptions thrown
IllegalArgumentExceptionif df is null or usingColumns is null/emptydev.cjfravel.ariadne.exceptions.UnsupportedJoinTypeExceptionifjoinTypeis "left", "full" or "left_anti" (or an alias), which depend on unmatched index-side rows
Example: -
def
joinDf(df: DataFrame, usingColumns: Seq[String]): DataFrame
Locates and reads indexed data files relevant to the given DataFrame.
Locates and reads indexed data files relevant to the given DataFrame.
Uses the index to identify which files contain matching values, then reads those files into a lazy DataFrame. If no matching files are found, returns an empty DataFrame produced through the same read path as the populated branch, so computed indexes, exploded fields and any active
select()are applied and the schema matches either way. The actual row-level filtering happens in the subsequent join in join.Logs data pruning metrics (file count, data size saved) when available, and includes detailed file-level debug information when debug mode is enabled.
- df
The DataFrame to match against the index.
- usingColumns
The columns used for the join.
- returns
A lazy DataFrame containing data from indexed files, or an empty DataFrame if no files match.
- Attributes
- protected
- Definition Classes
- IndexJoinOperations
- Exceptions thrown
IllegalArgumentExceptionif none of the join columns have indexes, or if df/usingColumns are null/emptydev.cjfravel.ariadne.exceptions.ColumnNotFoundExceptionif join columns are not in the selected columns, schema, or indexes
-
def
largeIndexColumns: Set[String]
Returns the set of column names that have large index Delta tables.
Returns the set of column names that have large index Delta tables.
- returns
Set of column names with large index storage
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
lazy val
largeIndexLimit: Long
Maximum number of distinct values per file before an index column is promoted to a "large index" (stored in a separate exploded Delta table).
Maximum number of distinct values per file before an index column is promoted to a "large index" (stored in a separate exploded Delta table). Reads from
spark.ariadne.largeIndexLimit(default: 500000).If the configured value is not a valid positive number, a warning is logged and the default of 500,000 is used.
- Definition Classes
- AriadneContextUser
-
def
largeIndexesFilePath: Path
Hadoop root path for large index Delta tables (
{storagePath}/large_indexes/).Hadoop root path for large index Delta tables (
{storagePath}/large_indexes/).- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
loadColumnIndex(indexDf: DataFrame, colName: String, bloomCandidateFiles: Option[Set[String]] = None): DataFrame
Returns a unified (filename, colName) DataFrame for a regular index column by exploding the main index arrays and unioning with any large index rows that exist for that column.
Returns a unified (filename, colName) DataFrame for a regular index column by exploding the main index arrays and unioning with any large index rows that exist for that column.
- indexDf
The main index DataFrame (may be repartitioned)
- colName
The storage column name
- bloomCandidateFiles
Optional set of files that passed auto-bloom pre-filtering. When provided, only large index rows for these files are included, and an empty set drops the large index entirely. Both restrictions are skipped while a staging table exists; see pruneLargeIndexRows.
- returns
DataFrame with (filename, colName) scalar rows
- Attributes
- protected
- Definition Classes
- IndexQueryOperations
-
def
loadTemporalColumnIndex(indexDf: DataFrame, colName: String, bloomCandidateFiles: Option[Set[String]] = None): DataFrame
Returns a unified (filename, _value, _max_ts) DataFrame for a temporal index column by exploding the main index struct arrays and unioning with any large index rows that exist for that column.
Returns a unified (filename, _value, _max_ts) DataFrame for a temporal index column by exploding the main index struct arrays and unioning with any large index rows that exist for that column.
- indexDf
The main index DataFrame (may be repartitioned)
- colName
The temporal column name
- bloomCandidateFiles
Optional set of files that passed auto-bloom pre-filtering. When provided, only large index rows for these files are included, and an empty set drops the large index entirely. Both restrictions are skipped while a staging table exists; see pruneLargeIndexRows.
- returns
DataFrame with (filename, _value, _max_ts) rows
- Attributes
- protected
- Definition Classes
- IndexQueryOperations
-
def
locateFiles(indexes: Map[String, Array[Any]]): Set[String]
Locates files matching the given index values across all index types.
Locates files matching the given index values across all index types.
Queries regular, bloom, temporal, range, and auto-bloom indexes in parallel and intersects results with AND semantics. Each index type independently returns candidate files, and only files present in all queried categories are included in the final result.
- indexes
A map of index column names to arrays of values to search for. Columns are automatically routed to the appropriate index type (bloom, temporal, range, or regular).
- returns
A set of file paths matching all query criteria, or an empty set if no matches are found or no index exists
- Definition Classes
- IndexQueryOperations
val matchingFiles = index.locateFiles(Map("userId" -> Array("u1", "u2")))
- Exceptions thrown
IllegalArgumentExceptionifindexesis null- Note
If a staging table exists, its contents are collected to the driver for merging. This may cause driver OOM for very large staging tables.
Example: -
def
locateFilesFromDataFrame(valuesDf: DataFrame, columnMappings: Map[String, String], joinColumns: Seq[String]): Set[String]
Locates files based on a DataFrame containing join column values.
Locates files based on a DataFrame containing join column values. Handles both regular indexes and bloom filter indexes.
- valuesDf
DataFrame containing the distinct values to search for
- columnMappings
Map from join column names to storage column names in the index
- joinColumns
The columns to use for filtering
- returns
A set of file names matching the criteria.
- Definition Classes
- IndexQueryOperations
val columnMappings = Map("userId" -> "userId") val files = index.locateFilesFromDataFrame(lookupDf, columnMappings, Seq("userId"))
- Exceptions thrown
IllegalArgumentExceptionifvaluesDfis null orcolumnMappingsis null or empty
Example: -
def
locateFilesWithBloom(column: String, values: Array[Any], indexDf: DataFrame): Set[String]
Locates files that might contain any of the given values using bloom filters.
Locates files that might contain any of the given values using bloom filters.
The probe runs as a distributed Spark filter via probeBloomFilters: candidate values are broadcast to the executors and only matching filenames are returned to the driver. Bloom filter binary data is never collected.
- column
the source column name (without
bloom_prefix)- values
array of candidate values to search for
- indexDf
the index DataFrame containing bloom filter binary columns
- returns
set of filenames whose bloom filters indicate a possible match; empty set if the bloom column does not exist in the index
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
locateFilesWithBloomFromDataFrame(column: String, valuesDf: DataFrame, indexDf: DataFrame): Set[String]
Locates files that might contain values from a DataFrame using bloom filters.
Locates files that might contain values from a DataFrame using bloom filters.
Collects distinct non-null values from the specified column of
valuesDf, then delegates to locateFilesWithBloom for the actual bloom filter check.- column
the source column name (without
bloom_prefix)- valuesDf
DataFrame containing the candidate values to search for
- indexDf
the index DataFrame containing bloom filter binary columns
- returns
set of filenames whose bloom filters indicate a possible match; empty set if no non-null values exist or the bloom column is missing
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
- Note
This method collects the distinct candidate values from
valuesDfto the driver so they can be broadcast for the probe. Bloom filter binary data is not collected — see probeBloomFilters. Driver memory is therefore bounded by the join-key cardinality ofvaluesDf. This value set is never truncated: dropping probe values would prune away files that hold them and silently omit matching rows. When the probe bound is exceeded the pre-filter is skipped instead, returning every indexed file — a safe superset that the subsequent read and join narrow correctly.
-
lazy val
lockMaxWait: Long
Maximum total time in seconds to wait for lock acquisition before failing.
Maximum total time in seconds to wait for lock acquisition before failing. Reads from spark.ariadne.lockMaxWait configuration (default: 3600).
- Definition Classes
- AriadneContextUser
- Note
Unlike
lockTimeoutandlockRetryInterval, non-positive values are not validated — a value of 0 causes immediate timeout on first retry. Configure with care.
-
lazy val
lockRefreshInterval: Int
Refresh the update lock every N batches during updateBatched.
Refresh the update lock every N batches during updateBatched. Reads from spark.ariadne.lockRefreshInterval configuration (default: 1).
- Definition Classes
- AriadneContextUser
-
lazy val
lockRetryInterval: Long
Base interval in seconds between lock acquisition retries (exponential backoff applied).
Base interval in seconds between lock acquisition retries (exponential backoff applied). Reads from spark.ariadne.lockRetryInterval configuration (default: 60).
- Definition Classes
- AriadneContextUser
-
lazy val
lockTimeout: Long
Seconds since
lastRefreshedAtbefore a lock is considered stale and eligible for auto-heal (forcible acquisition by another process).Seconds since
lastRefreshedAtbefore a lock is considered stale and eligible for auto-heal (forcible acquisition by another process). Reads fromspark.ariadne.lockTimeout(default: 1800).- Definition Classes
- AriadneContextUser
-
lazy val
logger: Logger
Log4j logger shared by all Ariadne components.
Log4j logger shared by all Ariadne components. Uses the "ariadne" logger name. Marked
@transient lazyso thatIndex(a case class) remains serializable across Spark stages — Log4j loggers are notSerializable.- Definition Classes
- IndexMetadataOperations → AriadneContextUser
-
def
mapJoinColumnsToStorage(joinColumns: Seq[String]): Map[String, String]
Maps join column names to their corresponding storage column names.
Maps join column names to their corresponding storage column names.
Bloom filter columns are prefixed with the bloom prefix, range columns are prefixed with "range_", exploded field columns are mapped to their backing array column, and all other columns map to themselves.
- joinColumns
The column names used in joins
- returns
Map from each join column name to its storage column name in the index
- Attributes
- protected
- Definition Classes
- IndexJoinOperations
-
def
maybeAutoCompact(): Unit
Triggers compaction if the auto-compact threshold has been reached.
Triggers compaction if the auto-compact threshold has been reached.
Checks the batchesSinceCompact counter against the configured
autoCompactThreshold. If the threshold is met, compacts all Delta tables and resets the counter to zero. The counter is persisted to metadata so it survives across Spark jobs.- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
maybeRepartition(df: DataFrame): DataFrame
Conditionally repartitions a DataFrame if indexRepartitionCount is configured.
Conditionally repartitions a DataFrame if indexRepartitionCount is configured. This helps avoid FetchFailedExceptions when working with very large index DataFrames by spreading data across more partitions before expensive operations like explode.
- df
The DataFrame to potentially repartition
- returns
Repartitioned DataFrame if configured, otherwise the original DataFrame
- Attributes
- protected
- Definition Classes
- IndexQueryOperations
-
def
metadataExists: Boolean
Checks if metadata exists in the storage location.
Checks if metadata exists in the storage location.
- returns
True if metadata exists, otherwise false.
- Attributes
- protected
- Definition Classes
- IndexMetadataOperations
-
def
metadataFilePath: Path
Hadoop path for the metadata file.
Hadoop path for the metadata file.
- returns
the path to
metadata.jsonunder the index storage directory
- Attributes
- protected
- Definition Classes
- IndexMetadataOperations
-
def
migrateAutoBloomFilters(checkHeartbeat: () ⇒ Unit): Unit
Builds auto-bloom filters for columns that already have a
large_indexes/table but no filter in the main index.Builds auto-bloom filters for columns that already have a
large_indexes/table but no filter in the main index.Filters are folded from autoBloomValueRows, covering files whose values are inline as well as those in the large index. Files that hold at least one non-null value for the column end up with a filter. Files that hold none keep a null filter, which the query side treats as an unconditional candidate (
includeNullFilters = true), so every file is still represented and a query can treat an empty candidate set as a definitive no-match.Metadata is written only after every column has been backfilled, so a failure part way through leaves the storage version unchanged and the next preflight repeats the work.
- checkHeartbeat
callback that refreshes the update lock and aborts if ownership was lost
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
migrateExplodedFieldColumns(): Unit
Migrates pre-0.1.1 indexes that stored exploded field columns under
array_columnnames to the currentas_columnnaming convention.Migrates pre-0.1.1 indexes that stored exploded field columns under
array_columnnames to the currentas_columnnaming convention.Detection: reads the main and staging Delta table schemas and checks whether any
ExplodedFieldMapping.array_columnappears as a column while the correspondingas_columndoes not. When this pattern is found, each table is rewritten with renamed columns, and any large-index directories named byarray_columnare similarly migrated.This method is idempotent and is invoked only by storage migration preflight while the update lock is held.
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
migrateLargeIndexColumnNames(checkHeartbeat: () ⇒ Unit): Unit
Rewrites large index tables whose value column still uses a legacy name so they store it under the canonical name.
Rewrites large index tables whose value column still uses a legacy name so they store it under the canonical name.
loadLargeIndex resolves the legacy shape on read, but appendToLargeIndex writes under the canonical name with
mergeSchema, which would widen a legacy table into two value columns. Rewriting once removes that divergence.- checkHeartbeat
callback that refreshes the update lock and aborts if ownership was lost
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- val name: String
-
final
def
ne(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
-
final
def
notify(): Unit
- Definition Classes
- AnyRef
- Annotations
- @native() @HotSpotIntrinsicCandidate()
-
final
def
notifyAll(): Unit
- Definition Classes
- AnyRef
- Annotations
- @native() @HotSpotIntrinsicCandidate()
-
def
open(path: Path): FSDataInputStream
Opens an input stream to read from a path on the filesystem.
Opens an input stream to read from a path on the filesystem.
- path
The Hadoop Path to open for reading
- returns
An FSDataInputStream for reading the file contents
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulljava.io.IOExceptionif the file does not exist or cannot be opened
-
def
probeBloomFilters(indexDf: DataFrame, bloomColumn: String, values: Array[Any], includeNullFilters: Boolean): Set[String]
Probes every file's bloom filter for
bloomColumnagainstvalues, entirely on the executors.Probes every file's bloom filter for
bloomColumnagainstvalues, entirely on the executors.This is the memory-safe counterpart to collecting bloom filter binary data to the driver. The candidate values are broadcast, the
mightContaintest runs as a distributed filter, and only the surviving filenames are returned to the driver. Driver memory is therefore bounded by the number of matching files rather than by the total serialized size of the bloom filters.This distinction matters most for auto-bloom columns. Once any one file crosses
largeIndexLimitdistinct values (500,000 by default), every file in that column gets a filter, each sized from that file's own distinct count; the largest can reach roughly 585 KiB at the default 1% FPR. Collecting one per file would scale driver memory with the number of indexed files.- indexDf
the index DataFrame containing
filenameand the bloom binary column- bloomColumn
the fully prefixed storage column name (e.g.
bloom_user_idorauto_bloom_user_id)- values
candidate values to probe; nulls are ignored and the remainder de-duplicated before broadcast
- includeNullFilters
whether rows with a
nullbloom filter should be treated as candidates- returns
set of filenames whose bloom filter indicates a possible match
- Attributes
- protected
- Definition Classes
- BloomFilterOperations
-
def
rangeStorageColumns: Set[String]
Returns the set of range index storage column names (each prefixed with
range_).Returns the set of range index storage column names (each prefixed with
range_).- returns
set of prefixed column names (e.g.,
range_event_date)
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
refreshMetadata(): Unit
Forces a reload of metadata from the
metadata.jsonfile on disk.Forces a reload of metadata from the
metadata.jsonfile on disk.Replaces the in-memory cached metadata with a freshly deserialized copy. This is called automatically by metadata on first access; callers may invoke it explicitly to pick up external changes.
- Definition Classes
- IndexMetadataOperations
- Exceptions thrown
MetadataMissingOrCorruptExceptionif the metadata file does not exist, cannot be read, or contains invalid JSON
-
lazy val
repartitionDataFiles: Boolean
When true, applies indexRepartitionCount repartitioning to the data files read during joinDf.
When true, applies indexRepartitionCount repartitioning to the data files read during joinDf. When false, data files keep their natural parquet partitioning. Disable when column selection significantly reduces data volume, making the repartition shuffle more expensive than useful. Reads from spark.ariadne.repartitionDataFiles configuration (default: false).
- Definition Classes
- AriadneContextUser
-
def
requireBloomCompatibleType(column: String): Unit
Rejects a column whose values cannot be canonicalized into a stable bloom filter token.
Rejects a column whose values cannot be canonicalized into a stable bloom filter token.
Bloom filters hash the
toStringform of each value. Spark materializesBinaryTypeas a JVMArray[Byte], whosetoStringis identity-based ([B@1b6d3586), so the same bytes produce a different token on every occurrence — even for two references to the same object read at different times. A filter built over such a column is internally consistent but matches nothing, so every query prunes away every file. Nested binary is rejected for the same reason.The check is skipped when no schema is stored, matching the surrounding column validations: the type cannot be established without one.
- column
the candidate bloom index column name
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- Exceptions thrown
IllegalArgumentExceptionif the column's type is, or contains, a binary type
-
def
requireColumnNotAlreadyIndexed(column: String, newType: String): Unit
Enforces that a column carries at most one index type.
Enforces that a column carries at most one index type.
Checking every type in one place keeps the matrix symmetric by construction: a column carrying two index types can produce wrong results at query time rather than an error, and per-method lists of checks make detection depend on registration order wherever an entry is missing. A new index type therefore cannot reintroduce a one-directional gap.
Callers must invoke this only after confirming the column is not already registered as
newType, so that re-adding the same index type remains idempotent.- column
the column being indexed
- newType
the label of the index type being added, used in the error message
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- Exceptions thrown
IllegalArgumentExceptionif the column already carries a different index type
-
def
requireNonReservedStagingColumn(column: String): Unit
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
requireTopLevelIndexColumn(column: String, indexType: String): Unit
Rejects nested (dotted) paths for columns that become index column names.
Rejects nested (dotted) paths for columns that become index column names.
An indexed value column is persisted as a column of the index table under its own name, and is later read back with
col(name). A dotted path such asmeta.userIdwould be written as a literal column name but read back as nested field access, so such an index can never be built or queried.SchemaHelper.fieldExistsresolves dotted paths, so without this guard the configuration is accepted and persisted tometadata.json, and only fails later duringupdatewith an opaque Spark analysis error.This restriction applies to the indexed value column only. Temporal timestamp columns may be nested, because they are never persisted under their own name.
- column
the candidate index column name
- indexType
the index type, used to build the error message
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- Exceptions thrown
IllegalArgumentExceptionif the column is a nested path
-
def
safeDestroyBroadcast(broadcast: Broadcast[_]): Unit
Safely destroys a broadcast variable, logging but swallowing any exception.
Safely destroys a broadcast variable, logging but swallowing any exception.
Intended for use inside
finallyblocks where a cleanup failure must not mask the original exception propagating out of thetryblock. Tolerates a null broadcast.- broadcast
the broadcast to destroy; may be null
- Attributes
- protected
- Definition Classes
- AriadneContextUser
- val schema: Option[StructType]
-
def
select(columns: String*): Index
Selects specific columns for optimized reading.
Selects specific columns for optimized reading.
When set, only the selected columns (plus any join columns) are read from data files, reducing I/O. Returns this Index for method chaining.
- columns
The column names to select
- returns
This Index instance for method chaining
val result = index.select("name", "email").join(df, Seq("userId"))
- Exceptions thrown
ColumnNotFoundExceptionif any specified column doesn't exist in the schemaIllegalArgumentExceptionif columns is null/empty or any column is null/blank
Example: -
implicit
val
spark: SparkSession
Implicit SparkSession that must be provided by the mixing class.
Implicit SparkSession that must be provided by the mixing class.
- Definition Classes
- Index → AriadneContextUser
-
lazy val
stagingConsolidationThreshold: Int
Number of batches to process before consolidating staged data into the main index.
Number of batches to process before consolidating staged data into the main index. This provides fault tolerance for large index builds - if a job fails, work is preserved up to the last consolidation point. Reads from spark.ariadne.stagingConsolidationThreshold configuration (default: 50).
- Definition Classes
- AriadneContextUser
-
def
stagingFilePath: Path
Hadoop path for the staging Delta table (
{storagePath}/staging/).Hadoop path for the staging Delta table (
{storagePath}/staging/).The staging table accumulates batch results during update and is merged into the main index by consolidateStaging.
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
def
stats(): DataFrame
Returns a DataFrame of per-column index statistics and total file count.
Returns a DataFrame of per-column index statistics and total file count.
For each indexed column, computes statistics on the array length (number of distinct values per file): min, max, avg, median, and standard deviation. Also includes the total number of indexed files.
Returns an empty DataFrame if no index table exists. If
storageColumnsis empty (e.g., only bloom or range indexes), the result contains zero stat rows.- returns
Single-row DataFrame with FileCount and per-column stat structs, or an empty DataFrame if no index exists
- Definition Classes
- IndexQueryOperations
// Compute and display index statistics val statsDF = index.stats() statsDF.show() // +----------+---------+---------+---------+---------+------------+------+ // | Column|FileCount|MinValues|MaxValues|AvgValues|MedianValues|StdDev| // +----------+---------+---------+---------+---------+------------+------+ // | user_id | 150 | 1 | 5000 | 120.3 | 45 | 340.2| // +----------+---------+---------+---------+---------+------------+------+
Example: -
def
storageColumns: Set[String]
Returns the set of all storage column names across regular, computed, exploded-field, and temporal index types.
Returns the set of all storage column names across regular, computed, exploded-field, and temporal index types.
This is used internally to determine which columns to check for large-index separation and to build aggregation expressions.
- returns
set of column names used for index storage (excludes bloom, range, and auto-bloom)
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
-
lazy val
storagePath: Path
Path to the storage location of the index.
Path to the storage location of the index.
- Definition Classes
- Index → AriadneContextUser
-
def
storedSchema: StructType
Returns the stored schema of the index.
Returns the stored schema of the index.
Parses the JSON schema string from metadata into a Spark StructType.
- returns
The StructType schema parsed from metadata
- Definition Classes
- IndexFileOperations
val schema = index.storedSchema schema.printTreeString()- Exceptions thrown
SchemaParseExceptionif the schema string cannot be parsed as a StructTypedev.cjfravel.ariadne.exceptions.MissingSchemaExceptionif the schema is null in metadata
Example: -
final
def
synchronized[T0](arg0: ⇒ T0): T0
- Definition Classes
- AnyRef
-
def
update: Unit
Updates the index with new files and backfills newly added columns.
Updates the index with new files and backfills newly added columns.
Processes all unindexed files registered via addFile, using intelligent batching based on pre-flight analysis. Also backfills existing files when new index columns have been added since the last update.
Consistency
A column's values for a given file live in exactly one place: inline in the main index table, or — once the file contributes at least
largeIndexLimitdistinct values — inlarge_indexes/{column}/. Queries read both and union the results, so this exclusivity is what keeps a file from being counted twice or missed entirely.That invariant is guaranteed for an update that completes successfully. The two stores are separate Delta tables written in separate commits, so an update killed part-way can leave them disagreeing about a file. In the usual case the file is then still absent from the main index, so it remains unindexed and the next
updatereprocesses it and restores consistency by itself. The exception is a column backfill over a file that was already indexed: its stale large-index rows are deleted before the replacements are written, and a crash in that window leaves the file present in the main index — hence not reprocessed — with its large values gone. Re-running the update does not detect this; rebuilding the index does.This is a known and accepted limitation, recorded here rather than worked around. Callers that cannot tolerate it should treat an interrupted
updateas a signal to rebuild rather than resume.val index = Index("myIndex", schema, "parquet") index.addIndex("userId") index.addFile("/data/events/2024-01-01.parquet") index.update
- Exceptions thrown
IndexLockExceptionif the update lock cannot be acquired within the configured timeout
Example: -
def
updateLockPath: Path
Lock path for index update and storage migration operations.
Lock path for index update and storage migration operations.
- Attributes
- protected
-
def
vacuum(retentionHours: Int = 168): Unit
Vacuums all Delta tables belonging to this index to remove old files.
Vacuums all Delta tables belonging to this index to remove old files. Acquires the update lock to prevent concurrent modifications.
A
retentionHoursof zero or less additionally disables Delta's retention safety check (spark.databricks.delta.retentionDurationCheck.enabled) for the duration of the call, because Delta otherwise refuses any retention below 168 hours. That check exists to stop a vacuum from deleting files a concurrent reader or writer is still using, sovacuum(0)is only safe when nothing else is touching the index — the update lock excludes other Ariadne writers, but not readers and not non-Ariadne access to the same paths.- retentionHours
number of hours of history to retain (default 168 = 7 days); zero or less bypasses Delta's retention check
index.vacuum() // default 7 days retention index.vacuum(24) // 1 day retention
- Exceptions thrown
IndexLockExceptionif the update lock cannot be acquired within the configured timeout- Note
Thread-safety: The retention check is a
SparkSession-wide setting. It is restored in afinallyblock, but while avacuum(0)is in flight any other Delta operation on the sameSparkSessionalso runs without the guardrail.
Example: -
def
vacuumDeltaTables(retentionHours: Int = 168): Unit
Vacuums all Delta tables (main index and large index tables) to remove old files.
Vacuums all Delta tables (main index and large index tables) to remove old files.
Uses Delta Lake's VACUUM command to delete data files no longer referenced by the transaction log. Temporarily disables the retention duration safety check when
retentionHoursis zero or negative.- retentionHours
number of hours of history to retain (default 168 = 7 days)
- Attributes
- protected
- Definition Classes
- IndexBuildOperations
- Note
Thread-safety: This method temporarily mutates the shared SparkConf (
retentionDurationCheck.enabled). ConcurrentIndexinstances sharing the sameSparkSessionmay race on this setting. The value is restored in afinallyblock, but a TOCTOU window exists between set and restore.
-
final
def
wait(arg0: Long, arg1: Int): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws( ... )
-
final
def
wait(arg0: Long): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws( ... ) @native()
-
final
def
wait(): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws( ... )
-
def
writeMetadata(metadata: IndexMetadata): Unit
Writes metadata to the
metadata.jsonfile and updates the in-memory cache.Writes metadata to the
metadata.jsonfile and updates the in-memory cache.Creates the parent directory if it does not exist. JSON is written and validated at a unique sibling path, then replaced through Hadoop's overwrite rename operation. The replacement is atomic on filesystems with atomic rename semantics (such as HDFS); object-store connectors may implement rename as copy/delete and cannot provide the same guarantee.
- metadata
the metadata instance to serialize and persist
- Attributes
- protected
- Definition Classes
- IndexMetadataOperations
- Exceptions thrown
java.io.IOExceptionif the write fails (the error is logged before propagating)
Deprecated Value Members
-
def
finalize(): Unit
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws( classOf[java.lang.Throwable] ) @Deprecated
- Deprecated