Skip to content

Catalogue streamer

cutana.catalogue_streamer

Streaming catalogue loading infrastructure for Cutana.

Provides memory-efficient catalogue loading for large catalogues (10M+ sources) through a two-phase approach: 1. Index building: Stream through catalogue once building FITS-set-to-row-indices mapping 2. Batch reading: Read only specific rows on-demand using pyarrow's take()

This module maintains FITS set optimization (80-90% I/O reduction) while enabling O(index_size) + O(batch_size) memory usage instead of O(catalogue_size).

CatalogueIndex(fits_set_to_row_indices=dict(), row_count=0, catalogue_path='') dataclass

Lightweight index for streaming catalogue access.

Stores mapping from FITS file sets to row indices, enabling FITS-set-optimized batch creation without loading the full catalogue.

Memory usage: ~100 bytes per source (row_idx + fits_set reference) For 10M sources: ~1GB index vs ~100GB for full catalogue

build_from_path(path, batch_size=100000) classmethod

Build index by streaming through catalogue.

Extracts only (row_index, fits_file_paths) to minimize memory usage.

Parameters:

Name Type Description Default
path str

Path to catalogue file (CSV or Parquet)

required
batch_size int

Number of rows to process per chunk

100000

Returns:

Type Description
CatalogueIndex

CatalogueIndex with FITS set to row indices mapping

get_optimized_batch_ranges(max_sources_per_batch, min_sources_per_batch=500, max_fits_sets_per_batch=50)

Generate optimized batch row ranges grouped by FITS sets.

Groups sources by FITS file sets to maximize I/O efficiency using a greedy algorithm. The goal is to minimize FITS file loading by keeping sources using the same FITS files together.

Batching Strategy - Atomic Tile Sets: - A tile (FITS set) is NEVER split across batches UNLESS it individually exceeds max_sources_per_batch. - Multiple tiles CAN be combined into the same batch if their combined total is <= max_sources_per_batch. - Only tiles that exceed max_sources_per_batch are split into multiple consecutive batches.

This ensures that output files/folders contain complete tile sets, making downstream processing and organization more predictable.

Algorithm: 1. Sort FITS sets by size (largest first) 2. For sets > max_sources_per_batch: Split into max-sized chunks (unavoidable) 3. For sets <= max_sources_per_batch: Keep atomic, combine with others if room

Parameters:

Name Type Description Default
max_sources_per_batch int

Maximum sources per batch (hard limit for memory)

required
min_sources_per_batch int

Minimum sources before flushing (efficiency threshold)

500
max_fits_sets_per_batch int

Maximum FITS sets per batch (limits I/O complexity)

50

Returns:

Type Description
List[List[int]]

List of row index lists, each representing a batch

get_fits_set_statistics()

Return statistics about FITS set distribution.

CatalogueBatchReader(path)

Reads specific row ranges from catalogues efficiently.

For parquet: Uses pyarrow's take() for O(1) random access For CSV: Uses chunked reading with filtering (slower, recommends parquet)

Initialize batch reader.

Parameters:

Name Type Description Default
path str

Path to catalogue file

required

read_rows(row_indices)

Read specific rows from catalogue.

Parameters:

Name Type Description Default
row_indices List[int]

List of row indices to read (0-based)

required

Returns:

Type Description
DataFrame

DataFrame containing only the requested rows

close()

Release resources.

estimate_catalogue_size(path)

Estimate number of rows in catalogue without loading it fully.

For parquet files, returns exact row count from metadata. For CSV files, estimates by sampling first 10 rows to calculate average line size.

Parameters:

Name Type Description Default
path str

Path to catalogue file

required

Returns:

Type Description
int

Row count (exact for parquet, estimated for CSV)

Raises:

Type Description
ValueError

If file format is unsupported or CSV cannot be sampled