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