trait IndexBuildOperations extends BloomFilterOperations
Trait providing index building and maintenance operations for Index instances.
This trait implements the core data pipeline for building and maintaining file-level indexes. It is mixed into Index via the trait hierarchy and handles the complete update lifecycle:
Update pipeline (data flow):
- analyzeFiles — Pre-flight scan that computes per-file distinct value counts for all indexed columns. This information drives batching decisions.
- createOptimalBatches — Groups files into batches so that the sum of
maxDistinctCountper batch stays at or belowlargeIndexLimit. Files that individually exceed the limit are isolated into single-file batches. - Per-batch processing (in
updateSingleBatch):- Read source files, apply computed indexes, add filename column
- classifyLargeFiles — Decide, per
(file, column)pair, whether the file contributes enough distinct values to be stored out of line, reusing the counts from step 1 - Build regular indexes (array aggregation per file, skipping large pairs)
- Build exploded field, bloom filter, temporal, and range indexes
- Build auto-bloom filters for columns with at least one large file
- appendToStaging (inline columns to staging Delta table)
- appendToLargeIndex (large columns streamed to per-column Delta tables)
- Periodic consolidation — Every
stagingConsolidationThresholdbatches, consolidateStaging merges staging into the main index via Delta MERGE (upsert on filename). Auto-compaction may follow via maybeAutoCompact. - Final consolidation — After all batches complete, any remaining staged data is consolidated and the staging table is deleted.
Batching strategy: Files are sorted by maxDistinctCount (largest first) and packed sequentially into batches. This
greedy approach keeps each batch at or below the large index limit while maximizing batch sizes for efficiency.
Staging: Each batch is appended to a transient staging Delta table. Periodic consolidation merges staged rows into the main index, providing fault tolerance (partially processed updates can be recovered).
Large index handling: Columns whose array size exceeds largeIndexLimit for any file are stored in separate Delta
tables under large_indexes/{column}/ as exploded (one row per value) rather than array-per-file. Auto-bloom filters
are built for these columns in the main index to enable pre-filtering at query time.
- Self Type
- Index
- See also
BloomFilterOperations for bloom filter creation and querying
Index.update for the public entry point
- Alphabetic
- By Inheritance
- 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
-
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
Abstract Value Members
-
implicit abstract
def
spark: SparkSession
Implicit SparkSession that must be provided by the mixing class.
Implicit SparkSession that must be provided by the mixing class.
- Definition Classes
- AriadneContextUser
Concrete 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
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
analyzeFiles(files: Set[String]): Seq[Index.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
- 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: Index.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
- 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
-
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
-
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
-
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
-
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
-
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: Index.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
-
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: Index.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
-
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
-
def
buildRegularIndexes(df: DataFrame, large: Index.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
-
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: Index.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
-
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[Index.FileAnalysis]): Index.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
-
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
-
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
-
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
-
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[Index.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
-
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
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
-
final
def
eq(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
-
def
equals(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
-
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
-
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
-
def
handleLargeIndexes(sourceDf: DataFrame, large: Index.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
-
def
hashCode(): Int
- Definition Classes
- AnyRef → Any
- Annotations
- @native() @HotSpotIntrinsicCandidate()
-
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
-
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
-
final
def
isInstanceOf[T0]: Boolean
- Definition Classes
- Any
-
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
-
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
-
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
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
-
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
-
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
-
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
-
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
-
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
- 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
- Exceptions thrown
IllegalArgumentExceptionif the column already carries a different index type
-
def
requireNonReservedStagingColumn(column: String): Unit
- Attributes
- protected
-
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
- 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
-
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
-
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
-
lazy val
storagePath: Path
Base path on the Hadoop-compatible filesystem where all Ariadne index data is stored.
Base path on the Hadoop-compatible filesystem where all Ariadne index data is stored. Each index creates a subdirectory under this path.
Reads from
spark.ariadne.storagePath(required — no default).- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif the configuration key is not set
-
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
toString(): String
- Definition Classes
- AnyRef → Any
-
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
- 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