Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); GitHub - RBigData/tasktools: Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart. · GitHub
Skip to content

Repository files navigation

tasktools

  • Version: 0.2-0
  • Status:Build Status
  • License:BSD 2-Clause
  • Author: Drew Schmidt

Tools for task-based parallelism with MPI via pbdMPI. Currently we provide these basic functions:

  1. mpi_napply() --- a distributed lapply() that operates on an integer sequence. Supports checkpoint/restart and non-prescheduled workloads.
  2. mpi_lapply() --- a fully general, distributed lapply().

These functions are conceptually similar to pbdLapply() from the pbdMPI package, but with some key differences. The pbdMPI functions have more modes of operation, allowing for different kinds of distributions of the inputs for the more general mpi_lapply(). And naturally, the pbdMPI functions do not handle checkpoint/restart.

In addition to these "ply" functions, also offer mpi_progress() to check on the status of running jobs.

Installation

You can install the stable version from the HPCRAN using the usual install.packages():

install.packages("tasktools", repos=c("https://hpcran.org", "https://cran.rstudio.com"))

The development version is maintained on GitHub:

remotes::install_github("RBigData/tasktools")

Examples

Complete source code for all of these examples can be found in the inst/examples directory of the tasktools source tree. Here we'll take a look at them in pieces. Throughout, we'll use a (fake) "expensive" function for our evaluations:

costly=function(x, waittime)
{
Sys.sleep(waittime)
print(paste("iteration:", x))
sqrt(x)
}

We can run a checkpointed lapply() in serial via crlapply() from the crlapply package:

ret=crlapply::crlapply(1:10, costly, FILE="/tmp/cr.rdata", waittime=0.5)
unlist(ret)

If we save this source to the file crlapply.r. We can run it and kill it a few times to show its effectiveness:

$ r crlapply.r [1] "iteration: 1"
[1] "iteration: 2"
[1] "iteration: 3"
^C
$ r crlapply.r [1] "iteration: 4"
[1] "iteration: 5"
[1] "iteration: 6"
[1] "iteration: 7"
^C
$ r crlapply.r [1] "iteration: 8"
[1] "iteration: 9"
[1] "iteration: 10"
[1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

Since we are operating on the integer sequence of values 1 to 10, we can easily parallelize this, even distributing the work across multiple nodes, with mpi_napply():

ret= mpi_napply(10, costly, checkpoint_path="/tmp", waittime=1)
comm.print(unlist(ret))

To see exactly what happens during execution, we modify the printing in the "costly" function to be:

cat(paste("iter", i, "executed on rank", comm.rank(), "\n"))

Let's run this with 3 MPI ranks. We can again run and kill it a few times to demonstrate the checkpointing:

$ mpirun -np 3 r mpi_napply.r iter 4 executed on rank 1 iter 7 executed on rank 2 iter 1 executed on rank 0 ^Citer 2 executed on rank 0 iter 8 executed on rank 2 iter 5 executed on rank 1 $ mpirun -np 3 r mpi_napply.r iter 9 executed on rank 2 iter 3 executed on rank 0 iter 6 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

There is also a non-prescheduling variant. This can be useful if there is a lot of variance among function evaluation for the inputs, and you want the values to be executed on a "first come, first serve" basis. All we have to do is set preschedule=FALSE:

ret= mpi_napply(10, costly, preschedule=FALSE, waittime=1)
comm.print(unlist(ret))

Now, it's worth noting that in this case, rank 0 behaves as the manager, doling out work. So it is not used in computation:

iter 1 executed on rank 1 iter 2 executed on rank 2 iter 3 executed on rank 1 iter 4 executed on rank 2 iter 5 executed on rank 1 iter 6 executed on rank 2 iter 7 executed on rank 1 iter 8 executed on rank 2 iter 9 executed on rank 1 iter 10 executed on rank 2 [1] 1.000000 1.414214 1.732051 2.000000 2.236068 2.449490 2.645751 2.828427
[9] 3.000000 3.162278

This too supports checkpointing, but hopefully how that works is clear.

Progress Bar

We also support a kind of progress bar, but it's definitely not what you're thinking. Let's start with an example similar to the one above:

suppressMessages(library(tasktools))
f=function(i) {print(i); Sys.sleep(1); sqrt(i)}
ignore= mpi_napply(20, f, checkpoint_path="/tmp")
finalize()

We can put these into the file slow_sqrt.r and run it for a bit before manually killing it with Ctrl+c:

$ mpirun -np 3 Rscript slow_sqrt.r
[1] 14
[1] 7
[1] 1
[1] 15
^C[1] 16
[1] 8
[1] 2

We can check the progress by invoking mpi_progress():

$ Rscript -e "tasktools::mpi_progress('/tmp')"## [=================---------------------------------] (7/20)

The progress bar works by scanning the checkpoint files, so we don't actually have to kill the tasks to run the progress bar bit (and in fact for a real workflow, you wouldn't want to). But for the sake of demonstration, this is much simpler.

The above example was run with the default preschedule=TRUE, but it will also work if we have preschedule=FALSE. There are some caveats to the progress bar, however. Please carefully check the ?tasktools::mpi_progress documentation.

About

Tools for task-based parallelism with MPI via pbdMPI with automatic checkpoint/restart.

Topics

Resources

Stars

2 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages