| 1028 | |
| 1029 | |
| 1030 | class 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), |