Class to perform parallel execution of a Sequence evaluation.
| 35 | |
| 36 | |
| 37 | class ParallelExecutionEngine(ExecutionEngine): |
| 38 | """ |
| 39 | Class to perform parallel execution of a Sequence evaluation. |
| 40 | """ |
| 41 | |
| 42 | def __init__(self, processes=None, partition_size=None): |
| 43 | """ |
| 44 | Set the number of processes for parallel execution. |
| 45 | :param processes: Number of parallel Processes |
| 46 | """ |
| 47 | super(ParallelExecutionEngine, self).__init__() |
| 48 | self.processes = processes |
| 49 | self.partition_size = partition_size |
| 50 | |
| 51 | def evaluate(self, sequence, transformations): |
| 52 | """ |
| 53 | Execute the sequence of transformations in parallel |
| 54 | :param sequence: Sequence to evaluation |
| 55 | :param transformations: Transformations to apply |
| 56 | :return: Resulting sequence or value |
| 57 | """ |
| 58 | result = sequence |
| 59 | parallel = partial( |
| 60 | parallelize, processes=self.processes, partition_size=self.partition_size |
| 61 | ) |
| 62 | staged = [] |
| 63 | for transform in transformations: |
| 64 | strategies = transform.execution_strategies or {} |
| 65 | if ExecutionStrategies.PARALLEL in strategies: |
| 66 | staged.insert(0, transform.function) |
| 67 | else: |
| 68 | if staged: |
| 69 | result = parallel(compose(*staged), result) |
| 70 | staged = [] |
| 71 | if ExecutionStrategies.PRE_COMPUTE in strategies: |
| 72 | result = list(result) |
| 73 | result = transform.function(result) |
| 74 | if staged: |
| 75 | result = parallel(compose(*staged), result) |
| 76 | return iter(result) |