analyzer.core.executors.dask_exec

Attributes

Exceptions

AnalyzerRuntimeError

A combination of multiple unrelated exceptions.

Classes

RateColumn

Base class for a widget to use in progress display.

DaskRunException

DaskRunResult

LocalDaskExecutor

Helper class that provides a standard way to create an ABC using

LPCCondorDask

Helper class that provides a standard way to create an ABC using

Functions

configureDask()

callTimeoutCloud(process_timeout, function, *args, ...)

callTimeout(process_timeout, function, *args, **kwargs)

iaddMany(to_add)

reduceResults(client, reduction_function, futures[, ...])

runWithFinalize(analyzer, *args, **kwargs)

getAnalyzerRunFunc(analyzer, task[, timeout])

dumpAndComplete(metadata, output_name, dask_result)

processTask(client, analyzer, task, chunk_size, ...[, ...])

run(client, chunk_size, reduction_factor, analyzer, tasks)

Module Contents

analyzer.core.executors.dask_exec.logger[source]
class analyzer.core.executors.dask_exec.RateColumn(table_column: rich.table.Column | None = None)[source]

Bases: rich.progress.ProgressColumn

Base class for a widget to use in progress display.

render(task: rich.progress.Task) rich.text.Text[source]

Should return a renderable object.

analyzer.core.executors.dask_exec.configureDask()[source]
analyzer.core.executors.dask_exec.callTimeoutCloud(process_timeout, function, *args, **kwargs)[source]
analyzer.core.executors.dask_exec.callTimeout(process_timeout, function, *args, **kwargs)[source]
exception analyzer.core.executors.dask_exec.AnalyzerRuntimeError[source]

Bases: ExceptionGroup

A combination of multiple unrelated exceptions.

derive(excs)[source]
class analyzer.core.executors.dask_exec.DaskRunException[source]
chunk: analyzer.core.event_collection.FileChunk[source]
exception: Exception[source]
class analyzer.core.executors.dask_exec.DaskRunResult[source]
maybe_result: analyzer.core.executors.executor.CompletedTask | analyzer.core.results.ResultBase | None[source]
maybe_exceptions: list[DaskRunException][source]
events_processed: int = 0[source]
__iadd__(other: DaskRunResult)[source]
analyzer.core.executors.dask_exec.iaddMany(to_add)[source]
analyzer.core.executors.dask_exec.reduceResults(client, reduction_function, futures, reduction_factor=5, target_final_count=1, close_to_target_frac=0.8, key_suffix='')[source]
analyzer.core.executors.dask_exec.runWithFinalize(analyzer, *args, **kwargs)[source]
analyzer.core.executors.dask_exec.getAnalyzerRunFunc(analyzer, task, timeout=120)[source]
analyzer.core.executors.dask_exec.dumpAndComplete(metadata, output_name, dask_result)[source]
analyzer.core.executors.dask_exec.processTask(client, analyzer, task, chunk_size, reduction_factor, max_sample_events=None, timeout=120, target_final_count=1)[source]
analyzer.core.executors.dask_exec.run(client, chunk_size, reduction_factor, analyzer, tasks, max_sample_events=None, timeout=120, target_final_count=1)[source]
class analyzer.core.executors.dask_exec.LocalDaskExecutor[source]

Bases: analyzer.core.executors.executor.Executor

Helper class that provides a standard way to create an ABC using inheritance.

max_workers: int[source]
min_workers: int[source]
worker_memory: str = '4GB'[source]
dashboard_address: str = 'localhost:8789'[source]
schedd_address: str | None = 'localhost:12358'[source]
adapt: bool = True[source]
chunk_size: int | None = 100000[source]
processes: bool = True[source]
cluster: Any = None[source]
client: Any = None[source]
reduction_factor: int = 2[source]
target_final_count: int = 1[source]
timeout: int = 600[source]
setup(needed_resources)[source]
run(analyzer, tasks, max_sample_events=None)[source]
class analyzer.core.executors.dask_exec.LPCCondorDask[source]

Bases: analyzer.core.executors.executor.Executor

Helper class that provides a standard way to create an ABC using inheritance.

container: str[source]
venv_path: str = '.venv'[source]
x509_path: str | None = None[source]
log_path: str = 'logs/condor'[source]
worker_timeout: int | None = 7200[source]
min_workers: int = 1[source]
max_workers: int = 10[source]
worker_memory: str = '4GB'[source]
dashboard_address: str | None = 'localhost:8789'[source]
schedd_address: str | None = 'localhost:12358'[source]
adapt: bool = True[source]
chunk_size: int | None = 100000[source]
reduction_factor: int = 5[source]
timeout: int = 1200[source]
cluster: Any = None[source]
client: Any = None[source]
target_final_count: int = 1[source]
run(analyzer, tasks, max_sample_events=None)[source]
setup(needed_resources)[source]