MCPcopy Create free account
hub / github.com/dask/dask / HLGExpr

Class HLGExpr

dask/_expr.py:973–1081  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

971
972
973class HLGExpr(Expr):
974 _parameters = [
975 "dsk",
976 "low_level_optimizer",
977 "output_keys",
978 "postcompute",
979 "_cached_optimized",
980 ]
981 _defaults = {
982 "low_level_optimizer": None,
983 "output_keys": None,
984 "postcompute": None,
985 "_cached_optimized": None,
986 }
987
988 @property
989 def hlg(self):
990 return self.operand("dsk")
991
992 @staticmethod
993 def from_collection(collection, optimize_graph=True):
994 from dask.highlevelgraph import HighLevelGraph
995
996 if hasattr(collection, "dask"):
997 dsk = collection.dask.copy()
998 else:
999 dsk = collection.__dask_graph__()
1000
1001 # Delayed objects still ship with low level graphs as `dask` when going
1002 # through optimize / persist
1003 if not isinstance(dsk, HighLevelGraph):
1004
1005 dsk = HighLevelGraph.from_collections(
1006 str(id(collection)), dsk, dependencies=()
1007 )
1008 if optimize_graph and not hasattr(collection, "__dask_optimize__"):
1009 warnings.warn(
1010 f"Collection {type(collection)} does not define a "
1011 "`__dask_optimize__` method. In the future this will raise. "
1012 "If no optimization is desired, please set this to `None`.",
1013 PendingDeprecationWarning,
1014 )
1015 low_level_optimizer = None
1016 else:
1017 low_level_optimizer = (
1018 collection.__dask_optimize__ if optimize_graph else None
1019 )
1020 return HLGExpr(
1021 dsk=dsk,
1022 low_level_optimizer=low_level_optimizer,
1023 output_keys=collection.__dask_keys__(),
1024 postcompute=collection.__dask_postcompute__(),
1025 )
1026
1027 def finalize_compute(self):
1028 return HLGFinalizeCompute(
1029 self,
1030 low_level_optimizer=self.low_level_optimizer,

Calls

no outgoing calls