Reduce shard results without gathering all per-shard returns on the master.
shard_reduce() executes map() over shards in parallel and combines results
using an associative combine() function. Unlike shard_map(), it does not
accumulate all per-shard results on the master; it streams partials as chunks
complete.
Usage
shard_reduce(
shards,
map,
combine,
init,
borrow = list(),
out = list(),
workers = NULL,
chunk_size = "auto",
profile = c("default", "memory", "speed"),
mem_cap = "2GB",
recycle = TRUE,
cow = c("deny", "audit", "allow"),
seed = NULL,
diagnostics = TRUE,
packages = NULL,
init_expr = NULL,
timeout = 3600,
max_retries = 3L,
health_check_interval = 10L
)Arguments
- shards
A
shard_descriptorfromshards(), or an integer N.- map
Function executed per shard. Receives shard descriptor as first argument, followed by borrowed inputs and outputs.
- combine
Function
(acc, value) -> accused to combine results. Must be associative, and must accept two mapped values as arguments (worker partials start from a chunk's first mapped value; see Initial value semantics).- init
Initial accumulator value, combined exactly once on the master (see Initial value semantics).
- borrow
Named list of shared inputs (same semantics as
shard_map()).- out
Named list of output buffers/sinks (same semantics as
shard_map()).- workers
Number of worker processes.
- chunk_size
Shards to batch per worker dispatch. The default
"auto"targets roughly four chunks per worker,max(1, ceiling(num_shards / (workers * 4))), which amortizes dispatch round trips while retaining load balance. Supply an integer to control batching explicitly.- profile
Execution profile (same semantics as
shard_map()).- mem_cap
Memory cap per worker (same semantics as
shard_map()).- recycle
Worker recycling policy (same semantics as
shard_map()).- cow
Copy-on-write policy for borrowed inputs (same semantics as
shard_map()).- seed
RNG seed for reproducibility. When non-
NULL, one independent L'Ecuyer-CMRG stream per shard is derived on the master and installed in the worker immediately before eachmap()call, so per-shard RNG draws are reproducible regardless of worker count,chunk_size, or dynamic shard-to-worker assignment. The master's RNG state andRNGkind()are left exactly as found (seed = NULLtouches no RNG state). Note on floating-point results: partials are combined in chunk order, which is deterministic given identical chunking, so a given (seed, chunking,chunk_size) is exactly reproducible for anyworkers=; across differentchunk_sizevalues the combine order differs, so non-associative floating-point rounding may differ even though the per-shard RNG draws are identical. Whenshardsis a scalar N andseedis set, the shard decomposition is chosen independently of the worker count so the same seed gives identical results for anyworkers=.- diagnostics
Logical; collect diagnostics (default TRUE).
- packages
Additional packages to load in workers.
- init_expr
Expression to evaluate in each worker on startup.
- timeout
Seconds to wait for each chunk.
- max_retries
Maximum retries per chunk.
- health_check_interval
Check worker health every N completions.
Value
A shard_reduce_result with fields:
value: final accumulatorfailures: any permanently failed chunksdiagnostics: run telemetry including reduction statsqueue_status,pool_stats
Details
For performance and memory efficiency, reduction is performed in two stages:
per-chunk partial reduction inside each worker, and
streaming combine of partials on the master, folded in chunk order.
Initial value semantics
init is combined exactly once, on the master, at the start of the final
fold: the result is combine(combine(combine(init, p1), p2), ...) where
p1, p2, ... are per-chunk partials in chunk order. Worker-side partials
are built without init: each chunk's partial starts from the chunk's
first mapped value. A non-neutral init (e.g. init = 10 with +)
therefore contributes exactly once, regardless of chunk_size or
workers. This requires combine to be associative and able to combine
two mapped values (not just an accumulator with a value).