case class IndexLock(lockPath: Path, indexName: String)(implicit spark: SparkSession) extends AriadneContextUser with Product with Serializable
File-based distributed lock for index operations.
Provides mutual exclusion for index operations (update, compact, vacuum) using lock files on the configured Hadoop
FileSystem. Lock files are JSON documents containing a LockInfo payload.
Concurrency semantics
Mutual exclusion rests on a single filesystem primitive: the lock file is created with overwrite = false, which
HDFS resolves atomically. Exactly one concurrent writer creates the file; every other writer is rejected with
FileAlreadyExistsException and enters a retry loop with exponential back-off (capped at 60 s) up to lockMaxWait
seconds.
Filesystem requirements
That guarantee is HDFS semantics, and Ariadne relies on the configured FileSystem implementation to honour them.
Any filesystem that resolves create-if-absent atomically inherits the same mutual exclusion; any filesystem that does
not — because it emulates the check client-side, or reports metadata only eventually — weakens it, and two callers
may believe they hold the same lock. Ariadne cannot detect or compensate for this from above the FileSystem
interface, so locking is best effort: each filesystem is responsible for the atomicity of its own
create-if-absent operation, and correctness under concurrency is only as strong as that operation.
Stale-lock healing is not atomic on any filesystem, including HDFS. Deleting the stale file and creating a
replacement are separate operations, so two processes that both judge a lock stale can both delete and one can
overwrite the other's fresh lock. Only the create step is protected; see the TOCTOU note on handleExistingLock.
Stale lock healing
If the lastRefreshedAt timestamp is older than lockTimeout seconds, the lock is considered stale. A stale lock is
automatically deleted and re-acquired, with a warning logged identifying the previous holder.
Corrupt lock file handling
If the lock file exists but cannot be parsed (empty or invalid JSON), it is treated as corrupt: the file is deleted
and acquisition is retried, up to MaxCorruptLockRetries times to prevent infinite recursion.
Thread safety: IndexLock instances are not safe for concurrent use from multiple threads. The
readLock/writeLockFile methods perform unsynchronized filesystem I/O. Each thread should use its own IndexLock
instance or coordinate externally.
- lockPath
Path to the lock file on the filesystem
- indexName
Name of the index being locked (used in log messages)
- spark
Implicit SparkSession for filesystem access and configuration
- Alphabetic
- By Inheritance
- IndexLock
- Serializable
- Serializable
- Product
- Equals
- AriadneContextUser
- AnyRef
- Any
- Hide All
- Show All
- Public
- All
Instance Constructors
-
new
IndexLock(lockPath: Path, indexName: String)(implicit spark: SparkSession)
- lockPath
Path to the lock file on the filesystem
- indexName
Name of the index being locked (used in log messages)
- spark
Implicit SparkSession for filesystem access and configuration
Value Members
-
final
def
!=(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
-
final
def
##(): Int
- Definition Classes
- AnyRef → Any
-
final
def
==(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
-
def
acquire(correlationId: String): Unit
Acquires the lock for the given correlation ID.
Acquires the lock for the given correlation ID.
If the lock file already exists, this method will either heal a stale or corrupt lock, or enter a retry loop until the lock becomes available or
lockMaxWaitis exceeded.- correlationId
unique identifier to associate with this lock hold
- Exceptions thrown
IllegalArgumentExceptionif correlationId is null or blankIndexLockExceptionif the lock cannot be acquired within the configured timeout or after max retry attempts
-
final
def
asInstanceOf[T0]: T0
- Definition Classes
- Any
-
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
-
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
-
def
clone(): AnyRef
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws( ... ) @native() @HotSpotIntrinsicCandidate()
-
lazy val
debugEnabled: Boolean
When true, logs detailed diagnostics during join operations including partition counts, physical plans, and per-phase timing.
When true, logs detailed diagnostics during join operations including partition counts, physical plans, and per-phase timing. Reads from spark.ariadne.debug configuration (default: false).
- Definition Classes
- AriadneContextUser
-
def
delete(path: Path): Boolean
Deletes a path recursively from the filesystem.
Deletes a path recursively from the filesystem.
- path
The Hadoop Path to delete
- returns
true if the path was successfully deleted
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulljava.io.IOExceptionif the filesystem operation fails
-
def
delta(path: Path): Option[DeltaTable]
Returns a io.delta.tables.DeltaTable if the path exists and contains valid Delta metadata.
Returns a io.delta.tables.DeltaTable if the path exists and contains valid Delta metadata.
Three outcomes are possible:
- the path holds a valid Delta table —
Some(DeltaTable); - the path is absent, or exists but is an empty directory —
None, so the caller may create the table there; - the path exists, is non-empty, and holds no readable Delta metadata — dev.cjfravel.ariadne.exceptions.InvalidDeltaTableException.
The third case is reported rather than repaired. Ariadne cannot distinguish an abandoned partial write of its own from data written by something else, so it never deletes the directory; cleanup is left to the caller, who can inspect the contents first. An empty directory is treated as absent because there is nothing to inspect and nothing to lose.
- path
The Hadoop Path to check for a Delta table
- returns
Some(DeltaTable)if a valid Delta table exists,Noneif the path is absent or an empty directory
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulldev.cjfravel.ariadne.exceptions.InvalidDeltaTableExceptionif the path exists and is non-empty but is not a readable Delta table
- the path holds a valid Delta table —
-
final
def
eq(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
-
def
exists(path: Path): Boolean
Checks if a path exists on the filesystem.
Checks if a path exists on the filesystem.
- path
The Hadoop Path to check
- returns
true if the path exists
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulljava.io.IOExceptionif the filesystem operation fails
-
lazy val
fs: FileSystem
Hadoop
FileSysteminstance resolved from the storagePath URI and the Spark Hadoop configuration.Hadoop
FileSysteminstance resolved from the storagePath URI and the Spark Hadoop configuration. Used for all filesystem operations (existence checks, reads, deletes, lock files).- Definition Classes
- AriadneContextUser
-
final
def
getClass(): Class[_]
- Definition Classes
- AnyRef → Any
- Annotations
- @native() @HotSpotIntrinsicCandidate()
- val indexName: String
-
lazy val
indexRepartitionCount: Option[Int]
Optional number of partitions to use when repartitioning the index DataFrame during joins.
Optional number of partitions to use when repartitioning the index DataFrame during joins. When set, the index DataFrame is repartitioned before expensive operations like explode to reduce per-executor memory pressure and avoid FetchFailedExceptions on large indexes. Reads from spark.ariadne.indexRepartitionCount configuration (default: not set).
- Definition Classes
- AriadneContextUser
-
final
def
isInstanceOf[T0]: Boolean
- Definition Classes
- Any
-
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
-
lazy val
lockMaxWait: Long
Maximum total time in seconds to wait for lock acquisition before failing.
Maximum total time in seconds to wait for lock acquisition before failing. Reads from spark.ariadne.lockMaxWait configuration (default: 3600).
- Definition Classes
- AriadneContextUser
- Note
Unlike
lockTimeoutandlockRetryInterval, non-positive values are not validated — a value of 0 causes immediate timeout on first retry. Configure with care.
- val lockPath: Path
-
lazy val
lockRefreshInterval: Int
Refresh the update lock every N batches during updateBatched.
Refresh the update lock every N batches during updateBatched. Reads from spark.ariadne.lockRefreshInterval configuration (default: 1).
- Definition Classes
- AriadneContextUser
-
lazy val
lockRetryInterval: Long
Base interval in seconds between lock acquisition retries (exponential backoff applied).
Base interval in seconds between lock acquisition retries (exponential backoff applied). Reads from spark.ariadne.lockRetryInterval configuration (default: 60).
- Definition Classes
- AriadneContextUser
-
lazy val
lockTimeout: Long
Seconds since
lastRefreshedAtbefore a lock is considered stale and eligible for auto-heal (forcible acquisition by another process).Seconds since
lastRefreshedAtbefore a lock is considered stale and eligible for auto-heal (forcible acquisition by another process). Reads fromspark.ariadne.lockTimeout(default: 1800).- Definition Classes
- AriadneContextUser
-
lazy val
logger: Logger
Log4j logger shared by all Ariadne components.
Log4j logger shared by all Ariadne components. Uses the "ariadne" logger name. Marked
@transient lazyso thatIndex(a case class) remains serializable across Spark stages — Log4j loggers are notSerializable.- Definition Classes
- AriadneContextUser
- Annotations
- @transient()
-
final
def
ne(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
-
final
def
notify(): Unit
- Definition Classes
- AnyRef
- Annotations
- @native() @HotSpotIntrinsicCandidate()
-
final
def
notifyAll(): Unit
- Definition Classes
- AnyRef
- Annotations
- @native() @HotSpotIntrinsicCandidate()
-
def
open(path: Path): FSDataInputStream
Opens an input stream to read from a path on the filesystem.
Opens an input stream to read from a path on the filesystem.
- path
The Hadoop Path to open for reading
- returns
An FSDataInputStream for reading the file contents
- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif path is nulljava.io.IOExceptionif the file does not exist or cannot be opened
-
def
refresh(correlationId: String): Unit
Refreshes the lock's
lastRefreshedAttimestamp to prevent stale-lock healing.Refreshes the lock's
lastRefreshedAttimestamp to prevent stale-lock healing.Long-running operations should call this periodically (more frequently than
lockTimeout) to signal that the lock holder is still active.- correlationId
the correlation ID that should currently hold the lock
- Exceptions thrown
IllegalArgumentExceptionif correlationId is null or blank
-
def
release(correlationId: String): Unit
Releases the lock if it is held by the given correlation ID.
Releases the lock if it is held by the given correlation ID.
If the lock file does not exist or belongs to a different correlation ID, a warning is logged and no action is taken.
- correlationId
the correlation ID that should currently hold the lock
- Exceptions thrown
IllegalArgumentExceptionif correlationId is null or blank
-
lazy val
repartitionDataFiles: Boolean
When true, applies indexRepartitionCount repartitioning to the data files read during joinDf.
When true, applies indexRepartitionCount repartitioning to the data files read during joinDf. When false, data files keep their natural parquet partitioning. Disable when column selection significantly reduces data volume, making the repartition shuffle more expensive than useful. Reads from spark.ariadne.repartitionDataFiles configuration (default: false).
- Definition Classes
- AriadneContextUser
-
def
safeDestroyBroadcast(broadcast: Broadcast[_]): Unit
Safely destroys a broadcast variable, logging but swallowing any exception.
Safely destroys a broadcast variable, logging but swallowing any exception.
Intended for use inside
finallyblocks where a cleanup failure must not mask the original exception propagating out of thetryblock. Tolerates a null broadcast.- broadcast
the broadcast to destroy; may be null
- Attributes
- protected
- Definition Classes
- AriadneContextUser
-
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
- IndexLock → AriadneContextUser
-
lazy val
stagingConsolidationThreshold: Int
Number of batches to process before consolidating staged data into the main index.
Number of batches to process before consolidating staged data into the main index. This provides fault tolerance for large index builds - if a job fails, work is preserved up to the last consolidation point. Reads from spark.ariadne.stagingConsolidationThreshold configuration (default: 50).
- Definition Classes
- AriadneContextUser
-
lazy val
storagePath: Path
Base path on the Hadoop-compatible filesystem where all Ariadne index data is stored.
Base path on the Hadoop-compatible filesystem where all Ariadne index data is stored. Each index creates a subdirectory under this path.
Reads from
spark.ariadne.storagePath(required — no default).- Definition Classes
- AriadneContextUser
- Exceptions thrown
IllegalArgumentExceptionif the configuration key is not set
-
final
def
synchronized[T0](arg0: ⇒ T0): T0
- Definition Classes
- AnyRef
-
final
def
wait(arg0: Long, arg1: Int): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws( ... )
-
final
def
wait(arg0: Long): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws( ... ) @native()
-
final
def
wait(): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws( ... )
Deprecated Value Members
-
def
finalize(): Unit
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws( classOf[java.lang.Throwable] ) @Deprecated
- Deprecated