Skip to content

openstb.simulator.utils.reduction

Reduction tree utilities.

Reduction trees are a common pattern in parallel processing to combine results from multiple tasks. At each level of the tree, a set of results is aggregated to a single result. This is then passed to the next level of the tree. For addition, instead of computing the overall result as

result = r1 + r2 + r3 + r4 + r5 + r6 + r7 + r8

it might be computed as

result = ((r1 + r2) + (r3 + r4)) + ((r5 + r6) + (r7 + r8))

This reduces the number of tasks that the scheduler has to track, reducing both scheduler overhead and the amount of memory required to store intermediate results.

In this module, the size of the tree is parametrised by the number of levels (3 in the above example) and the number of results aggregared at each level (2 in the above example). This leads to a reduction factor which is calculated as the number of results per level raised to the power of the number of levels. In the above example this is 2^3 = 8, meaning that 8 inputs to the tree are reduced to a single output.

Classes:

Name Description
DaskReductionTree

Dask-based reduction tree to aggregate multiple futures to one.

DaskReductionTree

DaskReductionTree(client, output_func, reduce_func, reduce_args=None, reduce_kwargs=None, levels=3, futures=4)

Dask-based reduction tree to aggregate multiple futures to one.

This is given a reduction function which is scheduled to run on the cluster. Each output from the tree is passed to the given output function. Note that the output function is called directly, not scheduled to run on the cluster.

Note that no guarantee is made about the order the input futures are reduced in.

Parameters:

Name Type Description Default
client Client

The Dask client to submit reduction tasks to.

required
output_func Callable[[Future[T], Any], None]

The function to call when a reduction is complete. This will be given the future corresponding to the output of the final reduction and the current value of the tag property.

required
reduce_func Callable[Concatenate[list[T], ...], T]

The function to call to reduce some data. This will be scheduled through the Dask client. It will be passed a list of the results from the input futures and must returned the reduced result.

required
reduce_args Iterable[Any] | None

Positional arguments to pass to reduce_func.

None
reduce_kwargs dict[str, Any] | None

Keyword arguments to pass to reduce_func.

None
levels int

The number of levels of reduction to perform before calling output_func.

3
futures int

The number of futures to reduce at each level.

4

Methods:

Name Description
add_futures

Add futures to the reduction tree.

flush

Flush the reduction tree.

Attributes:

Name Type Description
tag Any

Tag associated with the data being reduced currently.

Source code in openstb/simulator/utils/reduction.py
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
def __init__(
    self,
    client: distributed.Client,
    output_func: Callable[[distributed.Future[T], Any], None],
    reduce_func: Callable[Concatenate[list[T], ...], T],
    reduce_args: Iterable[Any] | None = None,
    reduce_kwargs: dict[str, Any] | None = None,
    levels: int = 3,
    futures: int = 4,
):
    """
    Parameters
    ----------
    client
        The Dask client to submit reduction tasks to.
    output_func
        The function to call when a reduction is complete. This will be given the
        future corresponding to the output of the final reduction and the current
        value of the tag property.
    reduce_func
        The function to call to reduce some data. This will be scheduled through the
        Dask client. It will be passed a list of the results from the input futures
        and must returned the reduced result.
    reduce_args
        Positional arguments to pass to `reduce_func`.
    reduce_kwargs
        Keyword arguments to pass to `reduce_func`.
    levels
        The number of levels of reduction to perform before calling `output_func`.
    futures
        The number of futures to reduce at each level.

    """
    self.client = client
    self.output_func = output_func
    self.reduce_func = reduce_func
    self.reduce_args = reduce_args or []
    self.reduce_kwargs = reduce_kwargs or {}
    self.levels = levels
    self.futures = futures
    self._tag = None

    # Prepare storage for futures that have not been fully reduced.
    self._reducing: list[list[distributed.Future[T]]] = []
    for i in range(self.levels):
        self._reducing.append([])

tag property writable

tag

Tag associated with the data being reduced currently.

The current tag is passed to the output function. It can be used as metadata for the current calculations.

Note that changing the tag will result in an exception if there are pending reductions. It is recommended that you call the flush method prior to changing the tag.

add_futures

add_futures(*futures)

Add futures to the reduction tree.

Parameters:

Name Type Description Default
*futures Future[T]

Futures to reduce.

()

Returns:

Type Description
list[Future]

Any futures submitted to the cluster to perform the reduction.

Source code in openstb/simulator/utils/reduction.py
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
def add_futures(
    self, *futures: distributed.Future[T]
) -> list[distributed.Future[T]]:
    """Add futures to the reduction tree.

    Parameters
    ----------
    *futures
        Futures to reduce.

    Returns
    -------
    list[distributed.Future]
        Any futures submitted to the cluster to perform the reduction.

    """
    added = []

    # Insert at the first level.
    self._reducing[0].extend(futures)

    # Process each level which has enough futures to reduce.
    for level in range(self.levels):
        while len(self._reducing[level]) >= self.futures:
            # Apply the reduction function to the first N futures.
            to_reduce = self._reducing[level][: self.futures]
            reduced = self.client.submit(
                self.reduce_func,
                to_reduce,
                *self.reduce_args,
                key=f"reduction-{level}-{tokenize(to_reduce)}",
                **self.reduce_kwargs,
            )
            added.append(reduced)
            del self._reducing[level][: self.futures]

            # Output or add to the next level.
            if level == self.levels - 1:
                self.output_func(reduced, self._tag)
            else:
                self._reducing[level + 1].append(reduced)

    return added

flush

flush()

Flush the reduction tree.

In most situations, the number of futures to reduce won't be an exact multiple of the reduction factor. Calling this method will reduce any futures remaining in the tree. There are three cases:

  • If the tree is empty, do nothing.
  • If there is one future in the entire tree, output that.
  • If there are multiple futures in the tree, call the reduction function on all of them and output the result.

After this method has been called, the tree is guaranteed to be empty.

Returns:

Type Description
list[Future]

Any futures submitted to the cluster to perform the reduction.

Source code in openstb/simulator/utils/reduction.py
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
def flush(self) -> list[distributed.Future[T]]:
    """Flush the reduction tree.

    In most situations, the number of futures to reduce won't be an exact multiple
    of the reduction factor. Calling this method will reduce any futures remaining
    in the tree. There are three cases:

    * If the tree is empty, do nothing.
    * If there is one future in the entire tree, output that.
    * If there are multiple futures in the tree, call the reduction function on all
      of them and output the result.

    After this method has been called, the tree is guaranteed to be empty.

    Returns
    -------
    list[distributed.Future]
        Any futures submitted to the cluster to perform the reduction.

    """
    # Collect all remaining futures at any level.
    remaining = []
    for level in range(self.levels):
        remaining.extend(self._reducing[level])
        self._reducing[level].clear()

    if not remaining:
        return []

    if len(remaining) == 1:
        self.output_func(remaining[0], self._tag)
        return []

    # Schedule the reduction and output the future.
    reduced = self.client.submit(
        self.reduce_func,
        remaining,
        *self.reduce_args,
        key=f"reduction-finish-{tokenize(remaining)}",
        **self.reduce_kwargs,
    )
    self.output_func(reduced, self._tag)
    return [reduced]