Source code for pycif.plugins.obsoperators.standard.transforms.period_pipe
from logging import info
from .utils.default_subsimus import default_subsimus
from .utils.fwd_pipe import init_pipe_links, compute_pipe_order
[docs]
def period_pipe(self, all_transforms, mapper):
"""Arrange all transforms into ordered forward and adjoint execution pipes.
Determines the chronologically correct execution order for every
``(transform, sub-simulation date)`` pair by:
1. Propagating sub-simulation periods from each transform to its
precursors and successors via
:func:`~.utils.default_subsimus.default_subsimus`.
2. Building the graph of links between transforms in forward and
adjoint directions with :func:`~.utils.fwd_pipe.init_pipe_links`.
3. Without Dask, walking these graphs to get the sequential execution
order with :func:`~.utils.fwd_pipe.compute_pipe_order`.
Each returned pipe is a list of ``(date, transform_id, direction)``
tuples, where ``direction`` is either ``'forward'`` or ``'adjoint'``
and controls whether a transform runs in its normal or dry-run mode.
Args:
self (ObsOperator): the obs-operator plugin instance.
all_transforms: the :class:`~pycif.utils.classes.transforms.Transform`
object holding all initialized transforms.
mapper (dict): the pipeline mapper dictionary mapping transform IDs
to their sub-simulation, input/output and
precursor/successor metadata.
Returns:
tuple[list, list]: ``(pipe_fwd, pipe_adj)`` where each element is a
list of ``(datetime.datetime, str, str)`` tuples giving the execution
order for forward and adjoint runs respectively.
"""
info("Computing the optimal order of transformation. This can take a while")
# First update subsimulations from precursors and successors
default_subsimus(all_transforms, mapper)
for mode in ["forward", "adjoint"]:
info(f"Doing {mode} order")
pipe_links, transforms_ids = init_pipe_links(
self, all_transforms, mapper, mode=mode)
# Dask handles the order itself, but batch computation
# needs the forward order
if self.use_dask and not (
mode == "forward" and hasattr(self, "batch_computation")):
continue
compute_pipe_order(
self, all_transforms, mapper, pipe_links, transforms_ids,
mode=mode)