Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion DESCRIPTION
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
Package: filearray
Type: Package
Title: File-Backed Array for Out-of-Memory Computation
Version: 0.2.3.1
Version: 0.2.3.3
Language: en-US
Encoding: UTF-8
License: LGPL-3
Expand Down
1 change: 1 addition & 0 deletions NAMESPACE
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ export(fmap_element_wise)
export(fwhich)
export(mapreduce)
export(typeof)
export(with_filelock)
exportClasses(FileArray)
exportClasses(FileArrayProxy)
exportMethods(apply)
Expand Down
10 changes: 10 additions & 0 deletions NEWS.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,16 @@

* Partition files are written to a temporary file and then renamed into place, so other processes loading the same array no longer fail with "File size too small" while a partition is being created or filled
* `filearray_load` only treats files named `<n>.farr` as partitions
* Writes lock the partitions they change, so several processes (for example `parallel`, `callr` or `mirai` workers, or several applications) can write to one array without losing each other's data: sub-assignment (`x[...] <- value`), `fill_partition`, `initialize_partition`, `set_partition`, `fmap`, `fmap_element_wise` and `filearray_bind` wait while another process holds one of their partitions. Reading takes no locks, except that functions that must first create missing partitions (such as `max` or `sum`) lock the partitions they create
* Added `lock_partition(parts, timeout)` and `unlock_partition(parts)` methods to hold partitions over several operations (for example to check whether a partition has been computed and write it only if not); `lock_partition` waits and retries until the partitions are free or `timeout` seconds have passed
* Locks are byte-range locks on an empty file `.filearray.lock` in the array directory, released by the operating system when a process exits. `options(filearray.lock = FALSE)` (environment variable `FILEARRAY_LOCK`) turns off automatic locking; `options(filearray.lock.timeout = <seconds>)` (`FILEARRAY_LOCK_TIMEOUT`) limits how long writes wait, by default without limit
* A missing partition is created only if no other process has created it in the meantime, and every partition write uses its own temporary file, so concurrent writers no longer replace each other's partitions or share a temporary file
* `fmap` stops when the mapping function raises an error or is interrupted, instead of warning "cannot finish map" and returning a partly written result
* `delete()` releases the locks the session holds on the array before removing it
* Added `with_filelock(expr, lockfile, timeout, on_unsupported)` to evaluate code while holding an exclusive lock on any lock file, so that only one R session at a time runs it; other sessions wait. It is not tied to file arrays, so other packages can use it. `on_unsupported` chooses what happens when the lock cannot be taken: warn (default), stop, or ignore
* Partition locks never make a write fail when the lock itself cannot be taken: any error while locking (for example a stale network file handle or too many open files, not only file systems without lock support) now warns once and the write proceeds, and an invalid `filearray.lock.timeout` option or `FILEARRAY_LOCK_TIMEOUT` value warns once and is ignored instead of stopping every write
* Fixed `mapreduce` (and so `sum`, `max`, `min`, `range` and `fwhich`) reading only the first slice of every partition after the first when `partition_size` is greater than 1, which gave wrong results and wrong `fwhich` indices
* `fmap2` reports the error raised by the mapping function instead of "Unknown error."

# filearray 0.2.3

Expand Down
36 changes: 36 additions & 0 deletions R/RcppExports.R
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,42 @@ FARR_subset2 <- function(filebase, listOrEnv, reshape = NULL, drop = FALSE, use_
.Call(`_filearray_FARR_subset2`, filebase, listOrEnv, reshape, drop, use_dimnames, thread_buffer, split_dim, strict)
}

FARR_lock_acquire <- function(lock_path, parts, scope_id) {
.Call(`_filearray_FARR_lock_acquire`, lock_path, parts, scope_id)
}

FARR_lock_release_scope <- function(scope_id) {
.Call(`_filearray_FARR_lock_release_scope`, scope_id)
}

FARR_lock_commit_scope <- function(scope_id) {
.Call(`_filearray_FARR_lock_commit_scope`, scope_id)
}

FARR_lock_release_explicit <- function(lock_path, parts, all = FALSE) {
.Call(`_filearray_FARR_lock_release_explicit`, lock_path, parts, all)
}

FARR_lock_forget <- function(lock_path) {
.Call(`_filearray_FARR_lock_forget`, lock_path)
}

FARR_lock_release_all <- function() {
.Call(`_filearray_FARR_lock_release_all`)
}

FARR_link_noreplace <- function(from, to) {
.Call(`_filearray_FARR_link_noreplace`, from, to)
}

FARR_lock_status <- function(lock_path) {
.Call(`_filearray_FARR_lock_status`, lock_path)
}

FARR_lock_registry <- function() {
.Call(`_filearray_FARR_lock_registry`)
}

FARR_buffer_map <- function(input_filebases, output_filebase, map, buffer_nelems, result_nelems = 0L) {
.Call(`_filearray_FARR_buffer_map`, input_filebases, output_filebase, map, buffer_nelems, result_nelems)
}
Expand Down
7 changes: 6 additions & 1 deletion R/bind.R
Original file line number Diff line number Diff line change
Expand Up @@ -194,7 +194,12 @@ filearray_bind <- function(
unlink(filebase, recursive = TRUE)
}
re <- filearray_create(filebase = filebase, dimension = dim, type = type, partition_size = part_size)


# The partition files are copied or linked in below
lock_scope <- lock_scope_begin()
on.exit(lock_scope_end(lock_scope), add = TRUE)
lock_scope_acquire(lock_scope, re, re$.partition_info[, 1])

start <- 1
end <- 1

Expand Down
215 changes: 121 additions & 94 deletions R/class-filearray.R
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,74 @@
#' \item{\code{expand(n)}}{Expand array along the last margin; returns true if expanded; if the \code{dimnames} have been assigned prior to expansion, the last dimension names will be filled with \code{NA}}
#' \item{\code{initialize_partition()}}{Make sure a partition file exists; if not, create one and fill with \code{NA}s or 0 (\code{type='raw'})}
#' \item{\code{load(filebase, mode = c("readwrite", "readonly"))}}{Load file array from existing directory}
#' \item{\code{lock_partition(parts, timeout = NULL)}}{Lock partitions \code{parts} (default: all) for this \R session, so that other processes cannot write to them; waits and retries until all of them are locked or \code{timeout} seconds have passed (\code{NULL}: option \code{filearray.lock.timeout}, which defaults to waiting without limit; \code{0}: try once). Returns \code{TRUE} once locked, or \code{FALSE} after the time-out, with nothing locked; see section 'Partition locks'}
#' \item{\code{partition_path(part)}}{Get partition file path}
#' \item{\code{partition_size()}}{Get partition size; see \code{\link{filearray}}}
#' \item{\code{set_partition(part, value, ..., strict = TRUE)}}{Set partition value}
#' \item{\code{sexp_type()}}{Get data \code{SEXP} type; see R internal manuals}
#' \item{\code{show()}}{Print information}
#' \item{\code{type()}}{Get data type}
#' \item{\code{unlock_partition(parts)}}{Release one \code{lock_partition} lock on each of \code{parts}, or every \code{lock_partition} lock this \R session holds on the array if \code{parts} is missing; returns \code{TRUE} if a lock was released}
#' \item{\code{valid()}}{Check if the array is valid.}
#' }
#' @seealso \code{\link{filearray}}
#' @section Partition locks:
#' Several processes (for example \pkg{parallel}, \pkg{callr} or \pkg{mirai}
#' workers) can write to one array. Writes lock the partitions they change:
#' sub-assignment (\code{x[...] <- value}), \code{fill_partition},
#' \code{initialize_partition}, \code{set_partition}, \code{\link{fmap}},
#' \code{\link{fmap_element_wise}} and \code{\link{filearray_bind}} wait
#' while another process holds one of their partitions, so concurrent writers
#' never lose each other's data. Reading takes no locks, so a reader may see a
#' write that is still in progress. Functions that read but must first create
#' missing partitions (\code{max}, \code{min}, \code{range} and \code{sum}
#' with \code{na.rm = FALSE}, inputs of \code{fmap}, \code{filearray_bind},
#' evaluating array operations) lock the partitions they create, so they can
#' wait too.
#'
#' To keep partitions to yourself over several operations, for example to
#' check whether a partition has been computed and write it only if not, lock
#' them explicitly:
#' \preformatted{
#' if (!x$lock_partition(k, timeout = 60)) stop("partition ", k, " is busy")
#' on.exit(x$unlock_partition(k), add = TRUE)
#' if (is.na(x[1, 1, k])) x[, , k] <- compute(k)
#' }
#' Locks belong to the \R session, not to the object: every object of the
#' same array shares them, writes of this session to partitions it holds do
#' not wait, and each \code{lock_partition} call needs one
#' \code{unlock_partition} call. Lock everything you need in one call: two
#' sessions that each hold one partition and wait for the other's wait until
#' their time-outs. Likewise, do not hold partitions that forked workers
#' (\code{mclapply}) of the same session write. Non-interactive workers
#' should set a time-out.
#'
#' The locks are byte-range locks on an empty file \code{.filearray.lock} in
#' the array directory (byte \code{k} locks partition \code{k}). They are
#' advisory: older versions of this package and other programs ignore them.
#' The operating system releases them when a process exits, so a crashed or
#' killed process never leaves a stale lock. The locks guard writes but never
#' make them fail: when a lock cannot be taken for a reason other than another
#' process holding it (a file system without lock support such as some network
#' mounts, no permission to create the lock file, a network error), writes warn
#' once and proceed unlocked, and \code{lock_partition} warns and returns
#' \code{TRUE} with attribute \code{locked = FALSE}. Cloud-synchronized
#' folders lock on one computer only. Copying or archiving the array directory
#' in the session that holds its locks releases them on \verb{POSIX} systems.
#' To lock anything else the same way, see \code{\link{with_filelock}}.
#'
#' Options (each falls back to the environment variable, which worker
#' processes inherit):
#' \describe{
#' \item{\code{filearray.lock} (\code{FILEARRAY_LOCK})}{set to \code{FALSE}
#' to stop writes from locking automatically; \code{lock_partition} still
#' locks}
#' \item{\code{filearray.lock.timeout} (\code{FILEARRAY_LOCK_TIMEOUT})}{seconds
#' a write waits for a locked partition before it fails with an error of class
#' \code{filearray_lock_timeout}; default \code{Inf}, wait until the
#' partition is free. A value that is not a non-negative number warns once and
#' is ignored}
#' }
#' @seealso \code{\link{filearray}}, \code{\link{with_filelock}}
#' @exportClass FileArray
NULL

Expand Down Expand Up @@ -368,7 +427,12 @@ FileArray <- setRefClass(
return(FALSE)
}
}


# Hold the partition from reading it to writing it back
lock_scope <- lock_scope_begin()
on.exit(lock_scope_end(lock_scope), add = TRUE)
lock_scope_acquire(lock_scope, .self, part)

arglen <- ...length()
dim[[length(dim)]] <- .self$partition_size()
if ( arglen > 1 ) {
Expand Down Expand Up @@ -465,119 +529,79 @@ FileArray <- setRefClass(
if (length(value) > 1) {
quiet_warning("`fill_partition` value length coerced to first value")
}

value <- value[[1]]

type <- .self$type()
switch(
type,
"complex" = {
value <- cplxToReal2(as.complex(value))
},
"float" = {
value <- realToFloat2(as.double(value))
},
{
storage.mode(value) <- type
}
)
size <- get_elem_size(type)
file <- .self$partition_path(part)

# Write to a temporary file, then rename it to the partition:
# another process that loads this array while the partition is
# being written would reject a file shorter than the 1024-byte
# header. `load()` lists only `<n>.farr`, never `<n>..tmp`
tmp <- sprintf("%s.tmp%s", gsub("farr$", "", file, ignore.case = TRUE), Sys.getpid())
fid <- file(tmp, "wb")
fid_closed <- FALSE

on.exit({
# `tmp` still exists only if writing or renaming failed
if (file.exists(tmp)) {
if (!fid_closed) {
close(fid)
}
unlink(tmp)
}
}, add = TRUE)

if ( part <= nrow(.self$.partition_info)) {
partition_size <- .self$.partition_info[part, 2]
} else {
partition_size <- .self$partition_size()
}
dimension <- .self$dimension()
dimension[[length(dimension)]] <- partition_size
part_len <- prod(dimension)
buffer_len <- get_buffer_size() / size
if (buffer_len > part_len) {
buffer_len <- part_len
}

write_header(
fid = fid,
partition = part,
dimension = dimension,
type = type,
size = size
)
seek(con = fid, where = HEADER_SIZE, rw = "write")
buf <- writeBin(
con = raw(),
object = rep(value, buffer_len),
size = size,
endian = ENDIANNESS
)
nloop <- floor(part_len / buffer_len)
replicate(nloop, {
writeBin(con = fid, object = buf)
NULL
})
rest <- part_len - buffer_len * nloop
if ( rest > 0 ) {
writeBin(con = fid, object = buf[seq_len(rest * size)])
}
seek(con = fid, where = HEADER_SIZE - 8L, rw = "write")
writeBin(con = fid, object = part_len, size = 8L, endian = ENDIANNESS)

close(fid)
fid_closed <- TRUE

# rename() replaces `file` in one step, so other processes see
# either the old partition or the new one, never a partial one
if (!file.rename(tmp, file)) {
stop("Cannot write partition file: ", file)
}
lock_scope <- lock_scope_begin()
on.exit(lock_scope_end(lock_scope), add = TRUE)
locked <- lock_scope_acquire(lock_scope, .self, part)
write_partition_file(.self, part, value, replace = TRUE, locked = locked)
invisible()
},
initialize_partition = function(parts) {
if (!.self$valid()) {
stop("Invalid file array")
}
if (isTRUE(.self$.mode == "readonly")) {
.self$.mode <- "readwrite"
on.exit({
.self$.mode <- "readonly"
})
}
if (missing(parts)) {
parts <- .self$.partition_info[, 1]
}
parts <- parts[!is.na(parts)]
# Read APIs (max, min, range, sum, filearray_bind, ...) call this:
# when every partition exists, they return here without locking
parts <- parts[!file.exists(.self$partition_path(parts))]
if (!length(parts)) {
return(invisible())
}
parts <- sort(unique(as.integer(parts)))
if (any(parts <= 0)) {
stop("NA or non-positive partition are invalid")
}

lock_scope <- lock_scope_begin()
on.exit(lock_scope_end(lock_scope), add = TRUE)
locked <- lock_scope_acquire(lock_scope, .self, parts)
for (part in parts) {
file <- .self$partition_path(part)
if (!file.exists(file)) {
.self$fill_partition(part, NA)
# Another process may have created the partition meanwhile;
# write_partition_file() never replaces one that exists
if (!file.exists(.self$partition_path(part))) {
write_partition_file(.self, part, NA, replace = FALSE, locked = locked)
}
}
invisible()
},
can_write = function() {
if (isTRUE(.self$.mode == "readonly")) {
return(FALSE)
}
return(TRUE)
},
lock_partition = function(parts, timeout = NULL) {
if (!.self$valid()) {
stop("Invalid file array")
}
nparts <- nrow(.self$.partition_info)
if (missing(parts)) {
parts <- seq_len(nparts)
}
parts <- validate_lock_parts(parts, nparts)
if (is.null(timeout)) {
timeout <- lock_timeout_default()
} else {
timeout <- validate_lock_timeout(timeout)
}
locked <- lock_partition_explicit(.self$.filebase, parts, timeout)
if (isTRUE(locked)) {
return(invisible(locked))
}
locked
},
unlock_partition = function(parts) {
if (missing(parts)) {
released <- lock_release_explicit(.self$.filebase)
} else {
released <- lock_release_explicit(
.self$.filebase, validate_lock_parts(parts))
}
invisible(released > 0)
},
delete = function(force = FALSE) {
if (!.self$valid()) {
.self$.valid <- FALSE
Expand All @@ -591,6 +615,9 @@ FileArray <- setRefClass(
}
}
filebase <- .self$.filebase
# Drop this session's partition locks and close the lock file:
# on Windows an open file keeps the directory from being removed
lock_forget(filebase)
if (dir.exists(filebase)) {
unlink(filebase, recursive = TRUE, force = force)
}
Expand Down
6 changes: 6 additions & 0 deletions R/filearray-package.R
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,12 @@ symlink_enabled <- local({
ns$NA_float_ <- get_float_na()
}

.onUnload <- function(libpath) {
# A reloaded package starts a new lock registry, which could not see the
# locks of this one: release them all
try(FARR_lock_release_all(), silent = TRUE)
}

.onAttach <- function(libname, pkgname) {
if (Sys.getenv("_R_CHECK_LIMIT_CORES_") == "TRUE") {
packageStartupMessage(
Expand Down
Loading
Loading