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

Class CreateOverlappingPartitions

dask/dataframe/dask_expr/_expr.py:1030–1129  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

1028
1029
1030class CreateOverlappingPartitions(Expr):
1031 _parameters = ["frame", "before", "after"]
1032
1033 @functools.cached_property
1034 def _meta(self):
1035 return self.frame._meta
1036
1037 def _divisions(self):
1038 # Keep divisions alive, MapPartitions will handle the actual division logic
1039 return self.frame.divisions
1040
1041 def _layer(self) -> dict:
1042 dsk, prevs, nexts = {}, [], [] # type: ignore
1043
1044 name_prepend = "overlap-prepend" + self.frame._name
1045 if self.before:
1046 prevs.append(None)
1047 if isinstance(self.before, numbers.Integral):
1048 before = self.before
1049 for i in range(self.frame.npartitions - 1):
1050 dsk[(name_prepend, i)] = (M.tail, (self.frame._name, i), before)
1051 prevs.append((name_prepend, i))
1052 elif isinstance(self.before, datetime.timedelta):
1053 # Assumes monotonic (increasing?) index
1054 divs = pd.Series(self.frame.divisions)
1055 deltas = divs.diff().iloc[1:-1]
1056
1057 # In the first case window-size is larger than at least one partition, thus it is
1058 # necessary to calculate how many partitions must be used for each rolling task.
1059 # Otherwise, these calculations can be skipped (faster)
1060
1061 if (self.before > deltas).any():
1062 pt_z = divs[0]
1063 for i in range(self.frame.npartitions - 1):
1064 # Select all indexes of relevant partitions between the current partition and
1065 # the partition with the highest division outside the rolling window (before)
1066 pt_i = divs[i + 1]
1067
1068 # lower-bound the search to the first division
1069 lb = max(pt_i - self.before, pt_z)
1070
1071 first, j = divs[i], i
1072 while first > lb and j > 0:
1073 first = first - deltas[j]
1074 j = j - 1
1075
1076 dsk[(name_prepend, i)] = ( # type: ignore
1077 _tail_timedelta,
1078 (self.frame._name, i + 1),
1079 [(self.frame._name, k) for k in range(j, i + 1)],
1080 self.before,
1081 )
1082 prevs.append((name_prepend, i))
1083 else:
1084 for i in range(self.frame.npartitions - 1):
1085 dsk[(name_prepend, i)] = ( # type: ignore
1086 _tail_timedelta,
1087 (self.frame._name, i + 1),

Callers 1

_lowerMethod · 0.85

Calls

no outgoing calls

Tested by

no test coverage detected