@@ -70,14 +70,64 @@ def log_summary(self, prefix: str = "") -> None:
7070 " (created/updated/removed/unchanged)"
7171 )
7272
73+ def __iadd__ (self , other : "IngestionStats" ) -> "IngestionStats" :
74+ """Accumulate counts in place from another :class:`IngestionStats`."""
75+ self .datasets_created += other .datasets_created
76+ self .datasets_updated += other .datasets_updated
77+ self .datasets_unchanged += other .datasets_unchanged
78+ self .files_added += other .files_added
79+ self .files_updated += other .files_updated
80+ self .files_removed += other .files_removed
81+ self .files_unchanged += other .files_unchanged
82+ return self
83+
84+
85+ def _ingest_catalog (
86+ adapter : DatasetAdapter ,
87+ db : Database ,
88+ data_catalog : pd .DataFrame ,
89+ ) -> IngestionStats :
90+ """
91+ Register every dataset in ``data_catalog``, committing per-dataset.
92+
93+ The ORM session identity map is expired after each commit so memory
94+ use stays bounded by the largest single dataset, not by the size of
95+ ``data_catalog``.
96+ """
97+ stats = IngestionStats ()
98+
99+ for instance_id , data_catalog_dataset in data_catalog .groupby (adapter .slug_column ):
100+ logger .debug (f"Processing dataset { instance_id } " )
101+ with db .session .begin ():
102+ results = adapter .register_dataset (db , data_catalog_dataset )
103+
104+ if results .dataset_state == ModelState .CREATED :
105+ stats .datasets_created += 1
106+ elif results .dataset_state == ModelState .UPDATED :
107+ stats .datasets_updated += 1
108+ else :
109+ stats .datasets_unchanged += 1
110+ stats .files_added += len (results .files_added )
111+ stats .files_updated += len (results .files_updated )
112+ stats .files_removed += len (results .files_removed )
113+ stats .files_unchanged += len (results .files_unchanged )
114+
115+ # Release ORM objects from the session identity map after each commit.
116+ # Without this, all Dataset and DatasetFile objects accumulate in memory
117+ # across the entire ingestion loop.
118+ db .session .expire_all ()
119+
120+ return stats
73121
74- def ingest_datasets (
122+
123+ def ingest_datasets ( # noqa: PLR0913
75124 adapter : DatasetAdapter ,
76125 directory : Path | None ,
77126 db : Database ,
78127 * ,
79128 data_catalog : pd .DataFrame | None = None ,
80129 skip_invalid : bool = True ,
130+ chunk_size : int | None = None ,
81131) -> IngestionStats :
82132 """
83133 Ingest datasets from a directory into the database.
@@ -94,10 +144,18 @@ def ingest_datasets(
94144 db
95145 Database instance
96146 data_catalog
97- Optional pre-validated data catalog. If provided, directory is ignored and
98- the catalog is used directly. This avoids redundant find/validate operations.
147+ Optional pre-validated data catalog.
148+
149+ If provided, directory is ignored and the catalog is used directly.
150+ This avoids redundant find/validate operations.
151+ When supplied, ``chunk_size`` is ignored because the catalog is already fully materialised.
99152 skip_invalid
100153 If True, skip datasets that fail validation (default True)
154+ chunk_size
155+ When provided and ``data_catalog`` is None,
156+ stream the directory in batches of ``chunk_size`` files so peak memory is bounded regardless
157+ of how many files live under ``directory``.
158+ Requires the adapter to implement ``iter_local_datasets``.
101159
102160 Returns
103161 -------
@@ -109,51 +167,62 @@ def ingest_datasets(
109167 ValueError
110168 If no valid datasets are found in the directory
111169 """
112- if data_catalog is None :
113- if directory is None :
114- raise ValueError ("Either directory or data_catalog must be provided" )
115-
116- if not directory .exists ():
117- raise ValueError (f"Directory { directory } does not exist" )
118-
119- # Check for .nc files
120- if not list (directory .rglob ("*.nc" )):
121- raise ValueError (f"No .nc files found in { directory } " )
122-
123- data_catalog = adapter .find_local_datasets (directory )
124- data_catalog = adapter .validate_data_catalog (data_catalog , skip_invalid = skip_invalid )
125-
126- if data_catalog .empty :
170+ if data_catalog is not None :
171+ return _ingest_catalog (adapter , db , data_catalog )
172+
173+ if directory is None :
174+ raise ValueError ("Either directory or data_catalog must be provided" )
175+
176+ if not directory .exists ():
177+ raise ValueError (f"Directory { directory } does not exist" )
178+
179+ # Check for .nc files
180+ if not any (directory .rglob ("*.nc" )):
181+ raise ValueError (f"No .nc files found in { directory } " )
182+
183+ if chunk_size is not None :
184+ if chunk_size < 1 :
185+ raise ValueError (f"chunk_size must be >= 1, got { chunk_size } " )
186+ iter_fn = getattr (adapter , "iter_local_datasets" , None )
187+ if iter_fn is None :
188+ raise ValueError (
189+ f"Adapter { type (adapter ).__name__ } does not support streaming ingest "
190+ "(missing iter_local_datasets); omit chunk_size to use whole-catalog mode."
191+ )
192+
193+ stats = IngestionStats ()
194+ total_files = 0
195+ total_datasets = 0
196+ emitted = False
197+ for raw_chunk in iter_fn (directory , chunk_size = chunk_size ):
198+ validated_chunk = adapter .validate_data_catalog (raw_chunk , skip_invalid = skip_invalid )
199+ if validated_chunk .empty :
200+ continue
201+ emitted = True
202+ total_files += len (validated_chunk )
203+ total_datasets += validated_chunk [adapter .slug_column ].nunique ()
204+ stats += _ingest_catalog (adapter , db , validated_chunk )
205+ # Drop chunk references so the per-chunk pandas memory can be
206+ # reclaimed before the next chunk is parsed.
207+ del raw_chunk , validated_chunk
208+
209+ if not emitted :
127210 raise ValueError (f"No valid datasets found in { directory } " )
128211
129- logger .info (
130- f"Found { len (data_catalog )} files for { len (data_catalog [adapter .slug_column ].unique ())} datasets"
131- )
212+ logger .info (f"Ingested { total_files } files across approximately { total_datasets } datasets (streamed)" )
213+ return stats
132214
133- stats = IngestionStats ()
215+ data_catalog = adapter .find_local_datasets (directory )
216+ data_catalog = adapter .validate_data_catalog (data_catalog , skip_invalid = skip_invalid )
134217
135- for instance_id , data_catalog_dataset in data_catalog .groupby (adapter .slug_column ):
136- logger .debug (f"Processing dataset { instance_id } " )
137- with db .session .begin ():
138- results = adapter .register_dataset (db , data_catalog_dataset )
218+ if data_catalog .empty :
219+ raise ValueError (f"No valid datasets found in { directory } " )
139220
140- if results .dataset_state == ModelState .CREATED :
141- stats .datasets_created += 1
142- elif results .dataset_state == ModelState .UPDATED :
143- stats .datasets_updated += 1
144- else :
145- stats .datasets_unchanged += 1
146- stats .files_added += len (results .files_added )
147- stats .files_updated += len (results .files_updated )
148- stats .files_removed += len (results .files_removed )
149- stats .files_unchanged += len (results .files_unchanged )
221+ logger .info (
222+ f"Found { len (data_catalog )} files for { len (data_catalog [adapter .slug_column ].unique ())} datasets"
223+ )
150224
151- # Release ORM objects from the session identity map after each commit.
152- # Without this, all Dataset and DatasetFile objects accumulate in memory
153- # across the entire ingestion loop.
154- db .session .expire_all ()
155-
156- return stats
225+ return _ingest_catalog (adapter , db , data_catalog )
157226
158227
159228def get_dataset_adapter (source_type : str , ** kwargs : Any ) -> DatasetAdapter :
0 commit comments