knitr::opts_chunk$set( comment = "#", prompt = FALSE, tidy = FALSE, cache = FALSE, collapse = TRUE ) old <- options(width = 100L)
A pipeline step can contain another pipeline: instead of computing a result directly, the step builds an inner pipeline and returns it. This allows to reuse a standard analysis pipeline inside a larger workflow.
We start with the same pipeline as in the previous "Split, map, and reduce" vignette, fitting a linear model and returning its coefficients. This will serve as our inner pipeline.
library(pipeflow) inner <- pip_new("coefficients") |> pip_add("data", \(data = NULL) data) |> pip_add( "fit", \(data = ~data, xVar = "x", yVar = "y") { lm(paste(yVar, "~", xVar), data = data) } ) |> pip_add("coefs", \(fit = ~fit) coefficients(fit)) inner
The outer pipeline splits the data into subsets and derives the model coefficients by running the inner pipeline for each split.
# Helper to run inner pipeline run_inner_pip <- function(pip, name, data) { pip$name <- sprintf("coefs for *%s*", name) # Set data subset for inner pipeline and run it pip_set_params(pip, list(data = data)) |> pip_run() pip[["coefs", "out"]] } outer <- pip_new("full analysis") |> pip_add("data", \(data = NULL) data) |> pip_add( "split_data", \(data = ~data, byVar = "by") { split(data, f = data[[byVar]]) } ) |> pip_add( "inner_run", \(dataList = ~split_data, xVar = "x", yVar = "y") { p <- pip_clone(inner) # Forward parameters to inner pip_set_params(p, list(xVar = xVar, yVar = yVar)) Map( f = run_inner_pip, name = names(dataList), data = dataList, MoreArgs = list(pip = p) ) } ) |> pip_add( "combine", \(coefs = ~inner_run) as.data.frame(do.call(rbind, coefs)) ) outer
Note the inner_run step: its parameters xVar and yVar are forwarded
to the inner pipeline.
Let's now set the analysis parameters and run the full pipeline:
outer |> pip_set_params( list( data = iris, xVar = "Sepal.Length", yVar = "Sepal.Width", byVar = "Species" ) ) |> pip_run()
The output of the inner_run step is a list of coefficient vectors, one
for each species,
outer[["inner_run", "out"]]
and the combine step returns the expected combined table.
outer[["combine", "out"]]
Now suppose we want to change one of the model settings, say use
Petal.Length instead of Sepal.Length as the predictor.
pip_set_params(outer, params = list(xVar = "Petal.Length")) outer
Since the xVar parameter is part of the inner_run step's function arguments,
the inner_run step's state (and its downstream dependencies) correctly has
now been marked as "outdated".
While forwarding the parameters manually to the inner pipeline is straight-forward in our toy example, trying to manually synchronize real-world parameter sets between the outer and inner pipeline quickly becomes a unfeasible and a source for bugs that are hard to detect.
For this reason, in practice, the following pattern should be used, which
basically just forwards the combined set of all existing parameters.
To do this, we replace the inner_run step as follows:
outer |> pip_replace( "inner_run", \(dataList = ~split_data, ...) { p <- pip_clone(inner) # Forward all parameters (from outer and inner) all_params <- .self$get_params() pip_set_params(p, all_params) Map( f = run_inner_pip, name = names(dataList), data = dataList, MoreArgs = list(pip = p) ) }, params = pip_get_params(inner) # <--- default parameters of inner pipeline ) outer
In this version, the inner_run step no longer declares xVar and yVar as
its own arguments. Instead, its parameters are seeded with the inner
pipeline's parameters via params = pip_get_params(inner), and the step
forwards the combined parameter set at run time. Three aspects are worth
spelling out:
.self refers to the pipeline that is currently being run and is
available inside every step function without having to declare it. It
exposes the usual pipeline interface, so .self$get_params() returns the
current unbound parameters of the outer pipeline, and things like
.self$name or .self$run() work as expected. For more examples on this
mechanism see the
self-modifying pipelines vignette.params = pip_get_params(inner) registers the inner pipeline's unbound
parameters (data, xVar, yVar) as parameters of the inner_run step.
They therefore show up in the step's params column, become part of the
outer pipeline's parameter set (and can be updated with
pip_set_params(outer, ...)), and are passed to the step function when it
runs.params overlaps with the
arguments of the step function, the function's default values take
precedence. This is why the replacement function above declares only
dataList and ... — had it declared e.g. xVar = "x", that default would
win over the value coming from params, and the forwarding would silently
be ignored. As another example,
pip_add("s", \(x = 99, ...) x, params = list(x = 1)) stores x = 99,
not x = 1.Let's re-run the full pipeline.
outer |> pip_set_params( list( data = iris, xVar = "Sepal.Length", yVar = "Sepal.Width", byVar = "Species" ) ) |> pip_run()
Note that the warning in the log above is expected and harmless. Basically,
.self$get_params() returns the outer pipeline's entire parameter set,
which includes byVar (from the split_data step) while the inner
pipeline only defines data, xVar, and yVar. Since pip_set_params()
reports any parameters that are not defined in the target and simply leaves
them unset, the inner pipeline still receives all parameters it knows and the
result is unaffected.
If you need to omit the warning (e.g. in production code), just suppress it:
suppressWarnings(pip_set_params(p, all_params))
Alternatively, you could first restrict the forwarded parameters to those the inner pipeline actually knows. That has the same effect but adds code, and since forwarding the full parameter set is the whole point of this pattern, there is no need for it.
Again, changing one of the inner parameters will correctly outdate the
inner_run step plus downstream dependencies.
pip_set_params(outer, params = list(xVar = "Petal.Length")) outer
With the above pattern, you can now change both the inner and outer pipeline, adding and/or removing any steps or parameters, without having to worry about parameter synchronization.
Both this and the previous vignette solve the
same "split, apply, and combine" problem. The built-in execution modes
(exec = "split"/"reduce") are the recommended default whenever they fit
while nested pipelines can be considered the more general tool.
As a rule of thumb:
options(old)
Any scripts or data that you put into this service are public.
Add the following code to your website.
For more information on customizing the embed code, read Embedding Snippets.