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.

Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. Index
  2. Serializable
  3. Serializable
  4. Product
  5. Equals
  6. IndexQueryOperations
  7. IndexJoinOperations
  8. IndexBuildOperations
  9. BloomFilterOperations
  10. IndexFileOperations
  11. IndexMetadataOperations
  12. AriadneContextUser
  13. AnyRef
  14. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. All

Type Members

  1. 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
  2. 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 least largeIndexLimit distinct values to that column. Those values are streamed into large_indexes/{column} as individual rows instead of being collected into a per-file array, and the inline array column is left null.

    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

  1. final def !=(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  2. final def ##(): Int
    Definition Classes
    AnyRef → Any
  3. final def ==(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  4. 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 BinaryType as a JVM Array[Byte] whose toString is 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%)

    Example:
    1. index.addBloomIndex("sessionId")
      index.addBloomIndex("ipAddress", 0.001)
    Exceptions thrown

    ColumnNotFoundException if column doesn't exist in schema

    IllegalArgumentException if 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

  5. 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

    Example:
    1. index.addComputedIndex("yearMonth", "date_format(event_date, 'yyyy-MM')")
    Exceptions thrown

    IllegalArgumentException if name or sql_expression is null/blank, or name conflicts with another index type

  6. 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 arrayColumn is exploded, the fieldPath is extracted, and the distinct values are stored under asColumn in the index. Joins on asColumn will locate files containing any matching array element. Idempotent: calling again with the same asColumn is 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.

    Example:
    1. // Index the "id" field from each element of the "items" array column
      index.addExplodedFieldIndex("items", "id", "item_id")
    Exceptions thrown

    IllegalArgumentException if any parameter is null/blank, or asColumn conflicts with another index type

  7. 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

    Example:
    1. index.addFile("/data/events/2024-01-01.parquet")
      index.addFile("/data/events/2024-01-02.parquet", "/data/events/2024-01-03.parquet")
    Exceptions thrown

    IllegalArgumentException if fileNames is null/empty or any fileName is null/blank

  8. 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 filename column added

    Attributes
    protected
    Definition Classes
    IndexFileOperations
  9. 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.

    Example:
    1. val index = Index("myIndex", schema, "parquet")
      index.addIndex("userId")
    Exceptions thrown

    ColumnNotFoundException if the column doesn't exist in the schema

    IllegalArgumentException if the column name is null/blank, is a nested path, or is already indexed by any other type

  10. 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

    Example:
    1. index.addRangeIndex("timestamp")
    Exceptions thrown

    ColumnNotFoundException if column doesn't exist in schema

    IllegalArgumentException if column is null/blank or already indexed by another type

  11. 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

    Example:
    1. index.addTemporalIndex("userId", "updated_at")
      // the timestamp may also be a nested path
      index.addTemporalIndex("orderId", "meta.updatedAt")
    Exceptions thrown

    ColumnNotFoundException if either column doesn't exist in schema

    IllegalArgumentException if column or timestampColumn is null/blank, if column is a nested path, or if column is already indexed by another type

  12. 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: applyExplodedFields uses an inner explode, 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 the collect_set of struct(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 files is 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.

  13. 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 to large_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 filename and 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.

  14. 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 already null at this point: the index builders skip them and their values were streamed to large_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
  15. 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 df are ignored.

    returns

    DataFrame with selected columns or original DataFrame if no selection

    Attributes
    protected
    Definition Classes
    IndexFileOperations
  16. 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.AnalysisException if a computed index expression is invalid or references non-existent columns

  17. 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() (not explode_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
  18. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  19. 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
  20. 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
  21. 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
  22. 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 the value field 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
  23. 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
  24. 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 updateBatched and reset to zero after each compaction cycle or at the start of each update call to prevent stale counts from a previous update from triggering premature compaction.

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  25. 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
  26. 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
  27. 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
  28. 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
  29. 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 largeIndexLimit distinct values to it. From then on every file gets a filter, stored in the main index under the auto_bloom_ prefix, so queries can cheaply skip files before touching large_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 filename and all source columns (computed indexes already applied)

    large

    the classification identifying which files are large for which columns

    returns

    combinedDf with an auto_bloom_{column} binary column per auto-bloom column; unchanged if none qualify

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  30. 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:

    1. Reduces the input to distinct (filename, value) pairs
    2. Computes each file's distinct value count so its filter can be sized individually
    3. Folds the values into a bloom filter with BloomFilterAggregator
    4. 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 filename column and all bloom-indexed source columns

    returns

    DataFrame with filename plus one bloom_{column} binary column per configured bloom index

    Attributes
    protected
    Definition Classes
    BloomFilterOperations
  31. 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
  32. 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 filename and computes the min and max of the column per file. The result is a struct column named range_{column}.

    df

    the base DataFrame with filename column and range-indexed source columns

    returns

    DataFrame with filename plus one range struct column per configured range index; filename-only DataFrame if no range indexes are configured

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  33. 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 filename column and all indexed source columns

    returns

    DataFrame with filename plus one array column per regular/computed index

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  34. 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 filename column 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 filename plus the serialized filter column; files contributing no non-null values are absent

    Attributes
    protected
    Definition Classes
    BloomFilterOperations
  35. 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 filename plus one struct-array column per temporal index; filename-only DataFrame if no temporal indexes are configured

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  36. 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 applies toString inside 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 toString is not value-based. Array[Byte].toString is 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 as 2024-01-01 00:00:00 while java.sql.Timestamp.toString renders 2024-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, null where the source is null

    Attributes
    protected
    Definition Classes
    BloomFilterOperations
  37. 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
  38. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native() @HotSpotIntrinsicCandidate()
  39. 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 any distinct or 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 filename column, computed indexes applied, and all source columns present

    column

    the storage column name

    returns

    Some DataFrame of filename plus column, or None when column is not a known storage column

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  40. 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.

    Example:
    1. index.compact()
    Exceptions thrown

    IndexLockException if the update lock cannot be acquired within the configured timeout

  41. 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
  42. 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
  43. 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

    IllegalArgumentException if the stored format is not csv, parquet, or json

  44. def createOptimalBatches(fileAnalyses: Seq[FileAnalysis]): Seq[Set[String]]

    Groups files into batches whose aggregate maxDistinctCount stays at or below largeIndexLimit.

    Groups files into batches whose aggregate maxDistinctCount stays at or below largeIndexLimit.

    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 if fileAnalyses is empty

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  45. 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
  46. 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

    IllegalArgumentException if path is null

    java.io.IOException if the filesystem operation fails

  47. 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.

    Example:
    1. index.deleteFiles("/data/events/2024-01-01.parquet")
      index.deleteFiles("/data/old1.parquet", "/data/old2.parquet")
    Exceptions thrown

    IllegalArgumentException if filenames is null/empty or contains null/blank entries

    IndexLockException if the update lock cannot be acquired within the configured timeout

  48. 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, None if the path is absent or an empty directory

    Definition Classes
    AriadneContextUser
    Exceptions thrown

    IllegalArgumentException if path is null

    dev.cjfravel.ariadne.exceptions.InvalidDeltaTableException if the path exists and is non-empty but is not a readable Delta table

  49. 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
  50. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  51. 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

    IllegalArgumentException if path is null

    java.io.IOException if the filesystem operation fails

  52. 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
  53. lazy val fs: FileSystem

    Hadoop FileSystem instance resolved from the storagePath URI and the Spark Hadoop configuration.

    Hadoop FileSystem instance resolved from the storagePath URI and the Spark Hadoop configuration. Used for all filesystem operations (existence checks, reads, deletes, lock files).

    Definition Classes
    AriadneContextUser
  54. 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 null filter 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 valuesDf containing 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 None and 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.

  55. final def getClass(): Class[_]
    Definition Classes
    AnyRef → Any
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  56. 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
  57. 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 filename and all source columns (computed indexes already applied)

    large

    the classification identifying which files are large for which columns

    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  58. 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

    Example:
    1. val tracked = index.hasFile("s3a://bucket/data/file1.parquet")
  59. 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_column names).

    returns

    Some(DataFrame) containing the latest version of the index, or None if the index Delta table does not yet exist

    Attributes
    protected
    Definition Classes
    IndexQueryOperations
  60. 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
  61. 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
  62. 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

    Example:
    1. val allIndexedColumns: Set[String] = index.indexes
  63. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  64. 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 df is found, and index-side rows absent from df are 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, full and left_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
    Example:
    1. 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

    IllegalArgumentException if df is null or usingColumns is null/empty

    dev.cjfravel.ariadne.exceptions.UnsupportedJoinTypeException if joinType is "left", "full" or "left_anti" (or an alias), which depend on unmatched index-side rows

  65. 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

    IllegalArgumentException if none of the join columns have indexes, or if df/usingColumns are null/empty

    dev.cjfravel.ariadne.exceptions.ColumnNotFoundException if join columns are not in the selected columns, schema, or indexes

  66. 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
  67. 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
  68. 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
  69. 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
  70. 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
  71. 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
    Example:
    1. val matchingFiles = index.locateFiles(Map("userId" -> Array("u1", "u2")))
    Exceptions thrown

    IllegalArgumentException if indexes is 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.

  72. 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
    Example:
    1. val columnMappings = Map("userId" -> "userId")
      val files = index.locateFilesFromDataFrame(lookupDf, columnMappings, Seq("userId"))
    Exceptions thrown

    IllegalArgumentException if valuesDf is null or columnMappings is null or empty

  73. 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
  74. 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 valuesDf to 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 of valuesDf. 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.

  75. 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 lockTimeout and lockRetryInterval, non-positive values are not validated — a value of 0 causes immediate timeout on first retry. Configure with care.

  76. 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
  77. 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
  78. lazy val lockTimeout: Long

    Seconds since lastRefreshedAt before a lock is considered stale and eligible for auto-heal (forcible acquisition by another process).

    Seconds since lastRefreshedAt before a lock is considered stale and eligible for auto-heal (forcible acquisition by another process). Reads from spark.ariadne.lockTimeout (default: 1800).

    Definition Classes
    AriadneContextUser
  79. 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 lazy so that Index (a case class) remains serializable across Spark stages — Log4j loggers are not Serializable.

    Definition Classes
    IndexMetadataOperations → AriadneContextUser
  80. 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
  81. 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
  82. 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
  83. 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
  84. def metadataFilePath: Path

    Hadoop path for the metadata file.

    Hadoop path for the metadata file.

    returns

    the path to metadata.json under the index storage directory

    Attributes
    protected
    Definition Classes
    IndexMetadataOperations
  85. 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
  86. def migrateExplodedFieldColumns(): Unit

    Migrates pre-0.1.1 indexes that stored exploded field columns under array_column names to the current as_column naming convention.

    Migrates pre-0.1.1 indexes that stored exploded field columns under array_column names to the current as_column naming convention.

    Detection: reads the main and staging Delta table schemas and checks whether any ExplodedFieldMapping.array_column appears as a column while the corresponding as_column does not. When this pattern is found, each table is rewritten with renamed columns, and any large-index directories named by array_column are 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
  87. 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
  88. val name: String
  89. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  90. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  91. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  92. 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

    IllegalArgumentException if path is null

    java.io.IOException if the file does not exist or cannot be opened

  93. def probeBloomFilters(indexDf: DataFrame, bloomColumn: String, values: Array[Any], includeNullFilters: Boolean): Set[String]

    Probes every file's bloom filter for bloomColumn against values, entirely on the executors.

    Probes every file's bloom filter for bloomColumn against values, 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 mightContain test 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 largeIndexLimit distinct 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 filename and the bloom binary column

    bloomColumn

    the fully prefixed storage column name (e.g. bloom_user_id or auto_bloom_user_id)

    values

    candidate values to probe; nulls are ignored and the remainder de-duplicated before broadcast

    includeNullFilters

    whether rows with a null bloom filter should be treated as candidates

    returns

    set of filenames whose bloom filter indicates a possible match

    Attributes
    protected
    Definition Classes
    BloomFilterOperations
  94. 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
  95. def refreshMetadata(): Unit

    Forces a reload of metadata from the metadata.json file on disk.

    Forces a reload of metadata from the metadata.json file 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

    MetadataMissingOrCorruptException if the metadata file does not exist, cannot be read, or contains invalid JSON

  96. 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
  97. 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 toString form of each value. Spark materializes BinaryType as a JVM Array[Byte], whose toString is 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

    IllegalArgumentException if the column's type is, or contains, a binary type

  98. 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

    IllegalArgumentException if the column already carries a different index type

  99. def requireNonReservedStagingColumn(column: String): Unit
    Attributes
    protected
    Definition Classes
    IndexBuildOperations
  100. 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 as meta.userId would 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.fieldExists resolves dotted paths, so without this guard the configuration is accepted and persisted to metadata.json, and only fails later during update with 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

    IllegalArgumentException if the column is a nested path

  101. 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 finally blocks where a cleanup failure must not mask the original exception propagating out of the try block. Tolerates a null broadcast.

    broadcast

    the broadcast to destroy; may be null

    Attributes
    protected
    Definition Classes
    AriadneContextUser
  102. val schema: Option[StructType]
  103. 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

    Example:
    1. val result = index.select("name", "email").join(df, Seq("userId"))
    Exceptions thrown

    ColumnNotFoundException if any specified column doesn't exist in the schema

    IllegalArgumentException if columns is null/empty or any column is null/blank

  104. 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
  105. 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
  106. 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
  107. 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 storageColumns is 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
    Example:
    1. // 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|
      // +----------+---------+---------+---------+---------+------------+------+
  108. 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
  109. lazy val storagePath: Path

    Path to the storage location of the index.

    Path to the storage location of the index.

    Definition Classes
    Index → AriadneContextUser
  110. 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
    Example:
    1. val schema = index.storedSchema
      schema.printTreeString()
    Exceptions thrown

    SchemaParseException if the schema string cannot be parsed as a StructType

    dev.cjfravel.ariadne.exceptions.MissingSchemaException if the schema is null in metadata

  111. final def synchronized[T0](arg0: ⇒ T0): T0
    Definition Classes
    AnyRef
  112. 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 largeIndexLimit distinct values — in large_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 update reprocesses 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 update as a signal to rebuild rather than resume.

    Example:
    1. val index = Index("myIndex", schema, "parquet")
      index.addIndex("userId")
      index.addFile("/data/events/2024-01-01.parquet")
      index.update
    Exceptions thrown

    IndexLockException if the update lock cannot be acquired within the configured timeout

  113. def updateLockPath: Path

    Lock path for index update and storage migration operations.

    Lock path for index update and storage migration operations.

    Attributes
    protected
  114. 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 retentionHours of 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, so vacuum(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

    Example:
    1. index.vacuum()          // default 7 days retention
      index.vacuum(24)        // 1 day retention
    Exceptions thrown

    IndexLockException if 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 a finally block, but while a vacuum(0) is in flight any other Delta operation on the same SparkSession also runs without the guardrail.

  115. 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 retentionHours is 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). Concurrent Index instances sharing the same SparkSession may race on this setting. The value is restored in a finally block, but a TOCTOU window exists between set and restore.

  116. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  117. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native()
  118. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  119. def writeMetadata(metadata: IndexMetadata): Unit

    Writes metadata to the metadata.json file and updates the in-memory cache.

    Writes metadata to the metadata.json file 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.IOException if the write fails (the error is logged before propagating)

Deprecated Value Members

  1. def finalize(): Unit
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( classOf[java.lang.Throwable] ) @Deprecated
    Deprecated

Inherited from Serializable

Inherited from Serializable

Inherited from Product

Inherited from Equals

Inherited from IndexQueryOperations

Inherited from IndexJoinOperations

Inherited from IndexBuildOperations

Inherited from BloomFilterOperations

Inherited from IndexFileOperations

Inherited from IndexMetadataOperations

Inherited from AriadneContextUser

Inherited from AnyRef

Inherited from Any

Ungrouped