Skip to main content

heidelberg_source

Data source for loading ophthalmology files using private-eye.

Classes​

HeidelbergCSVColumns​

class HeidelbergCSVColumns(heidelberg_files_col: str = 'heidelberg_file'):

Arguments for ophthalmology columns in the csv.

Arguments

  • heidelberg_files_col: The name of the column that points to Heidelberg files in the CSV file. Defaults to 'heidelberg_file'. These files should be in .sdb format. For .e2e files, use HeidelbergE2ESource.

Variables​

  • static heidelberg_files_col : str

HeidelbergE2EAllSeriesSource​

class HeidelbergE2EAllSeriesSource(    private_eye_parser: PrivateEyeParser | Mapping[str, PrivateEyeParser],    path: os.PathLike[str] | str,    heidelberg_csv_columns: HeidelbergCSVColumns | _HeidelbergCSVColumnsTD | None = None,    required_fields: dict[str, Any] | None = None,    parsers: PrivateEyeParser | Mapping[str, PrivateEyeParser] | None = None,    ophthalmology_args: OphthalmologyDataSourceArgs | _OphthalmologyDataSourceArgsTD | None = None,    data_cache: DataPersister | None = None,    infer_class_labels_from_filepaths: bool = False,    output_path: os.PathLike[str] | str | None = None,    iterable: bool = True,    fast_load: bool = True,    cache_images: bool = False,    filter: FileSystemFilter | None = None,    data_splitter: DatasetSplitter | None = None,    seed: int | None = None,    ignore_cols: str | Sequence[str] | None = None,    modifiers: dict[str, DataPathModifiers] | None = None,    partition_size: int = 16,    name: str | None = None,):

Data source returning every series in a Heidelberg .e2e file.

The file-intrinsic counterpart to HeidelbergE2ESource: it applies no laterality/protocol selection and no one-row-per-file constraint (the inherited series hooks are pass-throughs), so a file yields the same rows no matter which task configuration triggered the read, and the selection is applied later at read time.

Used by the scan_metadata runtime, which records what a file intrinsically contains. It is deliberately not exposed via pod YAML: the one-row-per-file contract that the data cache, batch assembly and the downstream algorithms rely on holds only for HeidelbergE2ESource.

Arguments

  • path: The path to the directory containing the Heidelberg .e2e files or to a CSV file that includes Heidelberg files as one of its columns.
  • heidelberg_csv_columns: If path is a CSV file, this contains information about the specific columns that contain the path information for the Heidelberg files. Defaults to None.
  • **kwargs: Keyword arguments passed to the parent base classes.

See

Variables​

  • accessibility_details : AccessibilityDetails | None - Detailed accessibility status. None if accessible.

    Subclasses should override to perform lightweight connectivity check.

    Returns: None if accessible, or dict with 'error_code' and 'message' if not.

  • file_names : list[str] - Returns a list of file names in the specified directory.

    warning

Deprecated: The file_names property is deprecated and will be removed in a future release. Use file_names_iter(as_strs=True) for memory-efficient iteration, or list(file_names_iter(as_strs=True)) if you need a list. :::

This property accounts for files skipped at runtime by filtering them out of the list of cached file names. Files may get skipped at runtime due to errors or because they don't contain any image data and images_only is True. This allows us to skip these files again more quickly if they are still present in the directory.

  • is_accessible : bool - Check if datasource is currently accessible.

    Returns True if accessibility_details is None (no errors). This is a convenience property that wraps accessibility_details.

  • is_file_iterable : bool - Returns True since this source iterates over files.
  • is_initialised : bool - Checks if BaseSource was initialised.
  • is_task_running : bool - Returns True if a task is running.
  • name : str | None - The datasource's configured name.
  • path : pathlib.Path - Resolved absolute path to data.

    Provides a consistent version of the path provided by the user which should work throughout regardless of operating system and of directory structure.

  • selected_file_names : list[str] - Returns a list of selected file names as strings.

    Selected file names are affected by the selected_file_names_override and new_file_names_only attributes.

    WARNING: This method loads all filenames into memory. For large datasets, consider using selected_file_names_iter() instead.

  • selected_file_names_differ : bool - Returns True if selected_file_names will differ from default.

    In particular, returns True iff there is a selected file names override in place and/or there is filtering for new file names only present.

  • supports_project_db : bool - Whether the datasource supports the project database.

    Each datasource needs to implement its own methods to define how what its project database table should look like. If the datasource does not implement the methods to get the table creation query and columns, it does not support the projectdatabase.

  • task_skipped_file_names : set[str] - Return set of task-skipped filenames for set operations.

Static methods​


get_num_workers​

def get_num_workers(file_names: Sequence[str]) ‑> int:

Inherited from:

HeidelbergSource.get_num_workers :

Gets the number of workers to use for multiprocessing.

Ensures that the number of workers is at least 1 and at most equal to MAX_NUM_MULTIPROCESSING_WORKERS. If the number of files is less than MAX_NUM_MULTIPROCESSING_WORKERS, then we use the number of files as the number of workers. Unless the number of machine cores is also less than MAX_NUM_MULTIPROCESSING_WORKERS, in which case we use the lower of the two.

Arguments

  • file_names: The list of file names to load.

Returns The number of workers to use for multiprocessing.

Methods​


add_hook​

def add_hook(self, hook: DataSourceHook) ‑> None:

Inherited from:

HeidelbergSource.add_hook :

Add a hook to the datasource.

add_strategy_filter_flag_columns​

def add_strategy_filter_flag_columns(    self, df: pd.DataFrame,) ‑> pandas.core.frame.DataFrame:

Inherited from:

HeidelbergSource.add_strategy_filter_flag_columns :

Add task-scoped strategy filter flag columns to a report DataFrame.

Decodes the compact per-file bitmask into one boolean column per flag-only strategy. Only files present in df are decoded, keeping memory proportional to the report size rather than the total number of evaluated files.

Arguments

  • df: DataFrame containing at least the ORIGINAL_FILENAME_METADATA_COLUMN column.

Returns A new DataFrame with one boolean column per flag-only strategy. If no flag-only strategies are registered, returns df unchanged (same object, no copy). If df does not contain the ORIGINAL_FILENAME_METADATA_COLUMN (e.g. patient-level reports), flag injection is skipped gracefully and df is returned unchanged.

apply_ignore_cols​

def apply_ignore_cols(self, df: pd.DataFrame) ‑> pandas.core.frame.DataFrame:

Inherited from:

HeidelbergSource.apply_ignore_cols :

Apply ignored columns to dataframe, dropping columns as needed.

Returns A copy of the dataframe with ignored columns removed, or the original dataframe if this datasource does not specify any ignore columns.

apply_ignore_cols_iter​

def apply_ignore_cols_iter(    self, dfs: Iterator[pd.DataFrame],) ‑> collections.abc.Iterator[pandas.core.frame.DataFrame]:

Inherited from:

HeidelbergSource.apply_ignore_cols_iter :

Apply ignored columns to dataframes from iterator.

apply_merged_filter_config​

def apply_merged_filter_config(self, merged_filter_config: MergedFilterConfig) ‑> None:

Inherited from:

HeidelbergSource.apply_merged_filter_config :

Apply the filter configuration to the datasource.

apply_modifiers​

def apply_modifiers(self, df: pd.DataFrame) ‑> pandas.core.frame.DataFrame:

Inherited from:

HeidelbergSource.apply_modifiers :

Apply column modifiers to the dataframe.

If no modifiers are specified, returns the dataframe unchanged.

clear_dataset_cache​

def clear_dataset_cache(self) ‑> dict[str, typing.Any]:

Inherited from:

HeidelbergSource.clear_dataset_cache :

Clear all dataset cache for this data source.

This clears both:

  1. The file names cache (Python cached_property)
  2. The dataset cache file (deletes the SQLite database file completely)

Returns Dictionary with cache clearing results.

clear_file_names_cache​

def clear_file_names_cache(self) ‑> None:

Inherited from:

HeidelbergSource.clear_file_names_cache :

Clears the list of selected file names.

This allows the datasource to pick up any new files that have been added to the directory since the last time it was cached.

clear_task_specific_configs​

def clear_task_specific_configs(self) ‑> None:

Inherited from:

HeidelbergSource.clear_task_specific_configs :

Clear task-scoped state at task boundaries.

Sends the pending skipped-file telemetry summary, then resets skipped files, merged filter config, and strategy filter flags so that subsequent tasks start with a clean slate.

file_names_iter​

def file_names_iter(    self, as_strs: bool = False,) ‑> collections.abc.Iterator[pathlib.Path] | collections.abc.Iterator[str]:

Inherited from:

HeidelbergSource.file_names_iter :

Iterate over files in a directory, yielding those that match the criteria.

The datasource-specific filters are applied a chunk of files at a time rather than one file at a time, so each filter reads what it needs for the whole chunk in one query instead of one per file — see _DATASOURCE_FILTER_BATCH_SIZE. Files are still yielded individually, in walk order, and a file dropped by any filter is bookkept exactly as before: the filters record their own specific skip reason as they go, and the generic DATASOURCE_FILTER_FAILED below remains the backstop for a file that none of them accounted for.

Arguments

  • as_strs: By default the files yielded will be yielded as Path objects. If this is True, yield them as strings instead.

files_passing_filters​

def files_passing_filters(    self, file_names: Iterable[str | os.PathLike[str]],) ‑> list[str]:

Inherited from:

HeidelbergSource.files_passing_filters :

Return which of file_names this datasource's filters admit.

The same decision the file walk makes as it yields, for a caller holding a list of paths from elsewhere — the DAG's file-metadata inventory, which is keyed by path and cannot answer a filter on a file's contents (frame count, modality, inferred scan type). Includes any task-level filters currently applied via apply_merged_filter_config.

Arguments

  • file_names: Candidate paths.

Returns Those that pass, in the order the underlying filters return them.

get_all_cached_file_paths​

def get_all_cached_file_paths(self) ‑> list[str]:

Inherited from:

HeidelbergSource.get_all_cached_file_paths :

Get all file paths that are currently stored in the cache.

Returns A list of file paths that have cache entries, or an empty list if there is no cache or the cache hasn't been initialized.

get_data​

def get_data(    self,    data_keys: SingleOrMulti[str] | SingleOrMulti[int],    *,    use_cache: bool = True,    **kwargs: Any,) ‑> pandas.core.frame.DataFrame | None:

Inherited from:

HeidelbergSource.get_data :

Get data corresponding to the provided data key(s).

Can be used to return data for a single data key or for multiple at once. If used for multiple, the order of the output dataframe must match the order of the keys provided.

Arguments

  • data_keys: Key(s) for which to get the data of. These may be things such as file names, UUIDs, etc. Can also be a list of integers if the datasource has an integer index.
  • use_cache: Whether the cache should be used to retrieve data for these keys. Note that cached data may have some elements, particularly image-related fields such as image data or file paths, replaced with placeholder values when stored in the cache. If data_cache is set on the instance, data will be set in the cache, regardless of this argument.
  • **kwargs: Additional keyword arguments.

Returns A dataframe containing the data, ordered to match the order of keys in data_keys, or None if no data for those keys was available.

get_datasource_metrics​

def get_datasource_metrics(    self, use_skip_codes: bool = False, data: pd.DataFrame | None = None,) ‑> DatasourceSummaryStats:

Inherited from:

HeidelbergSource.get_datasource_metrics :

Get metadata about this datasource.

This can be used to store information about the datasource that may be useful for debugging or tracking purposes. The metadata will be stored in the project database.

Arguments

  • use_skip_codes: Whether to use the skip reason codes as the keys in the skip_reasons dictionary, rather than the existing reason descriptions.
  • data: The data to use for getting the metrics.

Returns A dictionary containing metadata about this datasource.

get_filter_config​

def get_filter_config(self) ‑> FilterConfig:

Inherited from:

HeidelbergSource.get_filter_config :

Get the filter configuration for the datasource.

get_project_db_sqlite_columns​

def get_project_db_sqlite_columns(self) ‑> list[str]:

Inherited from:

HeidelbergSource.get_project_db_sqlite_columns :

Returns the required columns to identify a data point.

The first value must be filename column, and second value must be the last modified datetime. These two are used to build the processed_file_cache for the worker execution.

get_project_db_sqlite_create_table_query​

def get_project_db_sqlite_create_table_query(self) ‑> str:

Inherited from:

HeidelbergSource.get_project_db_sqlite_create_table_query :

Returns the required columns and types to identify a data point.

The file name is used as the primary key and the last modified date is used to determine if the file has been updated since the last time it was processed. If there is a conflict on the file name, the row is replaced with the new data to ensure that the last modified date is always up to date.

get_schema​

def get_schema(self) ‑> dict[str, typing.Any]:

Inherited from:

HeidelbergSource.get_schema :

Get the pre-defined schema for this datasource.

This method should be overridden by datasources that have pre-defined schemas (i.e., those with has_predefined_schema = True).

Returns The schema as a dictionary.

Raises

  • NotImplementedError: If the datasource doesn't have a pre-defined schema.

get_task_skip_reason_summary​

def get_task_skip_reason_summary(self) ‑> dict[str, int]:

Inherited from:

HeidelbergSource.get_task_skip_reason_summary :

Get aggregated skip reasons for the current task.

Combines both task-only skips (files that failed task filters) and datasource skips that occurred during the current task execution. This provides a complete picture of all files skipped during a task run.

Returns Dict mapping reason codes (as strings) to file counts.

get_uncached_file_names​

def get_uncached_file_names(self) ‑> list[str]:

Inherited from:

HeidelbergSource.get_uncached_file_names :

Return potentially uncached files via fast raw filesystem scanning.

This fast path skips datasource filters and computes file-path set difference against cache/skipped tables. It may include files that are later filtered out during normal datasource processing.

has_uncached_files​

def has_uncached_files(self) ‑> bool:

Inherited from:

HeidelbergSource.has_uncached_files :

Returns True if there are any files in the datasource not yet cached.

Uses a fast path that skips the full filter pipeline: walks the filesystem with scantree and checks each path against cache + skipped files metadata.

is_file_skipped_at_task_level​

def is_file_skipped_at_task_level(self, filename: str) ‑> bool:

Inherited from:

HeidelbergSource.is_file_skipped_at_task_level :

Whether filename carries a task-level skip from the current task.

A task-level skip marks a failure the datasource judged transient, so it is deliberately not persisted to the cache. Callers that record their own per-file outcome need this to tell "read the file, it holds nothing" from "could not read the file this time", which both surface as an empty result from _process_file.

Arguments

  • filename: Path as passed to the call that may have skipped it.

Returns True when the current task recorded a task-level skip for it.

log_parallel_capacity​

def log_parallel_capacity(    self, file_names: Sequence[str], chunk_size: int | None = None,) ‑> None:

Inherited from:

HeidelbergSource.log_parallel_capacity :

Log what parallelism this machine and configuration allow.

None of it is otherwise recorded, and all of it bounds what any change to how files are read could be worth: a site with two cores cannot gain from a wider pool however the code is arranged, and a site whose partitions hold one file never reaches the parallel path at all. Without these figures a timing report from a deployment cannot be read.

active and workers describe one partition, not the whole run, because that is the unit the parallel decision is actually made on. A caller that narrows the datasource before reading - the inference step hands it one batch_size chunk at a time - must say so via chunk_size, or the figures describe a partition no read ever uses and report parallelism that never happens.

Arguments

  • file_names: Every file about to be processed, across all chunks.
  • chunk_size: How many of them the caller narrows the datasource to at a time. Defaults to all of them, for a caller that does not narrow.

merge_and_validate_filters​

def merge_and_validate_filters(    self, datasource_level_filters: FilterConfig, task_level_filters: list[TaskFilter],) ‑> MergedFilterConfig:

Inherited from:

HeidelbergSource.merge_and_validate_filters :

Merge and validate the filters from the datasource and the task.

Returns a MergedFilterConfig with resolved filter values from both datasource and task-level filters, using intersection logic (most restrictive wins).

partition​

def partition(    self, iterable: Iterable[_I], partition_size: int = 1,) ‑> collections.abc.Iterable[collections.abc.Sequence[~_I]]:

Inherited from:

HeidelbergSource.partition :

Partition the iterable into chunks of the given size.

process_file​

def process_file(    self, filename: str, skip_non_tabular_data: bool = False, **kwargs: Any,) ‑> list[dict[str, typing.Any]]:

Inherited from:

HeidelbergSource.process_file :

Parse a single file into records, outside the batch-loading path.

The entry point for callers that need one file's records, such as the scan_metadata runtime. Unlike loading a batch it applies no data cache and adds no metadata columns; the only bookkeeping is what the datasource itself does while processing, such as recording a skip. Subclasses implement _process_file, not this.

Arguments

  • filename: The name of the file to process.
  • skip_non_tabular_data: Whether non-tabular data, e.g. image data, can be left unloaded.
  • **kwargs: Additional keyword arguments for _process_file.

Returns One dictionary per datapoint in the file, as _process_file returns.

remove_hook​

def remove_hook(self, hook: DataSourceHook) ‑> None:

Inherited from:

HeidelbergSource.remove_hook :

Remove a hook from the datasource.

selected_file_names_iter​

def selected_file_names_iter(self) ‑> collections.abc.Iterator[str]:

Inherited from:

HeidelbergSource.selected_file_names_iter :

Returns an iterator over selected file names.

Selected file names are affected by the selected_file_names_override and new_file_names_only attributes.

Returns Iterator over selected file names.

set_strategy_filter_flags​

def set_strategy_filter_flags(    self, flags: dict[str, int], meta: list[tuple[int, str, str | None]],) ‑> None:

Inherited from:

HeidelbergSource.set_strategy_filter_flags :

Store strategy filter flag-only match bitmasks for the current task.

Arguments

  • flags: Mapping of filename to bitmask of matched flag-only strategies.
  • meta: Ordered list of (original_index, strategy_name, flag_column_name). Bit position in the bitmask corresponds to list index.

skip_file​

def skip_file(    self, filename: str, reason: FileSkipReason, data: dict[str, Any] | None = None,) ‑> None:

Inherited from:

HeidelbergSource.skip_file :

Skip a file by updating cache and skipped_files set.

The first reason is always the one recorded in the data cache. If the file was already skipped, the stored reason is used for reporting, logging, and telemetry instead of the caller's reason.

Arguments

  • filename: Path to the file being skipped
  • reason: Reason for skipping the file
  • data: Unused. Kept for callers that pass file metadata.

task_skip_file​

def task_skip_file(    self, filename: str, reason: FileSkipReason, data: dict[str, Any] | None = None,) ‑> None:

Inherited from:

HeidelbergSource.task_skip_file :

Skip a file due to task-level filter (not persisted to cache).

Unlike skip_file(), this does NOT persist to cache because the file may pass a different task's filters. The skip is tracked in-memory only and cleared between tasks.

Arguments

  • filename: Path to the file being skipped.
  • reason: Reason for skipping the file.
  • data: Unused. Kept for callers that pass file metadata.

task_skip_reason​

def task_skip_reason(    self, filename: str,) ‑> FileSkipReason | None:

Inherited from:

HeidelbergSource.task_skip_reason :

Why filename carries a task-level skip, or None if it does not.

Arguments

  • filename: Path as passed to the call that may have skipped it.

use_file_multiprocessing​

def use_file_multiprocessing(self, file_names: Sequence[str]) ‑> bool:

Inherited from:

HeidelbergSource.use_file_multiprocessing :

Check if file multiprocessing should be used.

Returns True if file multiprocessing has been enabled by the environment variable and the number of workers would be greater than 1, otherwise False. There is no need to use file multiprocessing if we are just going to use one worker - it would be slower than just loading the data in the main process.

Returns True if file multiprocessing should be used, otherwise False.

yield_data​

def yield_data(    self,    data_keys: SingleOrMulti[str] | SingleOrMulti[int] | None = None,    *,    use_cache: bool = True,    partition_size: int | None = None,    **kwargs: Any,) ‑> collections.abc.Iterator[pandas.core.frame.DataFrame]:

Inherited from:

HeidelbergSource.yield_data :

Yields data in batches from this source.

If data_keys is specified, only yield from that subset of the data. Otherwise, iterate through the whole datasource.

Arguments

  • data_keys: An optional list of data keys to use for yielding data. Otherwise, all data in the datasource will be considered. data_keys is always provided when this method is called from the Dataset as part of a task. Can also be a list of integers if the datasource has an integer index.
  • use_cache: Whether the cache should be used to retrieve data for these data points. Note that cached data may have some elements, particularly image-related fields such as image data or file paths, replaced with placeholder values when stored in the cache. If data_cache is set on the instance, data will be set in the cache, regardless of this argument.
  • partition_size: The number of data elements to load/yield in each iteration. If not provided, defaults to the partition size configured in the datasource.
  • **kwargs: Additional keyword arguments.

HeidelbergE2ESource​

class HeidelbergE2ESource(    private_eye_parser: PrivateEyeParser | Mapping[str, PrivateEyeParser],    path: os.PathLike[str] | str,    laterality: str,    series_protocol: str,    heidelberg_csv_columns: HeidelbergCSVColumns | _HeidelbergCSVColumnsTD | None = None,    required_fields: dict[str, Any] | None = None,    parsers: PrivateEyeParser | Mapping[str, PrivateEyeParser] | None = None,    ophthalmology_args: OphthalmologyDataSourceArgs | _OphthalmologyDataSourceArgsTD | None = None,    data_cache: DataPersister | None = None,    infer_class_labels_from_filepaths: bool = False,    output_path: os.PathLike[str] | str | None = None,    iterable: bool = True,    fast_load: bool = True,    cache_images: bool = False,    filter: FileSystemFilter | None = None,    data_splitter: DatasetSplitter | None = None,    seed: int | None = None,    ignore_cols: str | Sequence[str] | None = None,    modifiers: dict[str, DataPathModifiers] | None = None,    partition_size: int = 16,    name: str | None = None,):

Data source for loading Heidelberg .e2e files.

E2E files may contain multiple series (e.g. different lateralities or scan protocols within a single file). Because the datasource maintains a strict one-row-per-file contract (required by caching, batch assembly, and downstream algorithms), series selection is performed at parse time using the required laterality and series_protocol arguments.

These are datasource-level only and cannot be specified as task-level filters because the data cache is keyed by filename — a cached row corresponds to the series selected at parse time. To process both eyes from the same e2e file, create separate HeidelbergE2ESource instances with different laterality values.

If filters match multiple series within a single file, only the first matching series is used (with a warning).

Arguments

  • **kwargs: Keyword arguments passed to the parent base classes.
  • cache_images: Whether to cache images in the file system. Defaults to False. This is ignored if fast_load is True.
  • data_cache: A DataPersister instance to use for data caching.
  • data_splitter: Deprecated argument, will be removed in a future release. Defaults to None. Not used.
  • fast_load: Whether the data will be loaded in fast mode. This is used to determine whether the data will be iterated over during set up for schema generation and splitting (where necessary). Only relevant if iterable is True, otherwise it is ignored. Defaults to True.
  • heidelberg_csv_columns: If path is a CSV file, this contains information about the specific columns that contain the path information for the Heidelberg files. Defaults to None.
  • ignore_cols: Column/list of columns to be ignored from the data. Defaults to None.
  • infer_class_labels_from_filepaths: Whether class labels should be added to the data based on the filepath of the files. Defaults to the first directory within self.path, but can go a level deeper if the datasplitter is provided with infer_data_split_labels set to true
  • iterable: Whether the data source is iterable. This is used to determine whether the data source can be used in a streaming context during a task. Defaults to True.
  • laterality: The laterality to extract from multi-series e2e files. Must be one of 'L' (left), 'R' (right), 'B' (both), or 'U' (unknown).
  • modifiers: Dictionary used for modifying paths/ extensions in the dataframe. Defaults to None.
  • name: The name for the datasource. Optional, defaults to None.
  • ophthalmology_args: Arguments for ophthalmology modality data.
  • output_path: The path where to save intermediary output files. Defaults to 'preprocessed/'.
  • parsers: The private eye parsers to use for the different file extensions. Only needs to be supplied if file_extension filter is non-default. Can either be a single parser to use for all file extensions or a mapping of file extensions to parser type. Defaults to appropriate parser(s) for the default file extension(s).
  • partition_size: The size of each partition when iterating over the data in a batched fashion.
  • path: The path to the directory containing the Heidelberg .e2e files or to a CSV file that includes Heidelberg files as one of its columns.
  • private_eye_parser: Private-eye supported machine type(s). Can either be a single parser to use for all files or a mapping of file extension to the desired parser type. If private_eye_parser is a mapping of file extensions to parsers, there must be a parser for each file extension specified. If no file extensions are specified then mapping can exist in whatever form (warnings will be logged if we encounter an extension for which no parser is specified). If private_eye_parser is a single parser, file_extension can be anything, we simply try to use this parser against any extension.
  • required_fields: The fields a parsed row must contain. Defaults to _default_required_fields().
  • seed: Random number seed. Used for setting random seed for all libraries. Defaults to None.
  • series_protocol: The series protocol to extract from multi-series e2e files (e.g. 'OCT ART Volume', 'OCT B-Scan').

Attributes

  • seed: Random number seed. Used for setting random seed for all libraries.

Raises

  • ValueError: If laterality is not one of 'L', 'R', 'B', 'U'.

Variables​

  • accessibility_details : AccessibilityDetails | None - Detailed accessibility status. None if accessible.

    Subclasses should override to perform lightweight connectivity check.

    Returns: None if accessible, or dict with 'error_code' and 'message' if not.

  • file_names : list[str] - Returns a list of file names in the specified directory.

    warning

Deprecated: The file_names property is deprecated and will be removed in a future release. Use file_names_iter(as_strs=True) for memory-efficient iteration, or list(file_names_iter(as_strs=True)) if you need a list. :::

This property accounts for files skipped at runtime by filtering them out of the list of cached file names. Files may get skipped at runtime due to errors or because they don't contain any image data and images_only is True. This allows us to skip these files again more quickly if they are still present in the directory.

  • is_accessible : bool - Check if datasource is currently accessible.

    Returns True if accessibility_details is None (no errors). This is a convenience property that wraps accessibility_details.

  • is_file_iterable : bool - Returns True since this source iterates over files.
  • is_initialised : bool - Checks if BaseSource was initialised.
  • is_task_running : bool - Returns True if a task is running.
  • name : str | None - The datasource's configured name.
  • path : pathlib.Path - Resolved absolute path to data.

    Provides a consistent version of the path provided by the user which should work throughout regardless of operating system and of directory structure.

  • selected_file_names : list[str] - Returns a list of selected file names as strings.

    Selected file names are affected by the selected_file_names_override and new_file_names_only attributes.

    WARNING: This method loads all filenames into memory. For large datasets, consider using selected_file_names_iter() instead.

  • selected_file_names_differ : bool - Returns True if selected_file_names will differ from default.

    In particular, returns True iff there is a selected file names override in place and/or there is filtering for new file names only present.

  • supports_project_db : bool - Whether the datasource supports the project database.

    Each datasource needs to implement its own methods to define how what its project database table should look like. If the datasource does not implement the methods to get the table creation query and columns, it does not support the projectdatabase.

  • task_skipped_file_names : set[str] - Return set of task-skipped filenames for set operations.

Static methods​


get_num_workers​

def get_num_workers(file_names: Sequence[str]) ‑> int:

Inherited from:

HeidelbergSource.get_num_workers :

Gets the number of workers to use for multiprocessing.

Ensures that the number of workers is at least 1 and at most equal to MAX_NUM_MULTIPROCESSING_WORKERS. If the number of files is less than MAX_NUM_MULTIPROCESSING_WORKERS, then we use the number of files as the number of workers. Unless the number of machine cores is also less than MAX_NUM_MULTIPROCESSING_WORKERS, in which case we use the lower of the two.

Arguments

  • file_names: The list of file names to load.

Returns The number of workers to use for multiprocessing.

Methods​


add_hook​

def add_hook(self, hook: DataSourceHook) ‑> None:

Inherited from:

HeidelbergSource.add_hook :

Add a hook to the datasource.

add_strategy_filter_flag_columns​

def add_strategy_filter_flag_columns(    self, df: pd.DataFrame,) ‑> pandas.core.frame.DataFrame:

Inherited from:

HeidelbergSource.add_strategy_filter_flag_columns :

Add task-scoped strategy filter flag columns to a report DataFrame.

Decodes the compact per-file bitmask into one boolean column per flag-only strategy. Only files present in df are decoded, keeping memory proportional to the report size rather than the total number of evaluated files.

Arguments

  • df: DataFrame containing at least the ORIGINAL_FILENAME_METADATA_COLUMN column.

Returns A new DataFrame with one boolean column per flag-only strategy. If no flag-only strategies are registered, returns df unchanged (same object, no copy). If df does not contain the ORIGINAL_FILENAME_METADATA_COLUMN (e.g. patient-level reports), flag injection is skipped gracefully and df is returned unchanged.

apply_ignore_cols​

def apply_ignore_cols(self, df: pd.DataFrame) ‑> pandas.core.frame.DataFrame:

Inherited from:

HeidelbergSource.apply_ignore_cols :

Apply ignored columns to dataframe, dropping columns as needed.

Returns A copy of the dataframe with ignored columns removed, or the original dataframe if this datasource does not specify any ignore columns.

apply_ignore_cols_iter​

def apply_ignore_cols_iter(    self, dfs: Iterator[pd.DataFrame],) ‑> collections.abc.Iterator[pandas.core.frame.DataFrame]:

Inherited from:

HeidelbergSource.apply_ignore_cols_iter :

Apply ignored columns to dataframes from iterator.

apply_merged_filter_config​

def apply_merged_filter_config(self, merged_filter_config: MergedFilterConfig) ‑> None:

Inherited from:

HeidelbergSource.apply_merged_filter_config :

Apply the filter configuration to the datasource.

apply_modifiers​

def apply_modifiers(self, df: pd.DataFrame) ‑> pandas.core.frame.DataFrame:

Inherited from:

HeidelbergSource.apply_modifiers :

Apply column modifiers to the dataframe.

If no modifiers are specified, returns the dataframe unchanged.

clear_dataset_cache​

def clear_dataset_cache(self) ‑> dict[str, typing.Any]:

Inherited from:

HeidelbergSource.clear_dataset_cache :

Clear all dataset cache for this data source.

This clears both:

  1. The file names cache (Python cached_property)
  2. The dataset cache file (deletes the SQLite database file completely)

Returns Dictionary with cache clearing results.

clear_file_names_cache​

def clear_file_names_cache(self) ‑> None:

Inherited from:

HeidelbergSource.clear_file_names_cache :

Clears the list of selected file names.

This allows the datasource to pick up any new files that have been added to the directory since the last time it was cached.

clear_task_specific_configs​

def clear_task_specific_configs(self) ‑> None:

Inherited from:

HeidelbergSource.clear_task_specific_configs :

Clear task-scoped state at task boundaries.

Sends the pending skipped-file telemetry summary, then resets skipped files, merged filter config, and strategy filter flags so that subsequent tasks start with a clean slate.

file_names_iter​

def file_names_iter(    self, as_strs: bool = False,) ‑> collections.abc.Iterator[pathlib.Path] | collections.abc.Iterator[str]:

Inherited from:

HeidelbergSource.file_names_iter :

Iterate over files in a directory, yielding those that match the criteria.

The datasource-specific filters are applied a chunk of files at a time rather than one file at a time, so each filter reads what it needs for the whole chunk in one query instead of one per file — see _DATASOURCE_FILTER_BATCH_SIZE. Files are still yielded individually, in walk order, and a file dropped by any filter is bookkept exactly as before: the filters record their own specific skip reason as they go, and the generic DATASOURCE_FILTER_FAILED below remains the backstop for a file that none of them accounted for.

Arguments

  • as_strs: By default the files yielded will be yielded as Path objects. If this is True, yield them as strings instead.

files_passing_filters​

def files_passing_filters(    self, file_names: Iterable[str | os.PathLike[str]],) ‑> list[str]:

Inherited from:

HeidelbergSource.files_passing_filters :

Return which of file_names this datasource's filters admit.

The same decision the file walk makes as it yields, for a caller holding a list of paths from elsewhere — the DAG's file-metadata inventory, which is keyed by path and cannot answer a filter on a file's contents (frame count, modality, inferred scan type). Includes any task-level filters currently applied via apply_merged_filter_config.

Arguments

  • file_names: Candidate paths.

Returns Those that pass, in the order the underlying filters return them.

get_all_cached_file_paths​

def get_all_cached_file_paths(self) ‑> list[str]:

Inherited from:

HeidelbergSource.get_all_cached_file_paths :

Get all file paths that are currently stored in the cache.

Returns A list of file paths that have cache entries, or an empty list if there is no cache or the cache hasn't been initialized.

get_data​

def get_data(    self,    data_keys: SingleOrMulti[str] | SingleOrMulti[int],    *,    use_cache: bool = True,    **kwargs: Any,) ‑> pandas.core.frame.DataFrame | None:

Inherited from:

HeidelbergSource.get_data :

Get data corresponding to the provided data key(s).

Can be used to return data for a single data key or for multiple at once. If used for multiple, the order of the output dataframe must match the order of the keys provided.

Arguments

  • data_keys: Key(s) for which to get the data of. These may be things such as file names, UUIDs, etc. Can also be a list of integers if the datasource has an integer index.
  • use_cache: Whether the cache should be used to retrieve data for these keys. Note that cached data may have some elements, particularly image-related fields such as image data or file paths, replaced with placeholder values when stored in the cache. If data_cache is set on the instance, data will be set in the cache, regardless of this argument.
  • **kwargs: Additional keyword arguments.

Returns A dataframe containing the data, ordered to match the order of keys in data_keys, or None if no data for those keys was available.

get_datasource_metrics​

def get_datasource_metrics(    self, use_skip_codes: bool = False, data: pd.DataFrame | None = None,) ‑> DatasourceSummaryStats:

Inherited from:

HeidelbergSource.get_datasource_metrics :

Get metadata about this datasource.

This can be used to store information about the datasource that may be useful for debugging or tracking purposes. The metadata will be stored in the project database.

Arguments

  • use_skip_codes: Whether to use the skip reason codes as the keys in the skip_reasons dictionary, rather than the existing reason descriptions.
  • data: The data to use for getting the metrics.

Returns A dictionary containing metadata about this datasource.

get_filter_config​

def get_filter_config(self) ‑> FilterConfig:

Inherited from:

HeidelbergSource.get_filter_config :

Get the filter configuration for the datasource.

get_project_db_sqlite_columns​

def get_project_db_sqlite_columns(self) ‑> list[str]:

Inherited from:

HeidelbergSource.get_project_db_sqlite_columns :

Returns the required columns to identify a data point.

The first value must be filename column, and second value must be the last modified datetime. These two are used to build the processed_file_cache for the worker execution.

get_project_db_sqlite_create_table_query​

def get_project_db_sqlite_create_table_query(self) ‑> str:

Inherited from:

HeidelbergSource.get_project_db_sqlite_create_table_query :

Returns the required columns and types to identify a data point.

The file name is used as the primary key and the last modified date is used to determine if the file has been updated since the last time it was processed. If there is a conflict on the file name, the row is replaced with the new data to ensure that the last modified date is always up to date.

get_schema​

def get_schema(self) ‑> dict[str, typing.Any]:

Inherited from:

HeidelbergSource.get_schema :

Get the pre-defined schema for this datasource.

This method should be overridden by datasources that have pre-defined schemas (i.e., those with has_predefined_schema = True).

Returns The schema as a dictionary.

Raises

  • NotImplementedError: If the datasource doesn't have a pre-defined schema.

get_task_skip_reason_summary​

def get_task_skip_reason_summary(self) ‑> dict[str, int]:

Inherited from:

HeidelbergSource.get_task_skip_reason_summary :

Get aggregated skip reasons for the current task.

Combines both task-only skips (files that failed task filters) and datasource skips that occurred during the current task execution. This provides a complete picture of all files skipped during a task run.

Returns Dict mapping reason codes (as strings) to file counts.

get_uncached_file_names​

def get_uncached_file_names(self) ‑> list[str]:

Inherited from:

HeidelbergSource.get_uncached_file_names :

Return potentially uncached files via fast raw filesystem scanning.

This fast path skips datasource filters and computes file-path set difference against cache/skipped tables. It may include files that are later filtered out during normal datasource processing.

has_uncached_files​

def has_uncached_files(self) ‑> bool:

Inherited from:

HeidelbergSource.has_uncached_files :

Returns True if there are any files in the datasource not yet cached.

Uses a fast path that skips the full filter pipeline: walks the filesystem with scantree and checks each path against cache + skipped files metadata.

is_file_skipped_at_task_level​

def is_file_skipped_at_task_level(self, filename: str) ‑> bool:

Inherited from:

HeidelbergSource.is_file_skipped_at_task_level :

Whether filename carries a task-level skip from the current task.

A task-level skip marks a failure the datasource judged transient, so it is deliberately not persisted to the cache. Callers that record their own per-file outcome need this to tell "read the file, it holds nothing" from "could not read the file this time", which both surface as an empty result from _process_file.

Arguments

  • filename: Path as passed to the call that may have skipped it.

Returns True when the current task recorded a task-level skip for it.

log_parallel_capacity​

def log_parallel_capacity(    self, file_names: Sequence[str], chunk_size: int | None = None,) ‑> None:

Inherited from:

HeidelbergSource.log_parallel_capacity :

Log what parallelism this machine and configuration allow.

None of it is otherwise recorded, and all of it bounds what any change to how files are read could be worth: a site with two cores cannot gain from a wider pool however the code is arranged, and a site whose partitions hold one file never reaches the parallel path at all. Without these figures a timing report from a deployment cannot be read.

active and workers describe one partition, not the whole run, because that is the unit the parallel decision is actually made on. A caller that narrows the datasource before reading - the inference step hands it one batch_size chunk at a time - must say so via chunk_size, or the figures describe a partition no read ever uses and report parallelism that never happens.

Arguments

  • file_names: Every file about to be processed, across all chunks.
  • chunk_size: How many of them the caller narrows the datasource to at a time. Defaults to all of them, for a caller that does not narrow.

merge_and_validate_filters​

def merge_and_validate_filters(    self, datasource_level_filters: FilterConfig, task_level_filters: list[TaskFilter],) ‑> MergedFilterConfig:

Inherited from:

HeidelbergSource.merge_and_validate_filters :

Merge and validate the filters from the datasource and the task.

Returns a MergedFilterConfig with resolved filter values from both datasource and task-level filters, using intersection logic (most restrictive wins).

partition​

def partition(    self, iterable: Iterable[_I], partition_size: int = 1,) ‑> collections.abc.Iterable[collections.abc.Sequence[~_I]]:

Inherited from:

HeidelbergSource.partition :

Partition the iterable into chunks of the given size.

process_file​

def process_file(    self, filename: str, skip_non_tabular_data: bool = False, **kwargs: Any,) ‑> list[dict[str, typing.Any]]:

Inherited from:

HeidelbergSource.process_file :

Parse a single file into records, outside the batch-loading path.

The entry point for callers that need one file's records, such as the scan_metadata runtime. Unlike loading a batch it applies no data cache and adds no metadata columns; the only bookkeeping is what the datasource itself does while processing, such as recording a skip. Subclasses implement _process_file, not this.

Arguments

  • filename: The name of the file to process.
  • skip_non_tabular_data: Whether non-tabular data, e.g. image data, can be left unloaded.
  • **kwargs: Additional keyword arguments for _process_file.

Returns One dictionary per datapoint in the file, as _process_file returns.

remove_hook​

def remove_hook(self, hook: DataSourceHook) ‑> None:

Inherited from:

HeidelbergSource.remove_hook :

Remove a hook from the datasource.

selected_file_names_iter​

def selected_file_names_iter(self) ‑> collections.abc.Iterator[str]:

Inherited from:

HeidelbergSource.selected_file_names_iter :

Returns an iterator over selected file names.

Selected file names are affected by the selected_file_names_override and new_file_names_only attributes.

Returns Iterator over selected file names.

set_strategy_filter_flags​

def set_strategy_filter_flags(    self, flags: dict[str, int], meta: list[tuple[int, str, str | None]],) ‑> None:

Inherited from:

HeidelbergSource.set_strategy_filter_flags :

Store strategy filter flag-only match bitmasks for the current task.

Arguments

  • flags: Mapping of filename to bitmask of matched flag-only strategies.
  • meta: Ordered list of (original_index, strategy_name, flag_column_name). Bit position in the bitmask corresponds to list index.

skip_file​

def skip_file(    self, filename: str, reason: FileSkipReason, data: dict[str, Any] | None = None,) ‑> None:

Inherited from:

HeidelbergSource.skip_file :

Skip a file by updating cache and skipped_files set.

The first reason is always the one recorded in the data cache. If the file was already skipped, the stored reason is used for reporting, logging, and telemetry instead of the caller's reason.

Arguments

  • filename: Path to the file being skipped
  • reason: Reason for skipping the file
  • data: Unused. Kept for callers that pass file metadata.

task_skip_file​

def task_skip_file(    self, filename: str, reason: FileSkipReason, data: dict[str, Any] | None = None,) ‑> None:

Inherited from:

HeidelbergSource.task_skip_file :

Skip a file due to task-level filter (not persisted to cache).

Unlike skip_file(), this does NOT persist to cache because the file may pass a different task's filters. The skip is tracked in-memory only and cleared between tasks.

Arguments

  • filename: Path to the file being skipped.
  • reason: Reason for skipping the file.
  • data: Unused. Kept for callers that pass file metadata.

task_skip_reason​

def task_skip_reason(    self, filename: str,) ‑> FileSkipReason | None:

Inherited from:

HeidelbergSource.task_skip_reason :

Why filename carries a task-level skip, or None if it does not.

Arguments

  • filename: Path as passed to the call that may have skipped it.

use_file_multiprocessing​

def use_file_multiprocessing(self, file_names: Sequence[str]) ‑> bool:

Inherited from:

HeidelbergSource.use_file_multiprocessing :

Check if file multiprocessing should be used.

Returns True if file multiprocessing has been enabled by the environment variable and the number of workers would be greater than 1, otherwise False. There is no need to use file multiprocessing if we are just going to use one worker - it would be slower than just loading the data in the main process.

Returns True if file multiprocessing should be used, otherwise False.

yield_data​

def yield_data(    self,    data_keys: SingleOrMulti[str] | SingleOrMulti[int] | None = None,    *,    use_cache: bool = True,    partition_size: int | None = None,    **kwargs: Any,) ‑> collections.abc.Iterator[pandas.core.frame.DataFrame]:

Inherited from:

HeidelbergSource.yield_data :

Yields data in batches from this source.

If data_keys is specified, only yield from that subset of the data. Otherwise, iterate through the whole datasource.

Arguments

  • data_keys: An optional list of data keys to use for yielding data. Otherwise, all data in the datasource will be considered. data_keys is always provided when this method is called from the Dataset as part of a task. Can also be a list of integers if the datasource has an integer index.
  • use_cache: Whether the cache should be used to retrieve data for these data points. Note that cached data may have some elements, particularly image-related fields such as image data or file paths, replaced with placeholder values when stored in the cache. If data_cache is set on the instance, data will be set in the cache, regardless of this argument.
  • partition_size: The number of data elements to load/yield in each iteration. If not provided, defaults to the partition size configured in the datasource.
  • **kwargs: Additional keyword arguments.

HeidelbergSource​

class HeidelbergSource(    private_eye_parser: PrivateEyeParser | Mapping[str, PrivateEyeParser],    path: os.PathLike[str] | str,    parsers: PrivateEyeParser | Mapping[str, PrivateEyeParser] | None = None,    heidelberg_csv_columns: HeidelbergCSVColumns | _HeidelbergCSVColumnsTD | None = None,    required_fields: dict[str, Any] | None = None,    ophthalmology_args: OphthalmologyDataSourceArgs | _OphthalmologyDataSourceArgsTD | None = None,    data_cache: DataPersister | None = None,    infer_class_labels_from_filepaths: bool = False,    output_path: os.PathLike[str] | str | None = None,    iterable: bool = True,    fast_load: bool = True,    cache_images: bool = False,    filter: FileSystemFilter | None = None,    data_splitter: DatasetSplitter | None = None,    seed: int | None = None,    ignore_cols: str | Sequence[str] | None = None,    modifiers: dict[str, DataPathModifiers] | None = None,    partition_size: int = 16,    name: str | None = None,):

Data source for loading Heidelberg .sdb files.

Arguments

  • **kwargs: Keyword arguments passed to the parent base classes.
  • cache_images: Whether to cache images in the file system. Defaults to False. This is ignored if fast_load is True.
  • data_cache: A DataPersister instance to use for data caching.
  • data_splitter: Deprecated argument, will be removed in a future release. Defaults to None. Not used.
  • fast_load: Whether the data will be loaded in fast mode. This is used to determine whether the data will be iterated over during set up for schema generation and splitting (where necessary). Only relevant if iterable is True, otherwise it is ignored. Defaults to True.
  • heidelberg_csv_columns: If path is a CSV file, this contains information about the specific columns that contain the path information for the Heidelberg files. If not provided, it is assumed that the CSV file contains a column named 'heidelberg_file' that contains the paths to the Heidelberg files. Defaults to None.
  • ignore_cols: Column/list of columns to be ignored from the data. Defaults to None.
  • infer_class_labels_from_filepaths: Whether class labels should be added to the data based on the filepath of the files. Defaults to the first directory within self.path, but can go a level deeper if the datasplitter is provided with infer_data_split_labels set to true
  • iterable: Whether the data source is iterable. This is used to determine whether the data source can be used in a streaming context during a task. Defaults to True.
  • modifiers: Dictionary used for modifying paths/ extensions in the dataframe. Defaults to None.
  • name: The name for the datasource. Optional, defaults to None.
  • ophthalmology_args: Arguments for ophthalmology modality data.
  • output_path: The path where to save intermediary output files. Defaults to 'preprocessed/'.
  • parsers: The private eye parsers to use for the different file extensions. Only needs to be supplied if file_extension filter is non-default. Can either be a single parser to use for all file extensions or a mapping of file extensions to parser type. Defaults to appropriate parser(s) for the default file extension(s).
  • partition_size: The size of each partition when iterating over the data in a batched fashion.
  • path: The path to the directory containing the Heidelberg .sdb files or to a CSV file that includes Heidelberg files as one of its columns. If a CSV file is provided, the file extensions specified in file_extension will be ignored. If a CSV file is provided, the heidelberg_csv_columns argument should also be provided.
  • private_eye_parser: Private-eye supported machine type(s). Can either be a single parser to use for all files or a mapping of file extension to the desired parser type. If private_eye_parser is a mapping of file extensions to parsers, there must be a parser for each file extension specified. If no file extensions are specified then mapping can exist in whatever form (warnings will be logged if we encounter an extension for which no parser is specified). If private_eye_parser is a single parser, file_extension can be anything, we simply try to use this parser against any extension.
  • seed: Random number seed. Used for setting random seed for all libraries. Defaults to None.

Attributes

  • seed: Random number seed. Used for setting random seed for all libraries.

Raises

  • ValueError: If the minimum DOB is greater than the maximum DOB.
  • ValueError: If the minimum number of B-scans is greater than the maximum number of B-scans.
  • ValueError: If the minimum acquisition date is greater than the maximum acquisition date.

Variables​

  • accessibility_details : AccessibilityDetails | None - Detailed accessibility status. None if accessible.

    Subclasses should override to perform lightweight connectivity check.

    Returns: None if accessible, or dict with 'error_code' and 'message' if not.

  • file_names : list[str] - Returns a list of file names in the specified directory.

    warning

Deprecated: The file_names property is deprecated and will be removed in a future release. Use file_names_iter(as_strs=True) for memory-efficient iteration, or list(file_names_iter(as_strs=True)) if you need a list. :::

This property accounts for files skipped at runtime by filtering them out of the list of cached file names. Files may get skipped at runtime due to errors or because they don't contain any image data and images_only is True. This allows us to skip these files again more quickly if they are still present in the directory.

  • is_accessible : bool - Check if datasource is currently accessible.

    Returns True if accessibility_details is None (no errors). This is a convenience property that wraps accessibility_details.

  • is_file_iterable : bool - Returns True since this source iterates over files.
  • is_initialised : bool - Checks if BaseSource was initialised.
  • is_task_running : bool - Returns True if a task is running.
  • name : str | None - The datasource's configured name.
  • path : pathlib.Path - Resolved absolute path to data.

    Provides a consistent version of the path provided by the user which should work throughout regardless of operating system and of directory structure.

  • selected_file_names : list[str] - Returns a list of selected file names as strings.

    Selected file names are affected by the selected_file_names_override and new_file_names_only attributes.

    WARNING: This method loads all filenames into memory. For large datasets, consider using selected_file_names_iter() instead.

  • selected_file_names_differ : bool - Returns True if selected_file_names will differ from default.

    In particular, returns True iff there is a selected file names override in place and/or there is filtering for new file names only present.

  • supports_project_db : bool - Whether the datasource supports the project database.

    Each datasource needs to implement its own methods to define how what its project database table should look like. If the datasource does not implement the methods to get the table creation query and columns, it does not support the projectdatabase.

  • task_skipped_file_names : set[str] - Return set of task-skipped filenames for set operations.

Static methods​


get_num_workers​

def get_num_workers(file_names: Sequence[str]) ‑> int:

Inherited from:

FileSystemIterableSourceInferrable.get_num_workers :

Gets the number of workers to use for multiprocessing.

Ensures that the number of workers is at least 1 and at most equal to MAX_NUM_MULTIPROCESSING_WORKERS. If the number of files is less than MAX_NUM_MULTIPROCESSING_WORKERS, then we use the number of files as the number of workers. Unless the number of machine cores is also less than MAX_NUM_MULTIPROCESSING_WORKERS, in which case we use the lower of the two.

Arguments

  • file_names: The list of file names to load.

Returns The number of workers to use for multiprocessing.

Methods​


add_hook​

def add_hook(self, hook: DataSourceHook) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.add_hook :

Add a hook to the datasource.

add_strategy_filter_flag_columns​

def add_strategy_filter_flag_columns(    self, df: pd.DataFrame,) ‑> pandas.core.frame.DataFrame:

Inherited from:

FileSystemIterableSourceInferrable.add_strategy_filter_flag_columns :

Add task-scoped strategy filter flag columns to a report DataFrame.

Decodes the compact per-file bitmask into one boolean column per flag-only strategy. Only files present in df are decoded, keeping memory proportional to the report size rather than the total number of evaluated files.

Arguments

  • df: DataFrame containing at least the ORIGINAL_FILENAME_METADATA_COLUMN column.

Returns A new DataFrame with one boolean column per flag-only strategy. If no flag-only strategies are registered, returns df unchanged (same object, no copy). If df does not contain the ORIGINAL_FILENAME_METADATA_COLUMN (e.g. patient-level reports), flag injection is skipped gracefully and df is returned unchanged.

apply_ignore_cols​

def apply_ignore_cols(self, df: pd.DataFrame) ‑> pandas.core.frame.DataFrame:

Inherited from:

FileSystemIterableSourceInferrable.apply_ignore_cols :

Apply ignored columns to dataframe, dropping columns as needed.

Returns A copy of the dataframe with ignored columns removed, or the original dataframe if this datasource does not specify any ignore columns.

apply_ignore_cols_iter​

def apply_ignore_cols_iter(    self, dfs: Iterator[pd.DataFrame],) ‑> collections.abc.Iterator[pandas.core.frame.DataFrame]:

Inherited from:

FileSystemIterableSourceInferrable.apply_ignore_cols_iter :

Apply ignored columns to dataframes from iterator.

apply_merged_filter_config​

def apply_merged_filter_config(self, merged_filter_config: MergedFilterConfig) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.apply_merged_filter_config :

Apply the filter configuration to the datasource.

apply_modifiers​

def apply_modifiers(self, df: pd.DataFrame) ‑> pandas.core.frame.DataFrame:

Inherited from:

FileSystemIterableSourceInferrable.apply_modifiers :

Apply column modifiers to the dataframe.

If no modifiers are specified, returns the dataframe unchanged.

clear_dataset_cache​

def clear_dataset_cache(self) ‑> dict[str, typing.Any]:

Inherited from:

FileSystemIterableSourceInferrable.clear_dataset_cache :

Clear all dataset cache for this data source.

This clears both:

  1. The file names cache (Python cached_property)
  2. The dataset cache file (deletes the SQLite database file completely)

Returns Dictionary with cache clearing results.

clear_file_names_cache​

def clear_file_names_cache(self) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.clear_file_names_cache :

Clears the list of selected file names.

This allows the datasource to pick up any new files that have been added to the directory since the last time it was cached.

clear_task_specific_configs​

def clear_task_specific_configs(self) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.clear_task_specific_configs :

Clear task-scoped state at task boundaries.

Sends the pending skipped-file telemetry summary, then resets skipped files, merged filter config, and strategy filter flags so that subsequent tasks start with a clean slate.

file_names_iter​

def file_names_iter(    self, as_strs: bool = False,) ‑> collections.abc.Iterator[pathlib.Path] | collections.abc.Iterator[str]:

Inherited from:

FileSystemIterableSourceInferrable.file_names_iter :

Iterate over files in a directory, yielding those that match the criteria.

The datasource-specific filters are applied a chunk of files at a time rather than one file at a time, so each filter reads what it needs for the whole chunk in one query instead of one per file — see _DATASOURCE_FILTER_BATCH_SIZE. Files are still yielded individually, in walk order, and a file dropped by any filter is bookkept exactly as before: the filters record their own specific skip reason as they go, and the generic DATASOURCE_FILTER_FAILED below remains the backstop for a file that none of them accounted for.

Arguments

  • as_strs: By default the files yielded will be yielded as Path objects. If this is True, yield them as strings instead.

files_passing_filters​

def files_passing_filters(    self, file_names: Iterable[str | os.PathLike[str]],) ‑> list[str]:

Inherited from:

FileSystemIterableSourceInferrable.files_passing_filters :

Return which of file_names this datasource's filters admit.

The same decision the file walk makes as it yields, for a caller holding a list of paths from elsewhere — the DAG's file-metadata inventory, which is keyed by path and cannot answer a filter on a file's contents (frame count, modality, inferred scan type). Includes any task-level filters currently applied via apply_merged_filter_config.

Arguments

  • file_names: Candidate paths.

Returns Those that pass, in the order the underlying filters return them.

get_all_cached_file_paths​

def get_all_cached_file_paths(self) ‑> list[str]:

Inherited from:

FileSystemIterableSourceInferrable.get_all_cached_file_paths :

Get all file paths that are currently stored in the cache.

Returns A list of file paths that have cache entries, or an empty list if there is no cache or the cache hasn't been initialized.

get_data​

def get_data(    self,    data_keys: SingleOrMulti[str] | SingleOrMulti[int],    *,    use_cache: bool = True,    **kwargs: Any,) ‑> pandas.core.frame.DataFrame | None:

Inherited from:

FileSystemIterableSourceInferrable.get_data :

Get data corresponding to the provided data key(s).

Can be used to return data for a single data key or for multiple at once. If used for multiple, the order of the output dataframe must match the order of the keys provided.

Arguments

  • data_keys: Key(s) for which to get the data of. These may be things such as file names, UUIDs, etc. Can also be a list of integers if the datasource has an integer index.
  • use_cache: Whether the cache should be used to retrieve data for these keys. Note that cached data may have some elements, particularly image-related fields such as image data or file paths, replaced with placeholder values when stored in the cache. If data_cache is set on the instance, data will be set in the cache, regardless of this argument.
  • **kwargs: Additional keyword arguments.

Returns A dataframe containing the data, ordered to match the order of keys in data_keys, or None if no data for those keys was available.

get_datasource_metrics​

def get_datasource_metrics(    self, use_skip_codes: bool = False, data: pd.DataFrame | None = None,) ‑> DatasourceSummaryStats:

Inherited from:

FileSystemIterableSourceInferrable.get_datasource_metrics :

Get metadata about this datasource.

This can be used to store information about the datasource that may be useful for debugging or tracking purposes. The metadata will be stored in the project database.

Arguments

  • use_skip_codes: Whether to use the skip reason codes as the keys in the skip_reasons dictionary, rather than the existing reason descriptions.
  • data: The data to use for getting the metrics.

Returns A dictionary containing metadata about this datasource.

get_filter_config​

def get_filter_config(self) ‑> FilterConfig:

Inherited from:

FileSystemIterableSourceInferrable.get_filter_config :

Get the filter configuration for the datasource.

get_project_db_sqlite_columns​

def get_project_db_sqlite_columns(self) ‑> list[str]:

Inherited from:

FileSystemIterableSourceInferrable.get_project_db_sqlite_columns :

Returns the required columns to identify a data point.

The first value must be filename column, and second value must be the last modified datetime. These two are used to build the processed_file_cache for the worker execution.

get_project_db_sqlite_create_table_query​

def get_project_db_sqlite_create_table_query(self) ‑> str:

Inherited from:

FileSystemIterableSourceInferrable.get_project_db_sqlite_create_table_query :

Returns the required columns and types to identify a data point.

The file name is used as the primary key and the last modified date is used to determine if the file has been updated since the last time it was processed. If there is a conflict on the file name, the row is replaced with the new data to ensure that the last modified date is always up to date.

get_schema​

def get_schema(self) ‑> dict[str, typing.Any]:

Inherited from:

FileSystemIterableSourceInferrable.get_schema :

Get the pre-defined schema for this datasource.

This method should be overridden by datasources that have pre-defined schemas (i.e., those with has_predefined_schema = True).

Returns The schema as a dictionary.

Raises

  • NotImplementedError: If the datasource doesn't have a pre-defined schema.

get_task_skip_reason_summary​

def get_task_skip_reason_summary(self) ‑> dict[str, int]:

Inherited from:

FileSystemIterableSourceInferrable.get_task_skip_reason_summary :

Get aggregated skip reasons for the current task.

Combines both task-only skips (files that failed task filters) and datasource skips that occurred during the current task execution. This provides a complete picture of all files skipped during a task run.

Returns Dict mapping reason codes (as strings) to file counts.

get_uncached_file_names​

def get_uncached_file_names(self) ‑> list[str]:

Inherited from:

FileSystemIterableSourceInferrable.get_uncached_file_names :

Return potentially uncached files via fast raw filesystem scanning.

This fast path skips datasource filters and computes file-path set difference against cache/skipped tables. It may include files that are later filtered out during normal datasource processing.

has_uncached_files​

def has_uncached_files(self) ‑> bool:

Inherited from:

FileSystemIterableSourceInferrable.has_uncached_files :

Returns True if there are any files in the datasource not yet cached.

Uses a fast path that skips the full filter pipeline: walks the filesystem with scantree and checks each path against cache + skipped files metadata.

is_file_skipped_at_task_level​

def is_file_skipped_at_task_level(self, filename: str) ‑> bool:

Inherited from:

FileSystemIterableSourceInferrable.is_file_skipped_at_task_level :

Whether filename carries a task-level skip from the current task.

A task-level skip marks a failure the datasource judged transient, so it is deliberately not persisted to the cache. Callers that record their own per-file outcome need this to tell "read the file, it holds nothing" from "could not read the file this time", which both surface as an empty result from _process_file.

Arguments

  • filename: Path as passed to the call that may have skipped it.

Returns True when the current task recorded a task-level skip for it.

log_parallel_capacity​

def log_parallel_capacity(    self, file_names: Sequence[str], chunk_size: int | None = None,) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.log_parallel_capacity :

Log what parallelism this machine and configuration allow.

None of it is otherwise recorded, and all of it bounds what any change to how files are read could be worth: a site with two cores cannot gain from a wider pool however the code is arranged, and a site whose partitions hold one file never reaches the parallel path at all. Without these figures a timing report from a deployment cannot be read.

active and workers describe one partition, not the whole run, because that is the unit the parallel decision is actually made on. A caller that narrows the datasource before reading - the inference step hands it one batch_size chunk at a time - must say so via chunk_size, or the figures describe a partition no read ever uses and report parallelism that never happens.

Arguments

  • file_names: Every file about to be processed, across all chunks.
  • chunk_size: How many of them the caller narrows the datasource to at a time. Defaults to all of them, for a caller that does not narrow.

merge_and_validate_filters​

def merge_and_validate_filters(    self, datasource_level_filters: FilterConfig, task_level_filters: list[TaskFilter],) ‑> MergedFilterConfig:

Inherited from:

FileSystemIterableSourceInferrable.merge_and_validate_filters :

Merge and validate the filters from the datasource and the task.

Returns a MergedFilterConfig with resolved filter values from both datasource and task-level filters, using intersection logic (most restrictive wins).

partition​

def partition(    self, iterable: Iterable[_I], partition_size: int = 1,) ‑> collections.abc.Iterable[collections.abc.Sequence[~_I]]:

Inherited from:

FileSystemIterableSourceInferrable.partition :

Partition the iterable into chunks of the given size.

process_file​

def process_file(    self, filename: str, skip_non_tabular_data: bool = False, **kwargs: Any,) ‑> list[dict[str, typing.Any]]:

Inherited from:

FileSystemIterableSourceInferrable.process_file :

Parse a single file into records, outside the batch-loading path.

The entry point for callers that need one file's records, such as the scan_metadata runtime. Unlike loading a batch it applies no data cache and adds no metadata columns; the only bookkeeping is what the datasource itself does while processing, such as recording a skip. Subclasses implement _process_file, not this.

Arguments

  • filename: The name of the file to process.
  • skip_non_tabular_data: Whether non-tabular data, e.g. image data, can be left unloaded.
  • **kwargs: Additional keyword arguments for _process_file.

Returns One dictionary per datapoint in the file, as _process_file returns.

remove_hook​

def remove_hook(self, hook: DataSourceHook) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.remove_hook :

Remove a hook from the datasource.

selected_file_names_iter​

def selected_file_names_iter(self) ‑> collections.abc.Iterator[str]:

Inherited from:

FileSystemIterableSourceInferrable.selected_file_names_iter :

Returns an iterator over selected file names.

Selected file names are affected by the selected_file_names_override and new_file_names_only attributes.

Returns Iterator over selected file names.

set_strategy_filter_flags​

def set_strategy_filter_flags(    self, flags: dict[str, int], meta: list[tuple[int, str, str | None]],) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.set_strategy_filter_flags :

Store strategy filter flag-only match bitmasks for the current task.

Arguments

  • flags: Mapping of filename to bitmask of matched flag-only strategies.
  • meta: Ordered list of (original_index, strategy_name, flag_column_name). Bit position in the bitmask corresponds to list index.

skip_file​

def skip_file(    self, filename: str, reason: FileSkipReason, data: dict[str, Any] | None = None,) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.skip_file :

Skip a file by updating cache and skipped_files set.

The first reason is always the one recorded in the data cache. If the file was already skipped, the stored reason is used for reporting, logging, and telemetry instead of the caller's reason.

Arguments

  • filename: Path to the file being skipped
  • reason: Reason for skipping the file
  • data: Unused. Kept for callers that pass file metadata.

task_skip_file​

def task_skip_file(    self, filename: str, reason: FileSkipReason, data: dict[str, Any] | None = None,) ‑> None:

Inherited from:

FileSystemIterableSourceInferrable.task_skip_file :

Skip a file due to task-level filter (not persisted to cache).

Unlike skip_file(), this does NOT persist to cache because the file may pass a different task's filters. The skip is tracked in-memory only and cleared between tasks.

Arguments

  • filename: Path to the file being skipped.
  • reason: Reason for skipping the file.
  • data: Unused. Kept for callers that pass file metadata.

task_skip_reason​

def task_skip_reason(    self, filename: str,) ‑> FileSkipReason | None:

Inherited from:

FileSystemIterableSourceInferrable.task_skip_reason :

Why filename carries a task-level skip, or None if it does not.

Arguments

  • filename: Path as passed to the call that may have skipped it.

use_file_multiprocessing​

def use_file_multiprocessing(self, file_names: Sequence[str]) ‑> bool:

Inherited from:

FileSystemIterableSourceInferrable.use_file_multiprocessing :

Check if file multiprocessing should be used.

Returns True if file multiprocessing has been enabled by the environment variable and the number of workers would be greater than 1, otherwise False. There is no need to use file multiprocessing if we are just going to use one worker - it would be slower than just loading the data in the main process.

Returns True if file multiprocessing should be used, otherwise False.

yield_data​

def yield_data(    self,    data_keys: SingleOrMulti[str] | SingleOrMulti[int] | None = None,    *,    use_cache: bool = True,    partition_size: int | None = None,    **kwargs: Any,) ‑> collections.abc.Iterator[pandas.core.frame.DataFrame]:

Inherited from:

FileSystemIterableSourceInferrable.yield_data :

Yields data in batches from this source.

If data_keys is specified, only yield from that subset of the data. Otherwise, iterate through the whole datasource.

Arguments

  • data_keys: An optional list of data keys to use for yielding data. Otherwise, all data in the datasource will be considered. data_keys is always provided when this method is called from the Dataset as part of a task. Can also be a list of integers if the datasource has an integer index.
  • use_cache: Whether the cache should be used to retrieve data for these data points. Note that cached data may have some elements, particularly image-related fields such as image data or file paths, replaced with placeholder values when stored in the cache. If data_cache is set on the instance, data will be set in the cache, regardless of this argument.
  • partition_size: The number of data elements to load/yield in each iteration. If not provided, defaults to the partition size configured in the datasource.
  • **kwargs: Additional keyword arguments.