t

dev.cjfravel.ariadne

IndexBuildOperations

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):

  1. analyzeFiles — Pre-flight scan that computes per-file distinct value counts for all indexed columns. This information drives batching decisions.
  2. createOptimalBatches — Groups files into batches so that the sum of maxDistinctCount per batch stays at or below largeIndexLimit. Files that individually exceed the limit are isolated into single-file batches.
  3. 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)
  4. Periodic consolidation — Every stagingConsolidationThreshold batches, consolidateStaging merges staging into the main index via Delta MERGE (upsert on filename). Auto-compaction may follow via maybeAutoCompact.
  5. 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

Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. IndexBuildOperations
  2. BloomFilterOperations
  3. IndexFileOperations
  4. IndexMetadataOperations
  5. AriadneContextUser
  6. AnyRef
  7. 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

  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

Abstract Value Members

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

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

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

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

  10. 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
  11. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  12. 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
  13. 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
  14. 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
  15. 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
  16. 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
  17. 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
  18. 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
  19. 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
  20. 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
  21. 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
  22. 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 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
  23. 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
  24. 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
  25. 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
  26. 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 filename column and all indexed source columns

    returns

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

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

    Attributes
    protected
  29. 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
  30. 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
  31. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native() @HotSpotIntrinsicCandidate()
  32. 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
  33. 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
  34. 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
  35. 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

  36. def createOptimalBatches(fileAnalyses: Seq[Index.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
  37. 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
  38. 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

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

  40. 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
  41. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  42. def equals(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  43. 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

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

    large

    the classification identifying which files are large for which columns

    Attributes
    protected
  49. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  50. 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
  51. 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
  52. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  53. 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
  54. 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
  55. 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
  56. 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
  57. 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.

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

  59. 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
  60. 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
  61. 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
  62. 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
  63. 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
  64. 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
  65. 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
  66. 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
  67. 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
  68. 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
  69. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  70. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  71. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  72. 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

  73. 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
  74. 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
  75. 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

  76. 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
  77. 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
    Exceptions thrown

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

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

    IllegalArgumentException if the column already carries a different index type

  79. def requireNonReservedStagingColumn(column: String): Unit
    Attributes
    protected
  80. 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
    Exceptions thrown

    IllegalArgumentException if the column is a nested path

  81. 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
  82. 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
  83. 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
  84. 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
  85. 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

    IllegalArgumentException if the configuration key is not set

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

  87. final def synchronized[T0](arg0: ⇒ T0): T0
    Definition Classes
    AnyRef
  88. def toString(): String
    Definition Classes
    AnyRef → Any
  89. 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
    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.

  90. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  91. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native()
  92. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  93. 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 BloomFilterOperations

Inherited from IndexFileOperations

Inherited from IndexMetadataOperations

Inherited from AriadneContextUser

Inherited from AnyRef

Inherited from Any

Ungrouped