Coverage for python/lsst/daf/butler/direct_butler/_direct_butler.py: 88%
974 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-02 05:09 -0400
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-02 05:09 -0400
1# This file is part of daf_butler.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (http://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This software is dual licensed under the GNU General Public License and also
10# under a 3-clause BSD license. Recipients may choose which of these licenses
11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt,
12# respectively. If you choose the GPL option then the following text applies
13# (but note that there is still no warranty even if you opt for BSD instead):
14#
15# This program is free software: you can redistribute it and/or modify
16# it under the terms of the GNU General Public License as published by
17# the Free Software Foundation, either version 3 of the License, or
18# (at your option) any later version.
19#
20# This program is distributed in the hope that it will be useful,
21# but WITHOUT ANY WARRANTY; without even the implied warranty of
22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
23# GNU General Public License for more details.
24#
25# You should have received a copy of the GNU General Public License
26# along with this program. If not, see <http://www.gnu.org/licenses/>.
28"""Butler top level classes."""
30from __future__ import annotations
32__all__ = (
33 "ButlerValidationError",
34 "DirectButler",
35)
37import collections.abc
38import contextlib
39import io
40import itertools
41import math
42import numbers
43import os
44import uuid
45import warnings
46from collections import Counter, defaultdict
47from collections.abc import Collection, Iterable, Iterator, Mapping, MutableMapping, Sequence
48from functools import partial
49from types import EllipsisType
50from typing import TYPE_CHECKING, Any, ClassVar, NamedTuple, TextIO, cast
52from deprecated.sphinx import deprecated
53from sqlalchemy.exc import IntegrityError
55from lsst.resources import ResourcePath, ResourcePathExpression
56from lsst.utils.introspection import find_outside_stacklevel, get_class_of
57from lsst.utils.iteration import chunk_iterable
58from lsst.utils.logging import VERBOSE, getLogger
59from lsst.utils.timer import time_this
61from .._butler import Butler, _DeprecatedDefault
62from .._butler_config import ButlerConfig
63from .._butler_instance_options import ButlerInstanceOptions
64from .._butler_metrics import ButlerMetrics
65from .._collection_type import CollectionType
66from .._dataset_existence import DatasetExistence
67from .._dataset_ref import DatasetRef
68from .._dataset_type import DatasetType
69from .._deferredDatasetHandle import DeferredDatasetHandle
70from .._exceptions import (
71 DatasetNotFoundError,
72 DimensionValueError,
73 EmptyQueryResultError,
74 ValidationError,
75)
76from .._file_dataset import FileDataset
77from .._limited_butler import LimitedButler
78from .._query_all_datasets import QueryAllDatasetsParameters, query_all_datasets
79from .._registry_shim import RegistryShim
80from .._storage_class import StorageClass, StorageClassFactory
81from .._timespan import Timespan
82from ..datastore import Datastore, NullDatastore
83from ..datastores.file_datastore.retrieve_artifacts import ZipIndex, retrieve_and_zip
84from ..datastores.file_datastore.transfer import retrieve_file_transfer_records
85from ..dimensions import DataCoordinate, Dimension, DimensionGroup
86from ..direct_query_driver import DirectQueryDriver
87from ..progress import Progress
88from ..queries import Query
89from ..registry import (
90 ConflictingDefinitionError,
91 DataIdError,
92 MissingDatasetTypeError,
93 RegistryDefaults,
94 _RegistryFactory,
95)
96from ..registry.sql_registry import SqlRegistry
97from ..transfers import RepoExportContext
98from ..utils import transactional
99from ._direct_butler_collections import DirectButlerCollections
101if TYPE_CHECKING:
102 from lsst.resources import ResourceHandleProtocol
104 from .._dataset_provenance import DatasetProvenance
105 from .._dataset_ref import DatasetId
106 from ..datastore import DatasetRefURIs
107 from ..dimensions import DataId, DataIdValue, DimensionElement, DimensionRecord, DimensionUniverse
108 from ..registry import CollectionArgType, Registry
109 from ..transfers import RepoImportBackend
111_LOG = getLogger(__name__)
114class ButlerValidationError(ValidationError):
115 """There is a problem with the Butler configuration."""
117 pass
120class DirectButler(Butler): # numpydoc ignore=PR02
121 """Main entry point for the data access system.
123 Parameters
124 ----------
125 config : `ButlerConfig`
126 The configuration for this Butler instance.
127 registry : `SqlRegistry`
128 The object that manages dataset metadata and relationships.
129 datastore : Datastore
130 The object that manages actual dataset storage.
131 storageClasses : StorageClassFactory
132 An object that maps known storage class names to objects that fully
133 describe them.
135 Notes
136 -----
137 Most users should call the top-level `Butler`.``from_config`` instead of
138 using this constructor directly.
139 """
141 # This is __new__ instead of __init__ because we have to support
142 # instantiation via the legacy constructor Butler.__new__(), which
143 # reads the configuration and selects which subclass to instantiate. The
144 # interaction between __new__ and __init__ is kind of wacky in Python. If
145 # we were using __init__ here, __init__ would be called twice (once when
146 # the DirectButler instance is constructed inside Butler.from_config(), and
147 # a second time with the original arguments to Butler() when the instance
148 # is returned from Butler.__new__()
149 def __new__(
150 cls,
151 *,
152 config: ButlerConfig,
153 registry: SqlRegistry,
154 datastore: Datastore,
155 storageClasses: StorageClassFactory,
156 metrics: ButlerMetrics | None = None,
157 ) -> DirectButler:
158 self = cast(DirectButler, super().__new__(cls))
159 self._config = config
160 self._registry = registry
161 self._datastore = datastore
162 self.storageClasses = storageClasses
163 self._metrics = metrics if metrics is not None else ButlerMetrics()
165 # For execution butler the datastore needs a special
166 # dependency-inversion trick. This is not used by regular butler,
167 # but we do not have a way to distinguish regular butler from execution
168 # butler.
169 self._datastore.set_retrieve_dataset_type_method(partial(_retrieve_dataset_type, registry))
171 self._closed = False
173 return self
175 @classmethod
176 def create_from_config(
177 cls,
178 config: ButlerConfig,
179 *,
180 options: ButlerInstanceOptions,
181 without_datastore: bool = False,
182 ) -> DirectButler:
183 """Construct a Butler instance from a configuration file.
185 Parameters
186 ----------
187 config : `ButlerConfig`
188 The configuration for this Butler instance.
189 options : `ButlerInstanceOptions`
190 Default values and other settings for the Butler instance.
191 without_datastore : `bool`, optional
192 If `True` do not attach a datastore to this butler. Any attempts
193 to use a datastore will fail.
195 Notes
196 -----
197 Most users should call the top-level `Butler`.``from_config``
198 instead of using this function directly.
199 """
200 if "run" in config or "collection" in config: 200 ↛ 201line 200 didn't jump to line 201 because the condition on line 200 was never true
201 raise ValueError("Passing a run or collection via configuration is no longer supported.")
203 defaults = RegistryDefaults.from_butler_instance_options(options)
204 try:
205 butlerRoot = config.get("root", config.configDir)
206 writeable = options.writeable
207 if writeable is None:
208 writeable = options.run is not None
209 registry = _RegistryFactory(config).from_config(
210 butlerRoot=butlerRoot, writeable=writeable, defaults=defaults
211 )
212 if without_datastore:
213 datastore: Datastore = NullDatastore(None, None)
214 else:
215 datastore = Datastore.fromConfig(
216 config, registry.getDatastoreBridgeManager(), butlerRoot=butlerRoot
217 )
218 # TODO: Once datastore drops dependency on registry we can
219 # construct datastore first and pass opaque tables to registry
220 # constructor.
221 registry.make_datastore_tables(datastore.get_opaque_table_definitions())
222 storageClasses = StorageClassFactory()
223 storageClasses.addFromConfig(config)
225 return DirectButler(
226 config=config,
227 registry=registry,
228 datastore=datastore,
229 storageClasses=storageClasses,
230 metrics=options.metrics,
231 )
232 except Exception:
233 # Failures here usually mean that configuration is incomplete,
234 # just issue an error message which includes config file URI.
235 _LOG.error(f"Failed to instantiate Butler from config {config.configFile}.")
236 raise
238 def clone(
239 self,
240 *,
241 collections: CollectionArgType | None | EllipsisType = ...,
242 run: str | None | EllipsisType = ...,
243 inferDefaults: bool | EllipsisType = ...,
244 dataId: dict[str, str] | EllipsisType = ...,
245 metrics: ButlerMetrics | None = None,
246 ) -> DirectButler:
247 # Docstring inherited
248 defaults = self._registry.defaults.clone(collections, run, inferDefaults, dataId)
249 registry = self._registry.copy(defaults)
251 return DirectButler(
252 registry=registry,
253 config=self._config,
254 datastore=self._datastore.clone(registry.getDatastoreBridgeManager()),
255 storageClasses=self.storageClasses,
256 metrics=metrics,
257 )
259 def close(self) -> None:
260 if not self._closed:
261 self._closed = True
262 self._registry.close()
263 # Cause exceptions to be raised if a user attempts to use the
264 # instance after closing it. Without this, Butler would still
265 # work after being closed because of implementation details
266 # of SqlAlchemy, but this may not continue to be the case in the
267 # future and we don't want users to get in the habit of doing this.
268 self._registry = _BUTLER_CLOSED_INSTANCE
269 self._datastore = _BUTLER_CLOSED_INSTANCE
271 GENERATION: ClassVar[int] = 3
272 """This is a Generation 3 Butler.
274 This attribute may be removed in the future, once the Generation 2 Butler
275 interface has been fully retired; it should only be used in transitional
276 code.
277 """
279 @classmethod
280 def _unpickle(
281 cls,
282 config: ButlerConfig,
283 collections: tuple[str, ...] | None,
284 run: str | None,
285 defaultDataId: dict[str, str],
286 writeable: bool,
287 ) -> DirectButler:
288 """Callable used to unpickle a Butler.
290 We prefer not to use ``Butler.__init__`` directly so we can force some
291 of its many arguments to be keyword-only (note that ``__reduce__``
292 can only invoke callables with positional arguments).
294 Parameters
295 ----------
296 config : `ButlerConfig`
297 Butler configuration, already coerced into a true `ButlerConfig`
298 instance (and hence after any search paths for overrides have been
299 utilized).
300 collections : `tuple` [ `str` ]
301 Names of the default collections to read from.
302 run : `str`, optional
303 Name of the default `~CollectionType.RUN` collection to write to.
304 defaultDataId : `dict` [ `str`, `str` ]
305 Default data ID values.
306 writeable : `bool`
307 Whether the Butler should support write operations.
309 Returns
310 -------
311 butler : `Butler`
312 A new `Butler` instance.
313 """
314 return cls.create_from_config(
315 config=config,
316 options=ButlerInstanceOptions(
317 collections=collections, run=run, writeable=writeable, kwargs=defaultDataId
318 ),
319 )
321 def __reduce__(self) -> tuple:
322 """Support pickling."""
323 return (
324 DirectButler._unpickle,
325 (
326 self._config,
327 self.collections.defaults,
328 self.run,
329 dict(self._registry.defaults.dataId.required),
330 self._registry.isWriteable(),
331 ),
332 )
334 def __str__(self) -> str:
335 return (
336 f"Butler(collections={self.collections}, run={self.run}, "
337 f"datastore='{self._datastore}', registry='{self._registry}')"
338 )
340 def isWriteable(self) -> bool:
341 # Docstring inherited.
342 return self._registry.isWriteable()
344 def _caching_context(self) -> contextlib.AbstractContextManager[None]:
345 """Context manager that enables caching."""
346 return self._registry.caching_context()
348 @contextlib.contextmanager
349 def transaction(self) -> Iterator[None]:
350 """Context manager supporting `Butler` transactions.
352 Transactions can be nested.
353 """
354 with self._registry.transaction(), self._datastore.transaction():
355 yield
357 def _standardizeArgs(
358 self,
359 datasetRefOrType: DatasetRef | DatasetType | str,
360 dataId: DataId | None = None,
361 for_put: bool = True,
362 **kwargs: Any,
363 ) -> tuple[DatasetType, DataId | None]:
364 """Standardize the arguments passed to several Butler APIs.
366 Parameters
367 ----------
368 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str`
369 When `DatasetRef` the `dataId` should be `None`.
370 Otherwise the `DatasetType` or name thereof.
371 dataId : `dict` or `DataCoordinate`
372 A `dict` of `Dimension` link name, value pairs that label the
373 `DatasetRef` within a Collection. When `None`, a `DatasetRef`
374 should be provided as the second argument.
375 for_put : `bool`, optional
376 If `True` this call is invoked as part of a `Butler.put`.
377 Otherwise it is assumed to be part of a `Butler.get()`. This
378 parameter is only relevant if there is dataset type
379 inconsistency.
380 **kwargs
381 Additional keyword arguments used to augment or construct a
382 `DataCoordinate`. See `DataCoordinate.standardize`
383 parameters.
385 Returns
386 -------
387 datasetType : `DatasetType`
388 A `DatasetType` instance extracted from ``datasetRefOrType``.
389 dataId : `dict` or `DataId`, optional
390 Argument that can be used (along with ``kwargs``) to construct a
391 `DataId`.
393 Notes
394 -----
395 Butler APIs that conceptually need a DatasetRef also allow passing a
396 `DatasetType` (or the name of one) and a `DataId` (or a dict and
397 keyword arguments that can be used to construct one) separately. This
398 method accepts those arguments and always returns a true `DatasetType`
399 and a `DataId` or `dict`.
401 Standardization of `dict` vs `DataId` is best handled by passing the
402 returned ``dataId`` (and ``kwargs``) to `Registry` APIs, which are
403 generally similarly flexible.
404 """
405 externalDatasetType: DatasetType | None = None
406 internalDatasetType: DatasetType | None = None
407 if isinstance(datasetRefOrType, DatasetRef):
408 if dataId is not None or kwargs:
409 raise ValueError("DatasetRef given, cannot use dataId as well")
410 externalDatasetType = datasetRefOrType.datasetType
411 dataId = datasetRefOrType.dataId
412 else:
413 # Don't check whether DataId is provided, because Registry APIs
414 # can usually construct a better error message when it wasn't.
415 if isinstance(datasetRefOrType, DatasetType):
416 externalDatasetType = datasetRefOrType
417 else:
418 internalDatasetType = self.get_dataset_type(datasetRefOrType)
420 # Check that they are self-consistent
421 if externalDatasetType is not None:
422 registryDatasetType = self._get_registry_dataset_type(externalDatasetType)
423 if registryDatasetType is None:
424 # The caller has asked for a component that only their own
425 # storage class defines, so there is no registry definition to
426 # check against and their definition has to be used as given.
427 internalDatasetType = externalDatasetType
428 else:
429 internalDatasetType = registryDatasetType
430 if externalDatasetType != internalDatasetType:
431 # We can allow differences if they are compatible,
432 # depending on whether this is a get or a put. A get
433 # requires that the python type associated with the
434 # datastore can be converted to the user type. A put
435 # requires that the user supplied python type can be
436 # converted to the internal type expected by registry.
437 relevantDatasetType = internalDatasetType
438 if for_put:
439 is_compatible = internalDatasetType.is_compatible_with(externalDatasetType)
440 else:
441 is_compatible = externalDatasetType.is_compatible_with(internalDatasetType)
442 relevantDatasetType = externalDatasetType
443 if not is_compatible:
444 raise ValueError(
445 f"Supplied dataset type ({externalDatasetType}) inconsistent with "
446 f"registry definition ({internalDatasetType})"
447 )
448 # Override the internal definition.
449 internalDatasetType = relevantDatasetType
451 assert internalDatasetType is not None
452 return internalDatasetType, dataId
454 def _get_registry_dataset_type(self, datasetType: DatasetType) -> DatasetType | None:
455 """Return the registry definition corresponding to the given dataset
456 type.
458 Parameters
459 ----------
460 datasetType : `DatasetType`
461 Dataset type, possibly a component, supplied by the caller.
463 Returns
464 -------
465 registry_type : `DatasetType` or `None`
466 The registry definition of ``datasetType``, or `None` if the
467 registry has no definition for it.
469 Raises
470 ------
471 MissingDatasetTypeError
472 Raised if the dataset type is not registered. For a component
473 dataset type this refers to the composite.
475 Notes
476 -----
477 Only composites are registered: a component dataset type is derived
478 from the storage class of its composite. A component that is defined
479 only by a read-time storage class override therefore has no registry
480 definition at all, even though its composite is registered, and `None`
481 is returned to say so rather than something that only resembles a
482 registry definition.
483 """
484 parent_name, component = DatasetType.splitDatasetTypeName(datasetType.name)
485 parent = self.get_dataset_type(parent_name)
486 if component is None:
487 return parent
488 if component in parent.storageClass.allComponents():
489 return parent.makeComponentDatasetType(component)
490 return None
492 def _rewrite_data_id(
493 self, dataId: DataId | None, datasetType: DatasetType, **kwargs: Any
494 ) -> tuple[DataId | None, dict[str, Any]]:
495 """Rewrite a data ID taking into account dimension records.
497 Take a Data ID and keyword args and rewrite it if necessary to
498 allow the user to specify dimension records rather than dimension
499 primary values.
501 This allows a user to include a dataId dict with keys of
502 ``exposure.day_obs`` and ``exposure.seq_num`` instead of giving
503 the integer exposure ID. It also allows a string to be given
504 for a dimension value rather than the integer ID if that is more
505 convenient. For example, rather than having to specifying the
506 detector with ``detector.full_name``, a string given for ``detector``
507 will be interpreted as the full name and converted to the integer
508 value.
510 Keyword arguments can also use strings for dimensions like detector
511 and exposure but python does not allow them to include ``.`` and
512 so the ``exposure.day_obs`` syntax can not be used in a keyword
513 argument.
515 Parameters
516 ----------
517 dataId : `dict` or `DataCoordinate`
518 A `dict` of `Dimension` link name, value pairs that will label the
519 `DatasetRef` within a Collection.
520 datasetType : `DatasetType`
521 The dataset type associated with this dataId. Required to
522 determine the relevant dimensions.
523 **kwargs
524 Additional keyword arguments used to augment or construct a
525 `DataId`. See `DataId` parameters.
527 Returns
528 -------
529 dataId : `dict` or `DataCoordinate`
530 The, possibly rewritten, dataId. If given a `DataCoordinate` and
531 no keyword arguments, the original dataId will be returned
532 unchanged.
533 **kwargs : `dict`
534 Any unused keyword arguments (would normally be empty dict).
535 """
536 # Process dimension records that are using record information
537 # rather than ids
538 newDataId: dict[str, DataIdValue] = {}
539 byRecord: dict[str, dict[str, Any]] = defaultdict(dict)
541 if isinstance(dataId, DataCoordinate):
542 # Do nothing if we have a DataCoordinate and no kwargs.
543 if not kwargs: 543 ↛ 547line 543 didn't jump to line 547 because the condition on line 543 was always true
544 return dataId, kwargs
545 # If we have a DataCoordinate with kwargs, we know the
546 # DataCoordinate only has values for real dimensions.
547 newDataId.update(dataId.mapping)
548 elif dataId:
549 # The data is mapping, which means it might have keys like
550 # "exposure.obs_id" (unlike kwargs, because a "." is not allowed in
551 # a keyword parameter).
552 for k, v in dataId.items():
553 if isinstance(k, str) and "." in k:
554 # Someone is using a more human-readable dataId
555 dimensionName, record = k.split(".", 1)
556 byRecord[dimensionName][record] = v
557 else:
558 newDataId[k] = v
560 # Go through the updated dataId and check the type in case someone is
561 # using an alternate key. We have already filtered out the compound
562 # keys dimensions.record format.
563 not_dimensions = {}
565 # Will need to look in the dataId and the keyword arguments
566 # and will remove them if they need to be fixed or are unrecognized.
567 for dataIdDict in (newDataId, kwargs):
568 # Use a list so we can adjust the dict safely in the loop
569 for dimensionName in list(dataIdDict):
570 value = dataIdDict[dimensionName]
571 try:
572 dimension = self.dimensions.dimensions[dimensionName]
573 except KeyError:
574 # This is not a real dimension
575 not_dimensions[dimensionName] = value
576 del dataIdDict[dimensionName]
577 continue
579 # Convert an integral type to an explicit int to simplify
580 # comparisons here
581 if isinstance(value, numbers.Integral):
582 value = int(value)
584 if not isinstance(value, dimension.primaryKey.getPythonType()):
585 for alternate in dimension.alternateKeys: 585 ↛ 598line 585 didn't jump to line 598 because the loop on line 585 didn't complete
586 if isinstance(value, alternate.getPythonType()): 586 ↛ 585line 586 didn't jump to line 585 because the condition on line 586 was always true
587 byRecord[dimensionName][alternate.name] = value
588 del dataIdDict[dimensionName]
589 _LOG.debug(
590 "Converting dimension %s to %s.%s=%s",
591 dimensionName,
592 dimensionName,
593 alternate.name,
594 value,
595 )
596 break
597 else:
598 _LOG.warning(
599 "Type mismatch found for value '%r' provided for dimension %s. "
600 "Could not find matching alternative (primary key has type %s) "
601 "so attempting to use as-is.",
602 value,
603 dimensionName,
604 dimension.primaryKey.getPythonType(),
605 )
607 # By this point kwargs and newDataId should only include valid
608 # dimensions. Merge kwargs in to the new dataId and log if there
609 # are dimensions in both (rather than calling update).
610 for k, v in kwargs.items():
611 if k in newDataId and newDataId[k] != v:
612 _LOG.debug(
613 "Keyword arg %s overriding explicit value in dataId of %s with %s", k, newDataId[k], v
614 )
615 newDataId[k] = v
616 # No need to retain any values in kwargs now.
617 kwargs = {}
619 # If we have some unrecognized dimensions we have to try to connect
620 # them to records in other dimensions. This is made more complicated
621 # by some dimensions having records with clashing names. A mitigation
622 # is that we can tell by this point which dimensions are missing
623 # for the DatasetType but this does not work for calibrations
624 # where additional dimensions can be used to constrain the temporal
625 # axis.
626 if not_dimensions:
627 # Search for all dimensions even if we have been given a value
628 # explicitly. In some cases records are given as well as the
629 # actually dimension and this should not be an error if they
630 # match.
631 mandatoryDimensions = datasetType.dimensions.names # - provided
633 candidateDimensions: set[str] = set()
634 candidateDimensions.update(mandatoryDimensions)
636 # For calibrations we may well be needing temporal dimensions
637 # so rather than always including all dimensions in the scan
638 # restrict things a little. It is still possible for there
639 # to be confusion over day_obs in visit vs exposure for example.
640 # If we are not searching calibration collections things may
641 # fail but they are going to fail anyway because of the
642 # ambiguousness of the dataId...
643 if datasetType.isCalibration():
644 for dim in self.dimensions.dimensions:
645 if dim.temporal:
646 candidateDimensions.add(str(dim))
648 # Look up table for the first association with a dimension
649 guessedAssociation: dict[str, dict[str, Any]] = defaultdict(dict)
651 # Keep track of whether an item is associated with multiple
652 # dimensions.
653 counter: Counter[str] = Counter()
654 assigned: dict[str, set[str]] = defaultdict(set)
656 # Go through the missing dimensions and associate the
657 # given names with records within those dimensions
658 matched_dims = set()
659 for dimensionName in candidateDimensions:
660 dimension = self.dimensions.dimensions[dimensionName]
661 fields = dimension.metadata.names | dimension.uniqueKeys.names
662 for field in not_dimensions:
663 if field in fields:
664 guessedAssociation[dimensionName][field] = not_dimensions[field]
665 counter[dimensionName] += 1
666 assigned[field].add(dimensionName)
667 matched_dims.add(field)
669 # Calculate the fields that matched nothing.
670 never_found = set(not_dimensions) - matched_dims
672 if never_found:
673 raise DimensionValueError(f"Unrecognized keyword args given: {never_found}")
675 # There is a chance we have allocated a single dataId item
676 # to multiple dimensions. Need to decide which should be retained.
677 # For now assume that the most popular alternative wins.
678 # This means that day_obs with seq_num will result in
679 # exposure.day_obs and not visit.day_obs
680 # Also prefer an explicitly missing dimension over an inferred
681 # temporal dimension.
682 for fieldName, assignedDimensions in assigned.items():
683 if len(assignedDimensions) > 1:
684 # Pick the most popular (preferring mandatory dimensions)
685 requiredButMissing = assignedDimensions.intersection(mandatoryDimensions)
686 if requiredButMissing: 686 ↛ 687line 686 didn't jump to line 687 because the condition on line 686 was never true
687 candidateDimensions = requiredButMissing
688 else:
689 candidateDimensions = assignedDimensions
691 # If this is a choice between visit and exposure and
692 # neither was a required part of the dataset type,
693 # (hence in this branch) always prefer exposure over
694 # visit since exposures are always defined and visits
695 # are defined from exposures.
696 if candidateDimensions == {"exposure", "visit"}: 696 ↛ 701line 696 didn't jump to line 701 because the condition on line 696 was always true
697 candidateDimensions = {"exposure"}
699 # Select the relevant items and get a new restricted
700 # counter.
701 theseCounts = {k: v for k, v in counter.items() if k in candidateDimensions}
702 duplicatesCounter: Counter[str] = Counter()
703 duplicatesCounter.update(theseCounts)
705 # Choose the most common. If they are equally common
706 # we will pick the one that was found first.
707 # Returns a list of tuples
708 selected = duplicatesCounter.most_common(1)[0][0]
710 _LOG.debug(
711 "Ambiguous dataId entry '%s' associated with multiple dimensions: %s."
712 " Removed ambiguity by choosing dimension %s.",
713 fieldName,
714 ", ".join(assignedDimensions),
715 selected,
716 )
718 for candidateDimension in assignedDimensions:
719 if candidateDimension != selected:
720 del guessedAssociation[candidateDimension][fieldName]
722 # Update the record look up dict with the new associations
723 for dimensionName, values in guessedAssociation.items():
724 if values: # A dict might now be empty
725 _LOG.debug(
726 "Assigned non-dimension dataId keys to dimension %s: %s", dimensionName, values
727 )
728 byRecord[dimensionName].update(values)
730 if byRecord:
731 # Some record specifiers were found so we need to convert
732 # them to the Id form
733 for dimensionName, values in byRecord.items():
734 if dimensionName in newDataId:
735 _LOG.debug(
736 "DataId specified explicit %s dimension value of %s in addition to"
737 " general record specifiers for it of %s. Checking for self-consistency.",
738 dimensionName,
739 newDataId[dimensionName],
740 str(values),
741 )
742 # Get the actual record and compare with these values.
743 # Only query with relevant data ID values.
744 filtered_data_id = {
745 k: v for k, v in newDataId.items() if k in self.dimensions[dimensionName].required
746 }
747 try:
748 recs = self.query_dimension_records(
749 dimensionName,
750 data_id=filtered_data_id,
751 )
752 except (DataIdError, EmptyQueryResultError):
753 raise DimensionValueError(
754 f"Could not find dimension '{dimensionName}'"
755 f" with dataId {filtered_data_id} as part of comparing with"
756 f" record values {byRecord[dimensionName]}"
757 ) from None
758 if len(recs) == 1: 758 ↛ 772line 758 didn't jump to line 772 because the condition on line 758 was always true
759 errmsg: list[str] = []
760 for k, v in values.items():
761 if (recval := getattr(recs[0], k)) != v:
762 errmsg.append(f"{k} ({recval} != {v})")
763 if errmsg:
764 raise DimensionValueError(
765 f"Dimension {dimensionName} in dataId has explicit value"
766 f" {newDataId[dimensionName]} inconsistent with"
767 f" {dimensionName} dimension record: " + ", ".join(errmsg)
768 )
769 else:
770 # Multiple matches for an explicit dimension
771 # should never happen but let downstream complain.
772 pass
773 continue
775 # Do not use data ID keys in query that aren't relevant.
776 # Otherwise we can have detector queries being constrained
777 # by an exposure ID that doesn't exist and return no matches
778 # for a detector even though it's a good detector name.
779 filtered_data_id = {
780 k: v
781 for k, v in newDataId.items()
782 if k in self.dimensions[dimensionName].minimal_group.names
783 }
785 def _get_attr(obj: Any, attr: str) -> Any:
786 # Used to implement x.exposure.seq_num when given
787 # x and "exposure.seq_num".
788 for component in attr.split("."):
789 obj = getattr(obj, component)
790 return obj
792 with self.query() as q:
793 x = q.expression_factory
794 # Build up a WHERE expression.
795 predicates = tuple(_get_attr(x, f"{dimensionName}.{k}") == v for k, v in values.items())
796 extra_args: dict[str, Any] = {} # For mypy.
797 extra_args.update(filtered_data_id)
798 extra_args.update(kwargs)
799 q = q.where(x.all(*predicates), **extra_args)
800 records = set(q.dimension_records(dimensionName))
802 if len(records) != 1:
803 if len(records) > 1:
804 # visit can have an ambiguous answer without involving
805 # visit_system. The default visit_system is defined
806 # by the instrument.
807 if ( 807 ↛ 812line 807 didn't jump to line 812 because the condition on line 807 was never true
808 dimensionName == "visit"
809 and "visit_system_membership" in self.dimensions
810 and "visit_system" in self.dimensions["instrument"].metadata
811 ):
812 instrument_records = self.query_dimension_records(
813 "instrument",
814 data_id=newDataId,
815 explain=False,
816 **kwargs,
817 )
818 if len(instrument_records) == 1:
819 visit_system = instrument_records[0].visit_system
820 if visit_system is None:
821 # Set to a value that will never match.
822 visit_system = -1
824 # Look up each visit in the
825 # visit_system_membership records.
826 for rec in records:
827 membership = self.query_dimension_records(
828 # Use bind to allow zero results.
829 # This is a fully-specified query.
830 "visit_system_membership",
831 instrument=instrument_records[0].name,
832 visit_system=visit_system,
833 visit=rec.id,
834 explain=False,
835 )
836 if membership:
837 # This record is the right answer.
838 records = {rec}
839 break
841 # The ambiguity may have been resolved so check again.
842 if len(records) > 1: 842 ↛ 860line 842 didn't jump to line 860 because the condition on line 842 was always true
843 _LOG.debug(
844 "Received %d records from constraints of %s", len(records), str(values)
845 )
846 for r in records:
847 _LOG.debug("- %s", str(r))
848 raise DimensionValueError(
849 f"DataId specification for dimension {dimensionName} is not"
850 f" uniquely constrained to a single dataset by {values}."
851 f" Got {len(records)} results."
852 )
853 else:
854 raise DimensionValueError(
855 f"DataId specification for dimension {dimensionName} matched no"
856 f" records when constrained by {values}"
857 )
859 # Get the primary key from the real dimension object
860 dimension = self.dimensions.dimensions[dimensionName]
861 if not isinstance(dimension, Dimension): 861 ↛ 862line 861 didn't jump to line 862 because the condition on line 861 was never true
862 raise RuntimeError(
863 f"{dimension.name} is not a true dimension, and cannot be used in data IDs."
864 )
865 newDataId[dimensionName] = getattr(records.pop(), dimension.primaryKey.name)
867 return newDataId, kwargs
869 def _findDatasetRef(
870 self,
871 datasetRefOrType: DatasetRef | DatasetType | str,
872 dataId: DataId | None = None,
873 *,
874 collections: Any = None,
875 predict: bool = False,
876 run: str | None = None,
877 datastore_records: bool = False,
878 timespan: Timespan | None = None,
879 **kwargs: Any,
880 ) -> DatasetRef:
881 """Shared logic for methods that start with a search for a dataset in
882 the registry.
884 Parameters
885 ----------
886 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str`
887 When `DatasetRef` the `dataId` should be `None`.
888 Otherwise the `DatasetType` or name thereof.
889 dataId : `dict` or `DataCoordinate`, optional
890 A `dict` of `Dimension` link name, value pairs that label the
891 `DatasetRef` within a Collection. When `None`, a `DatasetRef`
892 should be provided as the first argument.
893 collections : Any, optional
894 Collections to be searched, overriding ``self.collections``.
895 Can be any of the types supported by the ``collections`` argument
896 to butler construction.
897 predict : `bool`, optional
898 If `True`, return a newly created `DatasetRef` with a unique
899 dataset ID if finding a reference in the `Registry` fails.
900 Defaults to `False`.
901 run : `str`, optional
902 Run collection name to use for creating `DatasetRef` for predicted
903 datasets. Only used if ``predict`` is `True`.
904 datastore_records : `bool`, optional
905 If `True` add datastore records to returned `DatasetRef`.
906 timespan : `Timespan` or `None`, optional
907 A timespan that the validity range of the dataset must overlap.
908 If not provided and this is a calibration dataset type, an attempt
909 will be made to find the timespan from any temporal coordinate
910 in the data ID.
911 **kwargs
912 Additional keyword arguments used to augment or construct a
913 `DataId`. See `DataId` parameters.
915 Returns
916 -------
917 ref : `DatasetRef`
918 A reference to the dataset identified by the given arguments.
919 This can be the same dataset reference as given if it was
920 resolved.
922 Raises
923 ------
924 LookupError
925 Raised if no matching dataset exists in the `Registry` (and
926 ``predict`` is `False`).
927 ValueError
928 Raised if a resolved `DatasetRef` was passed as an input, but it
929 differs from the one found in the registry.
930 TypeError
931 Raised if no collections were provided.
932 """
933 datasetType, dataId = self._standardizeArgs(datasetRefOrType, dataId, for_put=False, **kwargs)
934 if isinstance(datasetRefOrType, DatasetRef):
935 if collections is not None: 935 ↛ 936line 935 didn't jump to line 936 because the condition on line 935 was never true
936 warnings.warn("Collections should not be specified with DatasetRef", stacklevel=3)
937 if predict and not datasetRefOrType.dataId.hasRecords():
938 return datasetRefOrType.expanded(self.registry.expandDataId(datasetRefOrType.dataId))
939 # May need to retrieve datastore records if requested.
940 if datastore_records and datasetRefOrType._datastore_records is None:
941 datasetRefOrType = self._registry.get_datastore_records(datasetRefOrType)
942 return datasetRefOrType
944 dataId, kwargs = self._rewrite_data_id(dataId, datasetType, **kwargs)
946 if datasetType.isCalibration():
947 # Because this is a calibration dataset, first try to make a
948 # standardize the data ID without restricting the dimensions to
949 # those of the dataset type requested, because there may be extra
950 # dimensions that provide temporal information for a validity-range
951 # lookup.
952 dataId = DataCoordinate.standardize(
953 dataId, universe=self.dimensions, defaults=self._registry.defaults.dataId, **kwargs
954 )
955 if timespan is None:
956 if dataId.dimensions.temporal:
957 dataId = self._registry.expandDataId(dataId)
958 # Use the timespan from the data ID to constrain the
959 # calibration lookup, but only if the caller has not
960 # specified an explicit timespan.
961 timespan = dataId.timespan
962 else:
963 # Try an arbitrary timespan. Downstream will fail if this
964 # results in more than one matching dataset.
965 timespan = Timespan(None, None)
966 else:
967 # Standardize the data ID to just the dimensions of the dataset
968 # type instead of letting registry.findDataset do it, so we get the
969 # result even if no dataset is found.
970 dataId = DataCoordinate.standardize(
971 dataId,
972 dimensions=datasetType.dimensions,
973 defaults=self._registry.defaults.dataId,
974 **kwargs,
975 )
976 # Always lookup the DatasetRef, even if one is given, to ensure it is
977 # present in the current collection.
978 ref = self.find_dataset(
979 datasetType,
980 dataId,
981 collections=collections,
982 timespan=timespan,
983 datastore_records=datastore_records,
984 )
985 if ref is None:
986 if predict:
987 if run is None: 987 ↛ 991line 987 didn't jump to line 991 because the condition on line 987 was always true
988 run = self.run
989 if run is None: 989 ↛ 990line 989 didn't jump to line 990 because the condition on line 989 was never true
990 raise TypeError("Cannot predict dataset ID/location with run=None.")
991 dataId = self.registry.expandDataId(dataId)
992 return DatasetRef(datasetType, dataId, run=run)
993 else:
994 if collections is None:
995 collections = self._registry.defaults.collections
996 raise DatasetNotFoundError(
997 f"Dataset {datasetType.name} with data ID {dataId} "
998 f"could not be found in collections {collections}."
999 )
1000 if datasetType != ref.datasetType:
1001 # If they differ it is because the user explicitly specified
1002 # a compatible dataset type to this call rather than using the
1003 # registry definition. The DatasetRef must therefore be recreated
1004 # using the user definition such that the expected type is
1005 # returned.
1006 ref = DatasetRef(
1007 datasetType, ref.dataId, run=ref.run, id=ref.id, datastore_records=ref._datastore_records
1008 )
1010 return ref
1012 @transactional
1013 def put(
1014 self,
1015 obj: Any,
1016 datasetRefOrType: DatasetRef | DatasetType | str,
1017 /,
1018 dataId: DataId | None = None,
1019 *,
1020 run: str | None = None,
1021 provenance: DatasetProvenance | None = None,
1022 **kwargs: Any,
1023 ) -> DatasetRef:
1024 """Store and register a dataset.
1026 Parameters
1027 ----------
1028 obj : `object`
1029 The dataset.
1030 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str`
1031 When `DatasetRef` is provided, ``dataId`` should be `None`.
1032 Otherwise the `DatasetType` or name thereof. If a fully resolved
1033 `DatasetRef` is given the run and ID are used directly.
1034 dataId : `dict` or `DataCoordinate`
1035 A `dict` of `Dimension` link name, value pairs that label the
1036 `DatasetRef` within a Collection. When `None`, a `DatasetRef`
1037 should be provided as the second argument.
1038 run : `str`, optional
1039 The name of the run the dataset should be added to, overriding
1040 ``self.run``. Not used if a resolved `DatasetRef` is provided.
1041 provenance : `DatasetProvenance` or `None`, optional
1042 Any provenance that should be attached to the serialized dataset.
1043 Not supported by all serialization mechanisms.
1044 **kwargs
1045 Additional keyword arguments used to augment or construct a
1046 `DataCoordinate`. See `DataCoordinate.standardize`
1047 parameters. Not used if a resolve `DatasetRef` is provided.
1049 Returns
1050 -------
1051 ref : `DatasetRef`
1052 A reference to the stored dataset, updated with the correct id if
1053 given.
1055 Raises
1056 ------
1057 TypeError
1058 Raised if the butler is read-only or if no run has been provided.
1059 """
1060 if isinstance(datasetRefOrType, DatasetRef):
1061 # This is a direct put of predefined DatasetRef.
1062 _LOG.debug("Butler put direct: %s", datasetRefOrType)
1063 if run is not None: 1063 ↛ 1064line 1063 didn't jump to line 1064 because the condition on line 1063 was never true
1064 warnings.warn("Run collection is not used for DatasetRef", stacklevel=3)
1066 with self._metrics.instrument_put(_LOG, msg="Dataset put direct"):
1067 # If registry already has a dataset with the same dataset ID,
1068 # dataset type and DataId, then _importDatasets will do
1069 # nothing and just return an original ref. We have to raise in
1070 # this case, there is a datastore check below for that.
1071 self._registry._importDatasets([datasetRefOrType], expand=True)
1072 # Before trying to write to the datastore check that it does
1073 # not know this dataset. This is prone to races, of course.
1074 if self._datastore.knows(datasetRefOrType):
1075 raise ConflictingDefinitionError(
1076 f"Datastore already contains dataset: {datasetRefOrType}"
1077 )
1078 # Try to write dataset to the datastore, if it fails due to a
1079 # race with another write, the content of stored data may be
1080 # unpredictable.
1081 try:
1082 self._datastore.put(obj, datasetRefOrType, provenance=provenance)
1083 except IntegrityError as e:
1084 raise ConflictingDefinitionError(f"Datastore already contains dataset: {e}") from e
1086 return datasetRefOrType
1088 _LOG.debug("Butler put: %s, dataId=%s, run=%s", datasetRefOrType, dataId, run)
1089 if not self.isWriteable(): 1089 ↛ 1090line 1089 didn't jump to line 1090 because the condition on line 1089 was never true
1090 raise TypeError("Butler is read-only.")
1092 with self._metrics.instrument_put(_LOG, msg="Dataset put with dataID"):
1093 datasetType, dataId = self._standardizeArgs(datasetRefOrType, dataId, **kwargs)
1095 # Handle dimension records in dataId
1096 dataId, kwargs = self._rewrite_data_id(dataId, datasetType, **kwargs)
1098 # Add Registry Dataset entry.
1099 dataId = self._registry.expandDataId(dataId, dimensions=datasetType.dimensions, **kwargs)
1100 (ref,) = self._registry.insertDatasets(datasetType, run=run, dataIds=[dataId])
1101 self._datastore.put(obj, ref, provenance=provenance)
1103 return ref
1105 def getDeferred(
1106 self,
1107 datasetRefOrType: DatasetRef | DatasetType | str,
1108 /,
1109 dataId: DataId | None = None,
1110 *,
1111 parameters: dict | None = None,
1112 collections: Any = None,
1113 storageClass: str | StorageClass | None = None,
1114 timespan: Timespan | None = None,
1115 **kwargs: Any,
1116 ) -> DeferredDatasetHandle:
1117 """Create a `DeferredDatasetHandle` which can later retrieve a dataset,
1118 after an immediate registry lookup.
1120 Parameters
1121 ----------
1122 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str`
1123 When `DatasetRef` the `dataId` should be `None`.
1124 Otherwise the `DatasetType` or name thereof.
1125 dataId : `dict` or `DataCoordinate`, optional
1126 A `dict` of `Dimension` link name, value pairs that label the
1127 `DatasetRef` within a Collection. When `None`, a `DatasetRef`
1128 should be provided as the first argument.
1129 parameters : `dict`
1130 Additional StorageClass-defined options to control reading,
1131 typically used to efficiently read only a subset of the dataset.
1132 collections : Any, optional
1133 Collections to be searched, overriding ``self.collections``.
1134 Can be any of the types supported by the ``collections`` argument
1135 to butler construction.
1136 storageClass : `StorageClass` or `str`, optional
1137 The storage class to be used to override the Python type
1138 returned by this method. By default the returned type matches
1139 the dataset type definition for this dataset. Specifying a
1140 read `StorageClass` can force a different type to be returned.
1141 This type must be compatible with the original type.
1142 timespan : `Timespan` or `None`, optional
1143 A timespan that the validity range of the dataset must overlap.
1144 If not provided and this is a calibration dataset type, an attempt
1145 will be made to find the timespan from any temporal coordinate
1146 in the data ID.
1147 **kwargs
1148 Additional keyword arguments used to augment or construct a
1149 `DataId`. See `DataId` parameters.
1151 Returns
1152 -------
1153 obj : `DeferredDatasetHandle`
1154 A handle which can be used to retrieve a dataset at a later time.
1156 Raises
1157 ------
1158 LookupError
1159 Raised if no matching dataset exists in the `Registry` or
1160 datastore.
1161 ValueError
1162 Raised if a resolved `DatasetRef` was passed as an input, but it
1163 differs from the one found in the registry.
1164 TypeError
1165 Raised if no collections were provided.
1166 """
1167 if isinstance(datasetRefOrType, DatasetRef):
1168 # Do the quick check first and if that fails, check for artifact
1169 # existence. This is necessary for datastores that are configured
1170 # in trust mode where there won't be a record but there will be
1171 # a file.
1172 if self._datastore.knows(datasetRefOrType) or self._datastore.exists(datasetRefOrType):
1173 ref = datasetRefOrType
1174 else:
1175 raise LookupError(f"Dataset reference {datasetRefOrType} does not exist.")
1176 else:
1177 ref = self._findDatasetRef(
1178 datasetRefOrType, dataId, collections=collections, timespan=timespan, **kwargs
1179 )
1180 return DeferredDatasetHandle(butler=self, ref=ref, parameters=parameters, storageClass=storageClass)
1182 def get(
1183 self,
1184 datasetRefOrType: DatasetRef | DatasetType | str,
1185 /,
1186 dataId: DataId | None = None,
1187 *,
1188 parameters: dict[str, Any] | None = None,
1189 collections: Any = None,
1190 storageClass: StorageClass | str | None = None,
1191 timespan: Timespan | None = None,
1192 **kwargs: Any,
1193 ) -> Any:
1194 """Retrieve a stored dataset.
1196 Parameters
1197 ----------
1198 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str`
1199 When `DatasetRef` the `dataId` should be `None`.
1200 Otherwise the `DatasetType` or name thereof.
1201 If a resolved `DatasetRef`, the associated dataset
1202 is returned directly without additional querying.
1203 dataId : `dict` or `DataCoordinate`
1204 A `dict` of `Dimension` link name, value pairs that label the
1205 `DatasetRef` within a Collection. When `None`, a `DatasetRef`
1206 should be provided as the first argument.
1207 parameters : `dict`
1208 Additional StorageClass-defined options to control reading,
1209 typically used to efficiently read only a subset of the dataset.
1210 collections : Any, optional
1211 Collections to be searched, overriding ``self.collections``.
1212 Can be any of the types supported by the ``collections`` argument
1213 to butler construction.
1214 storageClass : `StorageClass` or `str`, optional
1215 The storage class to be used to override the Python type
1216 returned by this method. By default the returned type matches
1217 the dataset type definition for this dataset. Specifying a
1218 read `StorageClass` can force a different type to be returned.
1219 This type must be compatible with the original type.
1220 timespan : `Timespan` or `None`, optional
1221 A timespan that the validity range of the dataset must overlap.
1222 If not provided and this is a calibration dataset type, an attempt
1223 will be made to find the timespan from any temporal coordinate
1224 in the data ID.
1225 **kwargs
1226 Additional keyword arguments used to augment or construct a
1227 `DataCoordinate`. See `DataCoordinate.standardize`
1228 parameters.
1230 Returns
1231 -------
1232 obj : `object`
1233 The dataset.
1235 Raises
1236 ------
1237 LookupError
1238 Raised if no matching dataset exists in the `Registry`.
1239 TypeError
1240 Raised if no collections were provided.
1242 Notes
1243 -----
1244 When looking up datasets in a `~CollectionType.CALIBRATION` collection,
1245 this method requires that the given data ID include temporal dimensions
1246 beyond the dimensions of the dataset type itself, in order to find the
1247 dataset with the appropriate validity range. For example, a "bias"
1248 dataset with native dimensions ``{instrument, detector}`` could be
1249 fetched with a ``{instrument, detector, exposure}`` data ID, because
1250 ``exposure`` is a temporal dimension.
1251 """
1252 _LOG.debug("Butler get: %s, dataId=%s, parameters=%s", datasetRefOrType, dataId, parameters)
1253 with self._metrics.instrument_get(_LOG, msg="Retrieved dataset"):
1254 ref = self._findDatasetRef(
1255 datasetRefOrType,
1256 dataId,
1257 collections=collections,
1258 datastore_records=True,
1259 timespan=timespan,
1260 **kwargs,
1261 )
1262 return self._datastore.get(ref, parameters=parameters, storageClass=storageClass)
1264 def getURIs(
1265 self,
1266 datasetRefOrType: DatasetRef | DatasetType | str,
1267 /,
1268 dataId: DataId | None = None,
1269 *,
1270 predict: bool = False,
1271 collections: Any = None,
1272 run: str | None = None,
1273 **kwargs: Any,
1274 ) -> DatasetRefURIs:
1275 """Return the URIs associated with the dataset.
1277 Parameters
1278 ----------
1279 datasetRefOrType : `DatasetRef`, `DatasetType`, or `str`
1280 When `DatasetRef` the `dataId` should be `None`.
1281 Otherwise the `DatasetType` or name thereof.
1282 dataId : `dict` or `DataCoordinate`
1283 A `dict` of `Dimension` link name, value pairs that label the
1284 `DatasetRef` within a Collection. When `None`, a `DatasetRef`
1285 should be provided as the first argument.
1286 predict : `bool`
1287 If `True`, allow URIs to be returned of datasets that have not
1288 been written.
1289 collections : Any, optional
1290 Collections to be searched, overriding ``self.collections``.
1291 Can be any of the types supported by the ``collections`` argument
1292 to butler construction.
1293 run : `str`, optional
1294 Run to use for predictions, overriding ``self.run``.
1295 **kwargs
1296 Additional keyword arguments used to augment or construct a
1297 `DataCoordinate`. See `DataCoordinate.standardize`
1298 parameters.
1300 Returns
1301 -------
1302 uris : `DatasetRefURIs`
1303 The URI to the primary artifact associated with this dataset (if
1304 the dataset was disassembled within the datastore this may be
1305 `None`), and the URIs to any components associated with the dataset
1306 artifact. (can be empty if there are no components).
1307 """
1308 ref = self._findDatasetRef(
1309 datasetRefOrType, dataId, predict=predict, run=run, collections=collections, **kwargs
1310 )
1311 return self._datastore.getURIs(ref, predict)
1313 def get_dataset_type(self, name: str) -> DatasetType:
1314 return self._registry.getDatasetType(name)
1316 def get_dataset(
1317 self,
1318 id: DatasetId | str,
1319 *,
1320 storage_class: str | StorageClass | None = None,
1321 dimension_records: bool = False,
1322 datastore_records: bool = False,
1323 ) -> DatasetRef | None:
1324 id = _to_uuid(id)
1325 ref = self._registry.getDataset(id)
1326 if ref is not None:
1327 if dimension_records: 1327 ↛ 1328line 1327 didn't jump to line 1328 because the condition on line 1327 was never true
1328 ref = ref.expanded(
1329 self._registry.expandDataId(ref.dataId, dimensions=ref.datasetType.dimensions)
1330 )
1331 if storage_class: 1331 ↛ 1332line 1331 didn't jump to line 1332 because the condition on line 1331 was never true
1332 ref = ref.overrideStorageClass(storage_class)
1333 if datastore_records: 1333 ↛ 1334line 1333 didn't jump to line 1334 because the condition on line 1333 was never true
1334 ref = self._registry.get_datastore_records(ref)
1335 return ref
1337 def get_many_datasets(self, ids: Iterable[DatasetId | str]) -> list[DatasetRef]:
1338 uuids = [_to_uuid(id) for id in ids]
1339 return self._registry._managers.datasets.get_dataset_refs(uuids)
1341 def find_dataset(
1342 self,
1343 dataset_type: DatasetType | str,
1344 data_id: DataId | None = None,
1345 *,
1346 collections: str | Sequence[str] | None = None,
1347 timespan: Timespan | None = None,
1348 storage_class: str | StorageClass | None = None,
1349 dimension_records: bool = False,
1350 datastore_records: bool = False,
1351 **kwargs: Any,
1352 ) -> DatasetRef | None:
1353 # Handle any parts of the dataID that are not using primary dimension
1354 # keys.
1355 if isinstance(dataset_type, str):
1356 actual_type = self.get_dataset_type(dataset_type)
1357 else:
1358 actual_type = dataset_type
1360 # Store the component for later.
1361 component_name = actual_type.component()
1362 if actual_type.isComponent():
1363 parent_type = actual_type.makeCompositeDatasetType()
1364 else:
1365 parent_type = actual_type
1367 data_id, kwargs = self._rewrite_data_id(data_id, parent_type, **kwargs)
1369 ref = self.registry.findDataset(
1370 parent_type,
1371 data_id,
1372 collections=collections,
1373 timespan=timespan,
1374 datastore_records=datastore_records,
1375 **kwargs,
1376 )
1377 if ref is not None and dimension_records: 1377 ↛ 1378line 1377 didn't jump to line 1378 because the condition on line 1377 was never true
1378 ref = ref.expanded(self._registry.expandDataId(ref.dataId, dimensions=ref.datasetType.dimensions))
1379 if ref is not None and component_name:
1380 ref = ref.makeComponentRef(component_name)
1381 if ref is not None and storage_class is not None: 1381 ↛ 1382line 1381 didn't jump to line 1382 because the condition on line 1381 was never true
1382 ref = ref.overrideStorageClass(storage_class)
1384 return ref
1386 def retrieve_artifacts_zip(
1387 self,
1388 refs: Iterable[DatasetRef],
1389 destination: ResourcePathExpression,
1390 overwrite: bool = True,
1391 ) -> ResourcePath:
1392 return retrieve_and_zip(refs, destination, self._datastore.retrieveArtifacts, overwrite)
1394 def retrieveArtifacts(
1395 self,
1396 refs: Iterable[DatasetRef],
1397 destination: ResourcePathExpression,
1398 transfer: str = "auto",
1399 preserve_path: bool = True,
1400 overwrite: bool = False,
1401 ) -> list[ResourcePath]:
1402 # Docstring inherited.
1403 outdir = ResourcePath(destination)
1404 artifact_map = self._datastore.retrieveArtifacts(
1405 refs,
1406 outdir,
1407 transfer=transfer,
1408 preserve_path=preserve_path,
1409 overwrite=overwrite,
1410 write_index=True,
1411 )
1412 return list(artifact_map)
1414 def exists(
1415 self,
1416 dataset_ref_or_type: DatasetRef | DatasetType | str,
1417 /,
1418 data_id: DataId | None = None,
1419 *,
1420 full_check: bool = True,
1421 collections: Any = None,
1422 **kwargs: Any,
1423 ) -> DatasetExistence:
1424 # Docstring inherited.
1425 existence = DatasetExistence.UNRECOGNIZED
1427 if isinstance(dataset_ref_or_type, DatasetRef):
1428 if collections is not None: 1428 ↛ 1429line 1428 didn't jump to line 1429 because the condition on line 1428 was never true
1429 warnings.warn("Collections should not be specified with DatasetRef", stacklevel=2)
1430 if data_id is not None: 1430 ↛ 1431line 1430 didn't jump to line 1431 because the condition on line 1430 was never true
1431 warnings.warn("A DataID should not be specified with DatasetRef", stacklevel=2)
1432 ref = dataset_ref_or_type
1433 registry_ref = self._registry.getDataset(dataset_ref_or_type.id)
1434 if registry_ref is not None:
1435 existence |= DatasetExistence.RECORDED
1437 if dataset_ref_or_type != registry_ref:
1438 # This could mean that storage classes differ, so we should
1439 # check for that but use the registry ref for the rest of
1440 # the method.
1441 if registry_ref.is_compatible_with(dataset_ref_or_type):
1442 # Use the registry version from now on.
1443 ref = registry_ref
1444 else:
1445 raise ValueError(
1446 f"The ref given to exists() ({ref}) has the same dataset ID as one "
1447 f"in registry but has different incompatible values ({registry_ref})."
1448 )
1449 else:
1450 try:
1451 ref = self._findDatasetRef(dataset_ref_or_type, data_id, collections=collections, **kwargs)
1452 except (LookupError, TypeError):
1453 return existence
1454 existence |= DatasetExistence.RECORDED
1456 if self._datastore.knows(ref):
1457 existence |= DatasetExistence.DATASTORE
1459 if full_check:
1460 if self._datastore.exists(ref):
1461 existence |= DatasetExistence._ARTIFACT
1462 elif existence.value != DatasetExistence.UNRECOGNIZED.value:
1463 # Do not add this flag if we have no other idea about a dataset.
1464 existence |= DatasetExistence(DatasetExistence._ASSUMED)
1466 return existence
1468 def _exists_many(
1469 self,
1470 refs: Iterable[DatasetRef],
1471 /,
1472 *,
1473 full_check: bool = True,
1474 ) -> dict[DatasetRef, DatasetExistence]:
1475 # Docstring inherited.
1476 existence = {ref: DatasetExistence.UNRECOGNIZED for ref in refs}
1478 # Check which refs exist in the registry.
1479 id_map = {ref.id: ref for ref in existence.keys()}
1480 for registry_ref in self.get_many_datasets(id_map.keys()):
1481 # Consistency between the given DatasetRef and the information
1482 # recorded in the registry is not verified.
1483 existence[id_map[registry_ref.id]] |= DatasetExistence.RECORDED
1485 # Ask datastore if it knows about these refs.
1486 knows = self._datastore.knows_these(refs)
1487 for ref, known in knows.items():
1488 if known:
1489 existence[ref] |= DatasetExistence.DATASTORE
1491 if full_check:
1492 mexists = self._datastore.mexists(refs)
1493 for ref, exists in mexists.items():
1494 if exists:
1495 existence[ref] |= DatasetExistence._ARTIFACT
1496 else:
1497 # Do not set this flag if nothing is known about the dataset.
1498 for ref in existence:
1499 if existence[ref] != DatasetExistence.UNRECOGNIZED:
1500 existence[ref] |= DatasetExistence._ASSUMED
1502 return existence
1504 def removeRuns(
1505 self,
1506 names: Iterable[str],
1507 unstore: bool | type[_DeprecatedDefault] = _DeprecatedDefault,
1508 *,
1509 unlink_from_chains: bool = False,
1510 ) -> None:
1511 # Docstring inherited.
1512 if not self.isWriteable(): 1512 ↛ 1513line 1512 didn't jump to line 1513 because the condition on line 1512 was never true
1513 raise TypeError("Butler is read-only.")
1515 if unstore is not _DeprecatedDefault: 1515 ↛ 1518line 1515 didn't jump to line 1518 because the condition on line 1515 was never true
1516 # The value was passed in by a user. Must report it is now
1517 # ignored.
1518 if unstore is True:
1519 msg = "The unstore parameter is deprecated and is now always treated as True. "
1520 else:
1521 msg = "The unstore parameter for removeRuns can no longer be False and is now ignored. "
1522 warnings.warn(
1523 msg + " The parameter will be removed after v30.",
1524 category=FutureWarning,
1525 stacklevel=find_outside_stacklevel("lsst.daf.butler"),
1526 )
1528 names = list(names)
1529 refs: list[DatasetRef] = []
1530 # Map of the chained collections to the RUN children.
1531 parents_to_children: dict[str, set[str]] = defaultdict(set)
1533 with self._caching_context():
1534 # Get information about these RUNs.
1535 collections_info = self.collections.query_info(names, include_parents=unlink_from_chains)
1536 for info in collections_info:
1537 if info.type is not CollectionType.RUN:
1538 raise TypeError(f"The collection type of '{info.name}' is {info.type.name}, not RUN.")
1539 if unlink_from_chains:
1540 if info.parents is None: # For mypy.
1541 raise AssertionError("Internal error: Collection parents required but not received")
1542 for parent in info.parents:
1543 parents_to_children[parent].add(info.name)
1545 # Update the names in case the query unexpectedly had a wildcard.
1546 names = [info.name for info in collections_info]
1548 # Get all the datasets from these runs.
1549 refs = self.query_all_datasets(names, find_first=False, limit=None)
1551 # Call pruneDatasets since we are deliberately removing
1552 # datasets in chunks from the RUN collections rather than
1553 # attempting to remove everything at once.
1554 with time_this(
1555 _LOG,
1556 msg="Removing %d dataset%s from %s",
1557 args=(len(refs), "s" if len(refs) != 1 else "", ", ".join(names)),
1558 ):
1559 self.pruneDatasets(refs, unstore=True, purge=True, disassociate=True)
1561 # Now can remove the actual RUN collection and unlink from chains.
1562 with self._registry.transaction():
1563 # This will fail if caller is not unlinking from chains but the
1564 # RUN is in a chain -- but we have already deleted all the datasets
1565 # by this point.
1566 if unlink_from_chains:
1567 # Use deterministic order for deletions to attempt to minimize
1568 # risk of deadlocks for parallel deletes.
1569 for parent in sorted(parents_to_children):
1570 self.collections.remove_from_chain(parent, sorted(parents_to_children[parent]))
1571 # Sort to avoid potential deadlocks.
1572 for name in sorted(names):
1573 # This should be fast since the collection should be empty.
1574 with time_this(_LOG, msg="Removing RUN collection %s", args=(name,)):
1575 self._registry.removeCollection(name)
1576 _LOG.info("Completely removed the following RUN collections: %s", ", ".join(names))
1578 def pruneDatasets(
1579 self,
1580 refs: Iterable[DatasetRef],
1581 *,
1582 disassociate: bool = True,
1583 unstore: bool = False,
1584 tags: Iterable[str] = (),
1585 purge: bool = False,
1586 ) -> None:
1587 # docstring inherited from LimitedButler
1589 if not self.isWriteable(): 1589 ↛ 1590line 1589 didn't jump to line 1590 because the condition on line 1589 was never true
1590 raise TypeError("Butler is read-only.")
1591 if purge:
1592 if not disassociate: 1592 ↛ 1593line 1592 didn't jump to line 1593 because the condition on line 1592 was never true
1593 raise TypeError("Cannot pass purge=True without disassociate=True.")
1594 if not unstore: 1594 ↛ 1595line 1594 didn't jump to line 1595 because the condition on line 1594 was never true
1595 raise TypeError("Cannot pass purge=True without unstore=True.")
1596 elif disassociate:
1597 tags = tuple(tags)
1598 if not tags: 1598 ↛ 1599line 1598 didn't jump to line 1599 because the condition on line 1598 was never true
1599 raise TypeError("No tags provided but disassociate=True.")
1600 for tag in tags:
1601 collectionType = self._registry.getCollectionType(tag)
1602 if collectionType is not CollectionType.TAGGED: 1602 ↛ 1603line 1602 didn't jump to line 1603 because the condition on line 1602 was never true
1603 raise TypeError(
1604 f"Cannot disassociate from collection '{tag}' "
1605 f"of non-TAGGED type {collectionType.name}."
1606 )
1607 # Transform possibly-single-pass iterable into something we can iterate
1608 # over multiple times.
1609 refs = list(refs)
1610 # Pruning a component of a DatasetRef makes no sense since registry
1611 # doesn't know about components and datastore might not store
1612 # components in a separate file
1613 for ref in refs:
1614 if ref.datasetType.component(): 1614 ↛ 1615line 1614 didn't jump to line 1615 because the condition on line 1614 was never true
1615 raise ValueError(f"Can not prune a component of a dataset (ref={ref})")
1617 # Chunk the deletions using a size that is reasonably efficient whilst
1618 # also giving reasonable feedback to the user. Chunking also minimizes
1619 # what needs to rollback if there is a failure and should allow
1620 # incremental re-running of the pruning (assuming the query is
1621 # repeated). The only issue will be if the Ctrl-C comes during
1622 # emptyTrash since an admin command would need to run to finish the
1623 # emptying of that chunk.
1624 progress = Progress("lsst.daf.butler.Butler.pruneDatasets", level=_LOG.INFO)
1625 chunk_size = 50_000
1626 n_chunks = math.ceil(len(refs) / chunk_size)
1627 if n_chunks > 1: 1627 ↛ 1628line 1627 didn't jump to line 1628 because the condition on line 1627 was never true
1628 _LOG.verbose("Pruning a total of %d datasets", len(refs))
1629 chunk_num = 0
1630 for chunked_refs in progress.wrap(
1631 chunk_iterable(refs, chunk_size=chunk_size), desc="Deleting datasets", total=n_chunks
1632 ):
1633 chunk_num += 1
1634 _LOG.verbose(
1635 "Pruning %d dataset%s in chunk %d/%d",
1636 len(chunked_refs),
1637 "s" if len(chunked_refs) != 1 else "",
1638 chunk_num,
1639 n_chunks,
1640 )
1641 with time_this(
1642 _LOG,
1643 msg="Removing %d datasets for chunk %d/%d",
1644 args=(len(chunked_refs), chunk_num, n_chunks),
1645 ):
1646 self._prune_datasets(
1647 chunked_refs, tags=tags, unstore=unstore, purge=purge, disassociate=disassociate
1648 )
1650 def _prune_datasets(
1651 self,
1652 refs: Collection[DatasetRef],
1653 *,
1654 disassociate: bool = True,
1655 unstore: bool = False,
1656 tags: Iterable[str] = (),
1657 purge: bool = False,
1658 ) -> None:
1659 # We don't need an unreliable Datastore transaction for this, because
1660 # we've been extra careful to ensure that Datastore.trash only involves
1661 # mutating the Registry (it can _look_ at Datastore-specific things,
1662 # but shouldn't change them), and hence all operations here are
1663 # Registry operations.
1664 with self.transaction():
1665 plural = "s" if len(refs) != 1 else ""
1666 if unstore:
1667 with time_this(
1668 _LOG,
1669 msg="Marking %d dataset%s for removal during pruneDatasets",
1670 args=(len(refs), plural),
1671 ):
1672 self._datastore.trash(refs)
1673 if purge:
1674 with time_this(
1675 _LOG, msg="Removing %d pruned dataset%s from registry", args=(len(refs), plural)
1676 ):
1677 self._registry.removeDatasets(refs)
1678 elif disassociate:
1679 assert tags, "Guaranteed by earlier logic in this function."
1680 with time_this(
1681 _LOG, msg="Disassociating %d dataset%ss from tagged collections", args=(len(refs), plural)
1682 ):
1683 for tag in tags:
1684 self._registry.disassociate(tag, refs)
1685 # We've exited the Registry transaction, and apparently committed.
1686 # (if there was an exception, everything rolled back, and it's as if
1687 # nothing happened - and we never get here).
1688 # Datastore artifacts are not yet gone, but they're clearly marked
1689 # as trash, so if we fail to delete now because of (e.g.) filesystem
1690 # problems we can try again later, and if manual administrative
1691 # intervention is required, it's pretty clear what that should entail:
1692 # deleting everything on disk and in private Datastore tables that is
1693 # in the dataset_location_trash table.
1694 if unstore:
1695 # Point of no return for removing artifacts. Restrict the trash
1696 # emptying to the refs that this call trashed.
1697 with time_this(
1698 _LOG,
1699 msg="Attempting to remove artifacts for %d dataset%s associated with pruning",
1700 args=(len(refs), plural),
1701 ):
1702 self._datastore.emptyTrash(refs=refs)
1704 def ingest_zip(
1705 self,
1706 zip_file: ResourcePathExpression,
1707 transfer: str = "auto",
1708 *,
1709 transfer_dimensions: bool = False,
1710 dry_run: bool = False,
1711 skip_existing: bool = False,
1712 ) -> None:
1713 # Docstring inherited.
1714 if not self.isWriteable(): 1714 ↛ 1715line 1714 didn't jump to line 1715 because the condition on line 1714 was never true
1715 raise TypeError("Butler is read-only.")
1717 zip_path = ResourcePath(zip_file)
1718 index = ZipIndex.from_zip_file(zip_path)
1719 _LOG.verbose(
1720 "Ingesting %s containing %d datasets and %d files.", zip_path, len(index.refs), len(index)
1721 )
1723 # Need to ingest the refs into registry. Re-use the standard ingest
1724 # code by reconstructing FileDataset from the index.
1725 refs = index.refs.to_refs(universe=self.dimensions)
1726 id_to_ref = {ref.id: ref for ref in refs}
1727 datasets: list[FileDataset] = []
1728 processed_ids: set[uuid.UUID] = set()
1729 for path_in_zip, index_info in index.artifact_map.items():
1730 # Disassembled composites need to check this ref isn't already
1731 # included.
1732 unprocessed = {id_ for id_ in index_info.ids if id_ not in processed_ids}
1733 if not unprocessed: 1733 ↛ 1734line 1733 didn't jump to line 1734 because the condition on line 1733 was never true
1734 continue
1735 dataset = FileDataset(refs=[id_to_ref[id_] for id_ in unprocessed], path=path_in_zip)
1736 datasets.append(dataset)
1737 processed_ids.update(unprocessed)
1739 new_datasets, existing_datasets = self._partition_datasets_by_known(datasets)
1740 if existing_datasets:
1741 if skip_existing:
1742 _LOG.info(
1743 "Skipping %d datasets from zip file %s which already exist in the repository.",
1744 len(existing_datasets),
1745 zip_file,
1746 )
1747 else:
1748 raise ConflictingDefinitionError(
1749 f"Datastore already contains {len(existing_datasets)} of the given datasets."
1750 f" Example: {existing_datasets[0]}"
1751 )
1752 if new_datasets: 1752 ↛ 1755line 1752 didn't jump to line 1755 because the condition on line 1752 was never true
1753 # Can not yet support partial zip ingests where a zip contains
1754 # some datasets that are already in another zip.
1755 raise ValueError(
1756 f"The given zip file from {zip_file} contains {len(new_datasets)} datasets not known "
1757 f"to this butler but also contains {len(existing_datasets)} datasets already known to "
1758 "this butler. Currently butler can not ingest zip files with overlapping content."
1759 )
1760 return
1762 # Ingest doesn't create the RUN collections so we have to do that
1763 # here.
1764 #
1765 # Sort by run collection name to ensure Postgres takes locks in the
1766 # same order between different processes, to mitigate an issue
1767 # where Postgres can deadlock due to the unique index on collection
1768 # name. (See DM-47543).
1769 runs = {ref.run for ref in refs}
1770 for run in sorted(runs):
1771 registered = self.collections.register(run)
1772 if registered:
1773 _LOG.verbose("Created RUN collection %s as part of zip ingest", run)
1775 progress = Progress("lsst.daf.butler.Butler.ingest", level=VERBOSE)
1776 import_info = self._prepare_ingest_file_datasets(
1777 datasets, progress, dry_run=dry_run, transfer_dimensions=transfer_dimensions
1778 )
1780 # Calculate some statistics based on the given list of datasets.
1781 n_datasets = 0
1782 for d in datasets:
1783 n_datasets += len(d.refs)
1784 srefs = "s" if n_datasets != 1 else ""
1786 with (
1787 self._metrics.instrument_ingest(
1788 n_datasets,
1789 _LOG,
1790 msg="Ingesting zip file %s with %s dataset%s",
1791 args=(zip_file, n_datasets, srefs),
1792 ),
1793 self.transaction(),
1794 ):
1795 # Do not need expanded dataset refs so can ignore the return value.
1796 self._ingest_file_datasets(datasets, import_info, progress, dry_run=dry_run)
1798 try:
1799 self._datastore.ingest_zip(zip_path, transfer=transfer, dry_run=dry_run)
1800 except IntegrityError as e:
1801 raise ConflictingDefinitionError(
1802 f"Datastore already contains one or more datasets: {e}"
1803 ) from e
1805 def _prepare_ingest_file_datasets(
1806 self,
1807 datasets: Sequence[FileDataset],
1808 progress: Progress,
1809 *,
1810 transfer_dimensions: bool = False,
1811 dry_run: bool = False,
1812 ) -> _ImportDatasetsInfo:
1813 # Track DataIDs that are being ingested so we can spot issues early
1814 # with duplication. Retain previous FileDataset so we can report it.
1815 groupedDataIds: MutableMapping[tuple[DatasetType, str], dict[DataCoordinate, FileDataset]] = (
1816 defaultdict(dict)
1817 )
1819 # All the refs we need to import.
1820 refs: list[DatasetRef] = []
1822 for dataset in progress.wrap(datasets, desc="Validating dataIDs"):
1823 for ref in dataset.refs:
1824 group_key = (ref.datasetType, ref.run)
1826 if ref.dataId in groupedDataIds[group_key]: 1826 ↛ 1827line 1826 didn't jump to line 1827 because the condition on line 1826 was never true
1827 raise ConflictingDefinitionError(
1828 f"Ingest conflict. Dataset {dataset.path} has same"
1829 " DataId as other ingest dataset"
1830 f" {groupedDataIds[group_key][ref.dataId].path} "
1831 f" ({ref.dataId})"
1832 )
1834 groupedDataIds[group_key][ref.dataId] = dataset
1835 refs.extend(dataset.refs)
1837 # Ensure that dataset types are created and all ref information
1838 # extracted.
1839 import_info = self._prepare_for_import_refs(
1840 self,
1841 refs,
1842 register_dataset_types=True,
1843 dry_run=dry_run,
1844 transfer_dimensions=transfer_dimensions,
1845 )
1846 return import_info
1848 def _ingest_file_datasets(
1849 self,
1850 datasets: Sequence[FileDataset],
1851 import_info: _ImportDatasetsInfo,
1852 progress: Progress,
1853 *,
1854 dry_run: bool = False,
1855 ) -> None:
1856 self._import_dimension_records(import_info.dimension_records, dry_run=dry_run)
1857 imported_refs = self._import_grouped_refs(
1858 import_info.grouped_refs, None, progress, dry_run=dry_run, expand_refs=True
1859 )
1861 # The expanded refs need to be attached back to the original
1862 # FileDatasets for datastore to use.
1863 id_to_ref = {ref.id: ref for ref in imported_refs}
1865 for dataset in progress.wrap(datasets, desc="Re-attaching expanded refs"):
1866 dataset.refs = [id_to_ref[ref.id] for ref in dataset.refs]
1868 def ingest(
1869 self,
1870 *datasets: FileDataset,
1871 transfer: str | None = "auto",
1872 record_validation_info: bool = True,
1873 skip_existing: bool = False,
1874 ) -> None:
1875 # Docstring inherited.
1876 if not datasets:
1877 return
1878 if not self.isWriteable(): 1878 ↛ 1879line 1878 didn't jump to line 1879 because the condition on line 1878 was never true
1879 raise TypeError("Butler is read-only.")
1880 _LOG.verbose("Ingesting %d file dataset%s.", len(datasets), "" if len(datasets) == 1 else "s")
1881 progress = Progress("lsst.daf.butler.Butler.ingest", level=VERBOSE)
1883 new_datasets, existing_datasets = self._partition_datasets_by_known(datasets)
1884 if existing_datasets:
1885 if skip_existing:
1886 _LOG.info(
1887 "Skipping %d datasets which already exist in the repository.", len(existing_datasets)
1888 )
1889 else:
1890 raise ConflictingDefinitionError(
1891 f"Datastore already contains {len(existing_datasets)} of the given datasets."
1892 f" Example: {existing_datasets[0]}"
1893 )
1895 # Calculate some statistics based on the given list of datasets.
1896 n_files = len(datasets)
1897 n_datasets = 0
1898 for d in datasets:
1899 n_datasets += len(d.refs)
1900 sfiles = "s" if n_files != 1 else ""
1901 srefs = "s" if n_datasets != 1 else ""
1903 # We use `datasets` rather `new_datasets` for the Registry
1904 # portion of this, to let it confirm that everything matches the
1905 # existing datasets.
1906 import_info = self._prepare_ingest_file_datasets(datasets, progress)
1908 with (
1909 self._metrics.instrument_ingest(
1910 n_datasets,
1911 _LOG,
1912 msg="Ingesting %s file%s with %s dataset%s",
1913 args=(n_files, sfiles, n_datasets, srefs),
1914 ),
1915 self.transaction(),
1916 ):
1917 self._ingest_file_datasets(datasets, import_info, progress)
1919 # Bulk-insert everything into Datastore.
1920 # We do not know if any of the registry entries already existed
1921 # (_importDatasets only complains if they exist but differ).
1922 # The _partition_datasets_by_known logic above should catch most
1923 # instances where we attempt to re-ingest files that were already
1924 # ingested, but a concurrent writer could cause a unique constraint
1925 # violation here.
1926 try:
1927 self._datastore.ingest(
1928 *new_datasets, transfer=transfer, record_validation_info=record_validation_info
1929 )
1930 except IntegrityError as e:
1931 raise ConflictingDefinitionError(
1932 f"Datastore already contains one or more datasets: {e}"
1933 ) from e
1935 def _partition_datasets_by_known(
1936 self, datasets: Iterable[FileDataset]
1937 ) -> tuple[list[FileDataset], list[FileDataset]]:
1938 """Divides the given `FileDataset` objects into two groups: those for
1939 which the Datastore already has an entry, and those for which it does
1940 not.
1941 """
1942 new_datasets = []
1943 existing_datasets = []
1945 refs = itertools.chain.from_iterable(dataset.refs for dataset in datasets)
1946 known_refs = self._datastore.knows_these(refs)
1948 for dataset in datasets:
1949 if any(known_refs[ref] for ref in dataset.refs):
1950 existing_datasets.append(dataset)
1951 else:
1952 new_datasets.append(dataset)
1954 return new_datasets, existing_datasets
1956 @contextlib.contextmanager
1957 def export(
1958 self,
1959 *,
1960 directory: str | None = None,
1961 filename: str | None = None,
1962 format: str | None = None,
1963 transfer: str | None = None,
1964 ) -> Iterator[RepoExportContext]:
1965 # Docstring inherited.
1966 if directory is None and transfer is not None: 1966 ↛ 1967line 1966 didn't jump to line 1967 because the condition on line 1966 was never true
1967 raise TypeError("Cannot transfer without providing a directory.")
1968 if transfer == "move": 1968 ↛ 1969line 1968 didn't jump to line 1969 because the condition on line 1968 was never true
1969 raise TypeError("Transfer may not be 'move': export is read-only")
1970 if format is None:
1971 if filename is None: 1971 ↛ 1972line 1971 didn't jump to line 1972 because the condition on line 1971 was never true
1972 raise TypeError("At least one of 'filename' or 'format' must be provided.")
1973 else:
1974 _, format = os.path.splitext(filename)
1975 if not format:
1976 raise ValueError("Please specify a file extension to determine export format.")
1977 format = format[1:] # Strip leading ".""
1978 elif filename is None: 1978 ↛ 1980line 1978 didn't jump to line 1980 because the condition on line 1978 was always true
1979 filename = f"export.{format}"
1980 if directory is not None:
1981 filename = os.path.join(directory, filename)
1982 formats = self._config["repo_transfer_formats"]
1983 if format not in formats:
1984 raise ValueError(f"Unknown export format {format!r}, allowed: {','.join(formats.keys())}")
1985 BackendClass = get_class_of(formats[format, "export"])
1986 with open(filename, "w") as stream:
1987 backend = BackendClass(stream, universe=self.dimensions)
1988 try:
1989 helper = RepoExportContext(self, backend=backend, directory=directory, transfer=transfer)
1990 with self._caching_context():
1991 yield helper
1992 except BaseException:
1993 raise
1994 else:
1995 helper._finish()
1997 def import_(
1998 self,
1999 *,
2000 directory: ResourcePathExpression | None = None,
2001 filename: ResourcePathExpression | TextIO | None = None,
2002 format: str | None = None,
2003 transfer: str | None = None,
2004 skip_dimensions: set | None = None,
2005 record_validation_info: bool = True,
2006 without_datastore: bool = False,
2007 ) -> None:
2008 # Docstring inherited.
2009 if not self.isWriteable(): 2009 ↛ 2010line 2009 didn't jump to line 2010 because the condition on line 2009 was never true
2010 raise TypeError("Butler is read-only.")
2011 if filename is None and format is not None: 2011 ↛ 2012line 2011 didn't jump to line 2012 because the condition on line 2011 was never true
2012 filename = ResourcePath(f"export.{format}", forceAbsolute=False)
2013 if directory is not None:
2014 directory = ResourcePath(directory, forceDirectory=True)
2015 # mypy doesn't think this will work but it does in python >= 3.10.
2016 if isinstance(filename, ResourcePathExpression): # type: ignore
2017 filename = ResourcePath(filename, forceAbsolute=False) # type: ignore
2018 if format is None: 2018 ↛ 2020line 2018 didn't jump to line 2020 because the condition on line 2018 was always true
2019 format = filename.getExtension()
2020 if not filename.isabs() and directory is not None: 2020 ↛ 2021line 2020 didn't jump to line 2021 because the condition on line 2020 was never true
2021 potential = directory.join(filename)
2022 exists_in_cwd = filename.exists()
2023 exists_in_dir = potential.exists()
2024 if exists_in_cwd and exists_in_dir:
2025 _LOG.warning(
2026 "A relative path for filename was specified (%s) which exists relative to cwd. "
2027 "Additionally, the file exists relative to the given search directory (%s). "
2028 "Using the export file in the given directory.",
2029 filename,
2030 potential,
2031 )
2032 # Given they specified an explicit directory and that
2033 # directory has the export file in it, assume that that
2034 # is what was meant despite the file in cwd.
2035 filename = potential
2036 elif exists_in_dir:
2037 filename = potential
2038 elif not exists_in_cwd and not exists_in_dir:
2039 # Raise early.
2040 raise FileNotFoundError(
2041 f"Export file could not be found in {filename.abspath()} or {potential.abspath()}."
2042 )
2043 elif format is None: 2043 ↛ 2044line 2043 didn't jump to line 2044 because the condition on line 2043 was never true
2044 format = ".yaml"
2045 BackendClass: type[RepoImportBackend] = get_class_of(
2046 self._config["repo_transfer_formats"][format]["import"]
2047 )
2049 def doImport(importStream: TextIO | ResourceHandleProtocol) -> None:
2050 with self._caching_context():
2051 backend = BackendClass(importStream, self) # type: ignore[call-arg]
2052 backend.register()
2053 with self.transaction():
2054 backend.load(
2055 datastore=self._datastore if not without_datastore else None,
2056 directory=directory,
2057 transfer=transfer,
2058 skip_dimensions=skip_dimensions,
2059 record_validation_info=record_validation_info,
2060 )
2062 if isinstance(filename, ResourcePath):
2063 # We can not use open() here at the moment because of
2064 # DM-38589 since yaml does stream.read(8192) in a loop.
2065 stream = io.StringIO(filename.read().decode())
2066 doImport(stream)
2067 else:
2068 doImport(filename) # type: ignore
2070 def transfer_dimension_records_from(
2071 self, source_butler: LimitedButler | Butler, source_refs: Iterable[DatasetRef | DataCoordinate]
2072 ) -> None:
2073 # Allowed dimensions in the target butler.
2074 elements = frozenset(element for element in self.dimensions.elements if element.has_own_table)
2076 data_ids = {ref.dataId for ref in source_refs}
2078 dimension_records = self._extract_all_dimension_records_from_data_ids(
2079 source_butler, data_ids, elements
2080 )
2082 # Insert order is important.
2083 for element in self.dimensions.sorted(dimension_records.keys()):
2084 records = [r for r in dimension_records[element].values()]
2085 # Assume that if the record is already present that we can
2086 # use it without having to check that the record metadata
2087 # is consistent.
2088 self._registry.insertDimensionData(element, *records, skip_existing=True)
2089 _LOG.debug("Dimension '%s' -- number of records transferred: %d", element.name, len(records))
2091 def _extract_all_dimension_records_from_data_ids(
2092 self,
2093 source_butler: LimitedButler | Butler,
2094 data_ids: set[DataCoordinate],
2095 allowed_elements: frozenset[DimensionElement],
2096 ) -> dict[DimensionElement, dict[DataCoordinate, DimensionRecord]]:
2097 primary_records = self._extract_dimension_records_from_data_ids(
2098 source_butler, data_ids, allowed_elements
2099 )
2101 additional_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] = defaultdict(dict)
2102 for original_element, record_mapping in primary_records.items():
2103 # Get dimensions that depend on this dimension.
2104 populated_by = self.dimensions.get_elements_populated_by(
2105 self.dimensions[original_element.name] # type: ignore
2106 )
2107 if populated_by:
2108 for element in populated_by:
2109 if element not in allowed_elements: 2109 ↛ 2110line 2109 didn't jump to line 2110 because the condition on line 2109 was never true
2110 continue
2111 if element.name == original_element.name:
2112 continue
2114 if element.name in primary_records:
2115 # If this element has already been stored avoid
2116 # re-finding records since that may lead to additional
2117 # spurious records. e.g. visit is populated_by
2118 # visit_detector_region but querying
2119 # visit_detector_region by visit will return all the
2120 # detectors for this visit -- the visit dataId does not
2121 # constrain this.
2122 # To constrain the query the original dataIds would
2123 # have to be scanned.
2124 continue
2126 if record_mapping: 2126 ↛ 2108line 2126 didn't jump to line 2108 because the condition on line 2126 was always true
2127 if not isinstance(source_butler, Butler): 2127 ↛ 2128line 2127 didn't jump to line 2128 because the condition on line 2127 was never true
2128 raise RuntimeError(
2129 f"Transferring populated_by records like {element.name}"
2130 " requires a full Butler."
2131 )
2133 with source_butler.query() as query:
2134 records = query.join_data_coordinates(record_mapping.keys()).dimension_records(
2135 element.name
2136 )
2137 for record in records:
2138 additional_records[record.definition].setdefault(record.dataId, record)
2140 # The next step is to walk back through the additional records to
2141 # pick up any missing content (such as visit_definition needing to
2142 # know the exposure). Want to ensure we do not request records we
2143 # already have.
2144 missing_data_ids = set()
2145 for record_mapping in additional_records.values():
2146 for data_id in record_mapping.keys():
2147 for dimension in data_id.dimensions.required:
2148 element = source_butler.dimensions[dimension]
2149 dimension_key = data_id.subset(dimension)
2150 if dimension_key not in primary_records[element]:
2151 missing_data_ids.add(dimension_key)
2153 # Fill out the new records. Assume that these new records do not
2154 # also need to carry over additional populated_by records.
2155 secondary_records = self._extract_dimension_records_from_data_ids(
2156 source_butler, missing_data_ids, allowed_elements
2157 )
2159 # Merge the extra sets of records in with the original.
2160 for name, record_mapping in itertools.chain(additional_records.items(), secondary_records.items()):
2161 primary_records[name].update(record_mapping)
2163 return primary_records
2165 def _extract_dimension_records_from_data_ids(
2166 self,
2167 source_butler: LimitedButler | Butler,
2168 data_ids: Iterable[DataCoordinate],
2169 allowed_elements: frozenset[DimensionElement],
2170 ) -> dict[DimensionElement, dict[DataCoordinate, DimensionRecord]]:
2171 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] = defaultdict(dict)
2173 data_ids = set(data_ids)
2174 if not all(data_id.hasRecords() for data_id in data_ids):
2175 if isinstance(source_butler, Butler): 2175 ↛ 2178line 2175 didn't jump to line 2178 because the condition on line 2175 was always true
2176 data_ids = source_butler._expand_data_ids(data_ids)
2177 else:
2178 raise TypeError("Input butler needs to be a full butler to expand DataId.")
2180 for data_id in data_ids:
2181 # If this butler doesn't know about a dimension in the source
2182 # butler things will break later.
2183 for element_name in data_id.dimensions.elements:
2184 record = data_id.records[element_name]
2185 if record is not None and record.definition in allowed_elements:
2186 dimension_records[record.definition].setdefault(record.dataId, record)
2188 return dimension_records
2190 def _cast_universe_for_import_refs(
2191 self, source_refs: Iterable[DatasetRef]
2192 ) -> Mapping[DatasetType, list[DatasetRef]]:
2193 """Try to cast imported refs to the target universe if possible.
2195 Parameters
2196 ----------
2197 source_refs
2198 The refs to be imported.
2200 Returns
2201 -------
2202 refs
2203 The refs to be imported, grouped by dataset type, with the dataset
2204 types cast to the target universe.
2206 Raises
2207 ------
2208 InconsistentUniverseError
2209 Raised if any reference cannot be converted to target universe.
2211 Notes
2212 -----
2213 Potentially this method can perform a non-trivial migrations of the
2214 datasets by modifying dimensions and dataIds. Presently though it can
2215 only perform a trivial validation of the dataset types compatibility.
2216 Returned mapping will contain dataset types in the new universe, but
2217 returned references will still have the original dataset types as
2218 there is presently no easy way to replace dataset type in a reference.
2219 """
2220 # In theory input refs could come from multiple universes, but in
2221 # practice this will not happen, so just group everything by dataset
2222 # type.
2223 refs_by_source_type: defaultdict[DatasetType, list[DatasetRef]] = defaultdict(list)
2224 for ref in source_refs:
2225 refs_by_source_type[ref.datasetType].append(ref)
2227 refs_by_type: defaultdict[DatasetType, list[DatasetRef]] = defaultdict(list)
2228 for source_type, refs in refs_by_source_type.items():
2229 refs_by_type[source_type.conform_to(self.dimensions)] = refs
2231 return refs_by_type
2233 def _prepare_for_import_refs(
2234 self,
2235 source_butler: LimitedButler,
2236 source_refs: Iterable[DatasetRef],
2237 *,
2238 register_dataset_types: bool = False,
2239 transfer_dimensions: bool = False,
2240 dry_run: bool = False,
2241 ) -> _ImportDatasetsInfo:
2242 # Docstring inherited.
2243 if not self.isWriteable() and not dry_run: 2243 ↛ 2244line 2243 didn't jump to line 2244 because the condition on line 2243 was never true
2244 raise TypeError("Butler is read-only.")
2246 # Will iterate through the refs multiple times so need to convert
2247 # to a list if this isn't a collection.
2248 if not isinstance(source_refs, collections.abc.Collection): 2248 ↛ 2249line 2248 didn't jump to line 2249 because the condition on line 2248 was never true
2249 source_refs = list(source_refs)
2251 original_count = len(source_refs)
2252 log_level = _LOG.INFO if original_count > 1 else _LOG.VERBOSE
2253 _LOG.log(
2254 log_level,
2255 "Importing %d dataset%s into %s",
2256 original_count,
2257 "s" if original_count != 1 else "",
2258 str(self),
2259 )
2261 refs_by_type = self._cast_universe_for_import_refs(source_refs)
2263 # Importing requires that we group the refs by dimension group and run
2264 # before doing the import.
2265 grouped_refs: defaultdict[_RefGroup, list[DatasetRef]] = defaultdict(list)
2266 for ref in source_refs:
2267 grouped_refs[_RefGroup(ref.datasetType.dimensions, ref.run)].append(ref)
2269 # Check to see if the dataset type in the source butler has
2270 # the same definition in the target butler and register missing
2271 # ones if requested. Registration must happen outside a transaction.
2272 newly_registered_dataset_types = set()
2273 for datasetType in refs_by_type:
2274 if register_dataset_types:
2275 # Let this raise immediately if inconsistent. Continuing
2276 # on to find additional inconsistent dataset types
2277 # might result in additional unwanted dataset types being
2278 # registered.
2279 try:
2280 if not dry_run and self._registry.registerDatasetType(datasetType):
2281 newly_registered_dataset_types.add(datasetType)
2282 except ConflictingDefinitionError as e:
2283 # Be safe and require that conversions be bidirectional
2284 # when there are storage class mismatches. This is because
2285 # get() will have to support conversion from source to
2286 # target python type (the source formatter will be
2287 # returning source python type) but there also is an
2288 # expectation that people will want to be able to get() in
2289 # the target using the source python type, which will not
2290 # require conversion for transferred datasets but might
2291 # for target-native types. Additionally, butler.get does
2292 # not know that the formatter will return the wrong
2293 # python type and so will always check that the conversion
2294 # works even though it won't need it.
2295 target_dataset_type = self.get_dataset_type(datasetType.name)
2296 target_compatible_with_source = target_dataset_type.is_compatible_with(datasetType)
2297 source_compatible_with_target = datasetType.is_compatible_with(target_dataset_type)
2298 if not (target_compatible_with_source and source_compatible_with_target): 2298 ↛ 2299line 2298 didn't jump to line 2299 because the condition on line 2298 was never true
2299 if target_compatible_with_source:
2300 e.add_note(
2301 "Target dataset type storage class is compatible with source "
2302 "but the reverse is not true."
2303 )
2304 elif source_compatible_with_target:
2305 e.add_note(
2306 "Source dataset type storage class is compatible with target "
2307 "but the reverse is not true."
2308 )
2309 else:
2310 e.add_note("If storage classes differ, please register converters.")
2311 raise
2312 else:
2313 # If the dataset type is missing, let it fail immediately.
2314 target_dataset_type = self.get_dataset_type(datasetType.name)
2315 if target_dataset_type != datasetType:
2316 target_compatible_with_source = target_dataset_type.is_compatible_with(datasetType)
2317 source_compatible_with_target = datasetType.is_compatible_with(target_dataset_type)
2318 # Both conversion directions are currently required.
2319 if not (target_compatible_with_source and source_compatible_with_target):
2320 msg = ""
2321 if target_compatible_with_source: 2321 ↛ 2322line 2321 didn't jump to line 2322 because the condition on line 2321 was never true
2322 msg = (
2323 "Target storage class is compatible with the source storage class "
2324 "but the reverse is not true."
2325 )
2326 elif source_compatible_with_target: 2326 ↛ 2327line 2326 didn't jump to line 2327 because the condition on line 2326 was never true
2327 msg = (
2328 "Source storage class is compatible with the target storage class"
2329 " but the reverse is not true."
2330 )
2331 else:
2332 msg = "If storage classes differ register converters."
2333 if msg: 2333 ↛ 2335line 2333 didn't jump to line 2335 because the condition on line 2333 was always true
2334 msg = f"({msg})"
2335 raise ConflictingDefinitionError(
2336 "Source butler dataset type differs from definition"
2337 f" in target butler: {datasetType} !="
2338 f" {target_dataset_type} {msg}"
2339 )
2340 if newly_registered_dataset_types:
2341 # We may have registered some even if there were inconsistencies
2342 # but should let people know (or else remove them again).
2343 _LOG.verbose(
2344 "Registered the following dataset types in the target Butler: %s",
2345 ", ".join(d.name for d in newly_registered_dataset_types),
2346 )
2347 else:
2348 _LOG.verbose("All required dataset types are known to the target Butler")
2350 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]] = defaultdict(dict)
2351 if transfer_dimensions:
2352 # Collect all the dimension records for these refs.
2353 # All dimensions are to be copied but the list of valid dimensions
2354 # come from this butler's universe.
2355 elements = frozenset(element for element in self.dimensions.elements if element.has_own_table)
2356 dataIds = {ref.dataId for ref in source_refs}
2357 dimension_records = self._extract_all_dimension_records_from_data_ids(
2358 source_butler, dataIds, elements
2359 )
2360 return _ImportDatasetsInfo(grouped_refs, dimension_records)
2362 def _import_dimension_records(
2363 self,
2364 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]],
2365 *,
2366 dry_run: bool,
2367 ) -> None:
2368 """Import dimension records collected during import pre-process."""
2369 if dimension_records and not dry_run:
2370 _LOG.verbose("Ensuring that dimension records exist for transferred datasets.")
2371 # Order matters.
2372 for element in self.dimensions.sorted(dimension_records.keys()):
2373 records = list(dimension_records[element].values())
2374 # Assume that if the record is already present that we can
2375 # use it without having to check that the record metadata
2376 # is consistent.
2377 self._registry.insertDimensionData(element, *records, skip_existing=True)
2379 def _import_grouped_refs(
2380 self,
2381 grouped_refs: defaultdict[_RefGroup, list[DatasetRef]],
2382 source_butler: LimitedButler | None,
2383 progress: Progress,
2384 *,
2385 dry_run: bool = False,
2386 expand_refs: bool = False,
2387 ) -> list[DatasetRef]:
2388 handled_collections: set[str] = set()
2389 n_to_import = 0
2390 all_imported_refs: list[DatasetRef] = []
2391 # Sort by run collection name to ensure Postgres takes locks in the
2392 # same order between different processes, to mitigate an issue
2393 # where Postgres can deadlock due to the unique index on collection
2394 # name. (See DM-47543).
2395 groups = sorted(grouped_refs.items(), key=lambda item: item[0].run)
2396 for (dimension_group, run), refs_to_import in progress.iter_item_chunks(
2397 groups, desc="Importing to registry by run and dataset type"
2398 ):
2399 if run not in handled_collections:
2400 # May need to create output collection. If source butler
2401 # has a registry, ask for documentation string.
2402 run_doc = None
2403 if source_butler is not None and (registry := getattr(source_butler, "registry", None)):
2404 run_doc = registry.getCollectionDocumentation(run)
2405 if not dry_run:
2406 registered = self.collections.register(run, doc=run_doc)
2407 else:
2408 registered = True
2409 handled_collections.add(run)
2410 if registered:
2411 _LOG.verbose("Creating output run %s", run)
2413 n_refs = len(refs_to_import)
2414 n_to_import += n_refs
2415 _LOG.verbose(
2416 "Importing %d ref%s with dimensions %s into run %s",
2417 n_refs,
2418 "" if n_refs == 1 else "s",
2419 dimension_group.names,
2420 run,
2421 )
2423 # Assume we are using UUIDs and the source refs will match
2424 # those imported.
2425 if not dry_run:
2426 imported_refs = self._registry._importDatasets(refs_to_import, expand=expand_refs)
2427 else:
2428 imported_refs = refs_to_import
2430 all_imported_refs.extend(imported_refs)
2432 assert n_to_import == len(all_imported_refs)
2433 _LOG.verbose("Imported %d datasets into destination butler", n_to_import)
2434 return all_imported_refs
2436 def transfer_from(
2437 self,
2438 source_butler: LimitedButler,
2439 source_refs: Iterable[DatasetRef],
2440 transfer: str = "auto",
2441 skip_missing: bool = True,
2442 register_dataset_types: bool = False,
2443 transfer_dimensions: bool = False,
2444 dry_run: bool = False,
2445 ) -> collections.abc.Collection[DatasetRef]:
2446 # Docstring inherited.
2447 source_refs = list(source_refs)
2448 if not self.isWriteable() and not dry_run: 2448 ↛ 2449line 2448 didn't jump to line 2449 because the condition on line 2448 was never true
2449 raise TypeError("Butler is read-only.")
2451 progress = Progress("lsst.daf.butler.Butler.transfer_from", level=VERBOSE)
2453 artifact_existence: dict[ResourcePath, bool] = {}
2454 file_transfer_source = source_butler._file_transfer_source
2455 transfer_records = retrieve_file_transfer_records(
2456 file_transfer_source, source_refs, artifact_existence
2457 )
2458 # In some situations the datastore artifact may be missing and we do
2459 # not want that registry entry to be imported. For example, this can
2460 # happen if a file was removed but the dataset was left in the registry
2461 # for provenance, or if a pipeline task didn't create all of the
2462 # possible files in a QuantumBackedButler.
2463 if skip_missing:
2464 original_ids = {ref.id for ref in source_refs}
2465 missing_ids = original_ids - transfer_records.keys()
2466 if missing_ids:
2467 original_count = len(source_refs)
2468 source_refs = [ref for ref in source_refs if ref.id not in missing_ids]
2469 filtered_count = len(source_refs)
2470 n_missing = original_count - filtered_count
2471 _LOG.verbose(
2472 "%d dataset%s removed because the artifact does not exist. Now have %d.",
2473 n_missing,
2474 "" if n_missing == 1 else "s",
2475 filtered_count,
2476 )
2478 import_info = self._prepare_for_import_refs(
2479 source_butler,
2480 source_refs,
2481 register_dataset_types=register_dataset_types,
2482 dry_run=dry_run,
2483 transfer_dimensions=transfer_dimensions,
2484 )
2486 # Do all the importing in a single transaction.
2487 with self.transaction():
2488 self._import_dimension_records(import_info.dimension_records, dry_run=dry_run)
2489 imported_refs = self._import_grouped_refs(
2490 import_info.grouped_refs, source_butler, progress, dry_run=dry_run
2491 )
2493 # Ask the datastore to transfer. The datastore has to check that
2494 # the source datastore is compatible with the target datastore.
2495 _LOG.verbose("Transferring %d datasets from %s", len(transfer_records), file_transfer_source.name)
2496 accepted, rejected = self._datastore.transfer_from(
2497 transfer_records,
2498 imported_refs,
2499 transfer=transfer,
2500 artifact_existence=artifact_existence,
2501 dry_run=dry_run,
2502 )
2503 if rejected: 2503 ↛ 2505line 2503 didn't jump to line 2505 because the condition on line 2503 was never true
2504 # For now, accept the registry entries but not the files.
2505 _LOG.warning(
2506 "%d datasets were rejected and %d accepted for transfer.",
2507 len(rejected),
2508 len(accepted),
2509 )
2511 return imported_refs
2513 def validateConfiguration(
2514 self,
2515 logFailures: bool = False,
2516 datasetTypeNames: Iterable[str] | None = None,
2517 ignore: Iterable[str] | None = None,
2518 ) -> None:
2519 # Docstring inherited.
2520 if datasetTypeNames:
2521 datasetTypes = [self.get_dataset_type(name) for name in datasetTypeNames]
2522 else:
2523 datasetTypes = list(self._registry.queryDatasetTypes())
2525 # filter out anything from the ignore list
2526 if ignore:
2527 ignore = set(ignore)
2528 datasetTypes = [
2529 e for e in datasetTypes if e.name not in ignore and e.nameAndComponent()[0] not in ignore
2530 ]
2531 else:
2532 ignore = set()
2534 # For each datasetType that has an instrument dimension, create
2535 # a DatasetRef for each defined instrument
2536 datasetRefs = []
2538 # Find all the registered instruments (if "instrument" is in the
2539 # universe).
2540 instruments: set[str] = set()
2541 if "instrument" in self.dimensions: 2541 ↛ 2561line 2541 didn't jump to line 2561 because the condition on line 2541 was always true
2542 instruments = {rec.name for rec in self.query_dimension_records("instrument", explain=False)}
2544 for datasetType in datasetTypes:
2545 if "instrument" in datasetType.dimensions: 2545 ↛ 2544line 2545 didn't jump to line 2544 because the condition on line 2545 was always true
2546 # In order to create a conforming dataset ref, create
2547 # fake DataCoordinate values for the non-instrument
2548 # dimensions. The type of the value does not matter here.
2549 dataId = {dim: 1 for dim in datasetType.dimensions.names if dim != "instrument"}
2551 for instrument in instruments:
2552 datasetRef = DatasetRef(
2553 datasetType,
2554 DataCoordinate.standardize(
2555 dataId, instrument=instrument, dimensions=datasetType.dimensions
2556 ),
2557 run="validate",
2558 )
2559 datasetRefs.append(datasetRef)
2561 entities: list[DatasetType | DatasetRef] = []
2562 entities.extend(datasetTypes)
2563 entities.extend(datasetRefs)
2565 datastoreErrorStr = None
2566 try:
2567 self._datastore.validateConfiguration(entities, logFailures=logFailures)
2568 except ValidationError as e:
2569 datastoreErrorStr = str(e)
2571 # Also check that the LookupKeys used by the datastores match
2572 # registry and storage class definitions
2573 keys = self._datastore.getLookupKeys()
2575 failedNames = set()
2576 failedDataId = set()
2577 for key in keys:
2578 if key.name is not None:
2579 if key.name in ignore:
2580 continue
2582 # skip if specific datasetType names were requested and this
2583 # name does not match
2584 if datasetTypeNames and key.name not in datasetTypeNames:
2585 continue
2587 # See if it is a StorageClass or a DatasetType
2588 if key.name in self.storageClasses:
2589 pass
2590 else:
2591 try:
2592 self.get_dataset_type(key.name)
2593 except KeyError:
2594 if logFailures: 2594 ↛ 2595line 2594 didn't jump to line 2595 because the condition on line 2594 was never true
2595 _LOG.critical(
2596 "Key '%s' does not correspond to a DatasetType or StorageClass", key
2597 )
2598 failedNames.add(key)
2599 else:
2600 # Dimensions are checked for consistency when the Butler
2601 # is created and rendezvoused with a universe.
2602 pass
2604 # Check that the instrument is a valid instrument
2605 # Currently only support instrument so check for that
2606 if key.dataId:
2607 dataIdKeys = set(key.dataId)
2608 if {"instrument"} != dataIdKeys: 2608 ↛ 2609line 2608 didn't jump to line 2609 because the condition on line 2608 was never true
2609 if logFailures:
2610 _LOG.critical("Key '%s' has unsupported DataId override", key)
2611 failedDataId.add(key)
2612 elif key.dataId["instrument"] not in instruments: 2612 ↛ 2613line 2612 didn't jump to line 2613 because the condition on line 2612 was never true
2613 if logFailures:
2614 _LOG.critical("Key '%s' has unknown instrument", key)
2615 failedDataId.add(key)
2617 messages = []
2619 if datastoreErrorStr: 2619 ↛ 2620line 2619 didn't jump to line 2620 because the condition on line 2619 was never true
2620 messages.append(datastoreErrorStr)
2622 for failed, msg in (
2623 (failedNames, "Keys without corresponding DatasetType or StorageClass entry: "),
2624 (failedDataId, "Keys with bad DataId entries: "),
2625 ):
2626 if failed:
2627 msg += ", ".join(str(k) for k in failed)
2628 messages.append(msg)
2630 if messages:
2631 raise ValidationError(";\n".join(messages))
2633 @property
2634 @deprecated(
2635 "Please use 'collections' instead. collection_chains will be removed after v28.",
2636 version="v28",
2637 category=FutureWarning,
2638 )
2639 def collection_chains(self) -> DirectButlerCollections:
2640 """Object with methods for modifying collection chains."""
2641 return DirectButlerCollections(self._registry)
2643 @property
2644 def collections(self) -> DirectButlerCollections:
2645 """Object with methods for modifying and inspecting collections."""
2646 return DirectButlerCollections(self._registry)
2648 @property
2649 def run(self) -> str | None:
2650 """Name of the run this butler writes outputs to by default (`str` or
2651 `None`).
2653 This is an alias for ``self.registry.defaults.run``. It cannot be set
2654 directly in isolation, but all defaults may be changed together by
2655 assigning a new `RegistryDefaults` instance to
2656 ``self.registry.defaults``.
2657 """
2658 return self._registry.defaults.run
2660 @property
2661 def registry(self) -> Registry:
2662 """The object that manages dataset metadata and relationships
2663 (`Registry`).
2665 Many operations that don't involve reading or writing butler datasets
2666 are accessible only via `Registry` methods. Eventually these methods
2667 will be replaced by equivalent `Butler` methods.
2668 """
2669 return RegistryShim(self)
2671 @property
2672 def dimensions(self) -> DimensionUniverse:
2673 # Docstring inherited.
2674 return self._registry.dimensions
2676 def query(self) -> contextlib.AbstractContextManager[Query]:
2677 # Docstring inherited.
2678 return self._registry._query()
2680 def _query_driver(
2681 self,
2682 default_collections: Iterable[str],
2683 default_data_id: DataCoordinate,
2684 ) -> contextlib.AbstractContextManager[DirectQueryDriver]:
2685 """Set up a QueryDriver instance for use with this Butler. Although
2686 this is marked as a private method, it is also used by Butler server.
2687 """
2688 return self._registry._query_driver(default_collections, default_data_id)
2690 @contextlib.contextmanager
2691 def _query_all_datasets_by_page(
2692 self, args: QueryAllDatasetsParameters
2693 ) -> Iterator[Iterator[list[DatasetRef]]]:
2694 with self.query() as query:
2695 pages = query_all_datasets(self, query, args)
2696 yield iter(page.data for page in pages)
2698 def _preload_cache(self, *, load_dimension_record_cache: bool = True) -> None:
2699 """Immediately load caches that are used for common operations."""
2700 self._registry.preload_cache(load_dimension_record_cache=load_dimension_record_cache)
2702 def _expand_data_ids(self, data_ids: Iterable[DataCoordinate]) -> list[DataCoordinate]:
2703 return self._registry.expand_data_ids(data_ids)
2705 _config: ButlerConfig
2706 """Configuration for this Butler instance."""
2708 _registry: SqlRegistry
2709 """The object that manages dataset metadata and relationships
2710 (`SqlRegistry`).
2712 Most operations that don't involve reading or writing butler datasets are
2713 accessible only via `SqlRegistry` methods.
2714 """
2716 storageClasses: StorageClassFactory
2717 """An object that maps known storage class names to objects that fully
2718 describe them (`StorageClassFactory`).
2719 """
2721 _closed: bool
2722 """`True` if close() has already been called on this instance; `False`
2723 otherwise.
2724 """
2727class _RefGroup(NamedTuple):
2728 """Key identifying a batch of DatasetRefs to be inserted in
2729 `Butler.transfer_from`.
2730 """
2732 dimensions: DimensionGroup
2733 run: str
2736class _ImportDatasetsInfo(NamedTuple):
2737 """Information extracted from datasets to be imported."""
2739 grouped_refs: defaultdict[_RefGroup, list[DatasetRef]]
2740 dimension_records: dict[DimensionElement, dict[DataCoordinate, DimensionRecord]]
2743def _to_uuid(id: DatasetId | str) -> uuid.UUID:
2744 if isinstance(id, uuid.UUID):
2745 return id
2746 else:
2747 return uuid.UUID(id)
2750class _ButlerClosed:
2751 def __getattr__(self, name: str) -> Any:
2752 raise RuntimeError("Attempted to use a Butler instance which has been closed.")
2755_BUTLER_CLOSED_INSTANCE: Any = _ButlerClosed()
2758def _retrieve_dataset_type(registry: SqlRegistry, name: str) -> DatasetType | None:
2759 """Return DatasetType defined in registry given dataset type name."""
2760 try:
2761 return registry.getDatasetType(name)
2762 except MissingDatasetTypeError:
2763 return None