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

Method _lower

dask/dataframe/dask_expr/_expr.py:3524–3563  ·  view source on GitHub ↗
(self)

Source from the content-addressed store, hash-verified

3522 return [op for op in dfs if not is_broadcastable(dfs, op)]
3523
3524 def _lower(self):
3525 # This can be expensive when something that has expensive division
3526 # calculation is in the Expression
3527 dfs = self.args
3528 if (
3529 len(dfs) == 1
3530 or all(
3531 dfs[0].divisions == df.divisions and df.known_divisions for df in dfs
3532 )
3533 or len(self.divisions) == 2
3534 and max(map(lambda x: len(x.divisions), dfs)) == 2
3535 ):
3536 return self._expr_cls(*self.operands)
3537 elif self.divisions[0] is None:
3538 # We have to shuffle
3539 npartitions = max(df.npartitions for df in dfs)
3540 dtypes = {df._meta.index.dtype for df in dfs}
3541 if not _are_dtypes_shuffle_compatible(dtypes):
3542 raise TypeError(
3543 "DataFrames are not aligned. We need to shuffle to align partitions "
3544 "with each other. This is not possible because the indexes of the "
3545 f"DataFrames have differing dtypes={dtypes}. Please ensure that "
3546 "all Indexes have the same dtype or align manually for this to "
3547 "work."
3548 )
3549
3550 from dask.dataframe.dask_expr._shuffle import RearrangeByColumn
3551
3552 args = [
3553 (
3554 RearrangeByColumn(df, None, npartitions, index_shuffle=True)
3555 if isinstance(df, Expr)
3556 else df
3557 )
3558 for df in self.operands
3559 ]
3560 return self._expr_cls(*args)
3561
3562 args = maybe_align_partitions(*self.operands, divisions=self.divisions)
3563 return self._expr_cls(*args)
3564
3565 @functools.cached_property
3566 def _meta(self):

Callers

nothing calls this directly

Calls 5

RearrangeByColumnClass · 0.90
allFunction · 0.85
maxFunction · 0.85
maybe_align_partitionsFunction · 0.85

Tested by

no test coverage detected