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

Method _layer

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

Source from the content-addressed store, hash-verified

1043 return self.frame.divisions
1044
1045 def _layer(self) -> dict:
1046 dsk, prevs, nexts = {}, [], [] # type: ignore[var-annotated]
1047
1048 name_prepend = f"overlap-prepend-{self._name}"
1049 if self.before:
1050 prevs.append(None)
1051 if isinstance(self.before, numbers.Integral):
1052 before = self.before
1053 for i in range(self.frame.npartitions - 1):
1054 dsk[(name_prepend, i)] = (M.tail, (self.frame._name, i), before)
1055 prevs.append((name_prepend, i))
1056 elif isinstance(self.before, datetime.timedelta):
1057 # Assumes monotonic (increasing?) index
1058 divs = pd.Series(self.frame.divisions)
1059 deltas = divs.diff().iloc[1:-1]
1060
1061 # In the first case window-size is larger than at least one partition, thus it is
1062 # necessary to calculate how many partitions must be used for each rolling task.
1063 # Otherwise, these calculations can be skipped (faster)
1064
1065 if (self.before > deltas).any():
1066 pt_z = divs[0]
1067 for i in range(self.frame.npartitions - 1):
1068 # Select all indexes of relevant partitions between the current partition and
1069 # the partition with the highest division outside the rolling window (before)
1070 pt_i = divs[i + 1]
1071
1072 # lower-bound the search to the first division
1073 lb = max(pt_i - self.before, pt_z)
1074
1075 first, j = divs[i], i
1076 while first > lb and j > 0:
1077 first = first - deltas[j]
1078 j = j - 1
1079
1080 dsk[(name_prepend, i)] = ( # type: ignore[assignment]
1081 _tail_timedelta,
1082 (self.frame._name, i + 1),
1083 [(self.frame._name, k) for k in range(j, i + 1)],
1084 self.before,
1085 )
1086 prevs.append((name_prepend, i))
1087 else:
1088 for i in range(self.frame.npartitions - 1):
1089 dsk[(name_prepend, i)] = ( # type: ignore[assignment]
1090 _tail_timedelta,
1091 (self.frame._name, i + 1),
1092 [(self.frame._name, i)],
1093 self.before,
1094 )
1095 prevs.append((name_prepend, i))
1096 else:
1097 prevs.extend([None] * self.frame.npartitions) # type: ignore[list-item]
1098
1099 name_append = f"overlap-append-{self._name}"
1100 if self.after:
1101 if isinstance(self.after, numbers.Integral):
1102 after = self.after

Callers

nothing calls this directly

Calls 3

maxFunction · 0.85
diffMethod · 0.80
anyMethod · 0.45

Tested by

no test coverage detected