Assigning chunks to workers
Partitioning assigns whole chunks to logical groups that an executor can place on threads or processes. Use it to balance work and, with a suitable algorithm, keep computations that read the same source data together.
The same API accepts a neighbourhood MapChunkPlan, a GlobalRegridding.ChunkedPlan, or its destination-to-source dependency graph. Each result records the assigned chunks and the sources each partition needs. The result contains ordinary numeric arrays and can be serialized without an open store, a regridding plan, or the partitioning library.
Describe the work
A PartitionProblem has one row per work chunk. Each row lists the source resources that chunk reads. Work weights estimate computation; source weights estimate the cost of reading or transferring shared data.
import DiscreteGlobalGrids as DGG
problem = DGG.PartitionProblem([[1], [1], [2], [2]];
ids=[11, 12, 21, 22], sourceids=[101, 102],
weights=[1.0, 1.0, 1.0, 1.0])
assignment = DGG.partition(problem, 2)
[(chunks=DGG.partchunks(assignment, p),
sources=DGG.partsources(assignment, p)) for p in 1:2]2-element Vector{@NamedTuple{chunks::Vector{Int64}, sources::Vector{Int64}}}:
(chunks = [11, 12], sources = [101])
(chunks = [21, 22], sources = [102])Chunks 11 and 12 read source 101; chunks 21 and 22 read source 102. The default WeightedContiguous algorithm cuts the supplied traversal order into groups with approximately equal work.
capacities=[1, 2] asks the second partition to carry twice the work of the first. These are relative compute capacities. Chunks remain indivisible, so their sizes limit the balance any algorithm can achieve. Capacity targets are not memory limits.
unequal = DGG.partition(problem, 2; capacities=[1, 2])
DGG.partweights(unequal)2-element Vector{Float64}:
1.0
3.0The API keeps row positions separate from application IDs:
| value | meaning |
|---|---|
order | a permutation of input row positions |
partindices(assignment, p) | assigned input rows, in traversal order |
partchunks(assignment, p) | the corresponding application chunk IDs |
partsources(assignment, p) | unique application source IDs needed by the partition |
The requested number of logical partitions is preserved. Some can be empty when there is little work. Partition numbers are independent of process IDs; an executor chooses where each partition runs.
Partition a neighbourhood sweep
Build the full chunk plan, then partition it. partitionproblem identifies which storage chunks supply each chunk's owned cells and halo. Its default work estimate counts the owned and halo cells; override weights when measured kernel costs or other array dimensions matter.
plan = DGG.chunkplan(A; halo=1)
assignment = DGG.partition(plan, Threads.nthreads())
out = zeros(Float64, size(A))
@sync for p in 1:DGG.npartitions(assignment)
piece = plan[DGG.partindices(assignment, p)]
Threads.@spawn DGG.mapneighbors!(out, kernel, A, piece; threaded=false)
endEach piece writes its owned indices. For a stored destination, also align work ownership with its physical write chunks so concurrent writes do not update the same storage block.
The current chunk runner reads each chunk's halo separately. Grouping related chunks creates an opportunity for reuse; an executor needs a cache or halo exchange to realize it. See chunk sweeps for the loading and kernel contracts.
Partition a regridding run
A regridding partition assigns destination chunks. Source chunks remain read dependencies: several partitions may need the same source.
source = DGG.levelgrid(DGG.HEALPixSystem(), 1)
target = DGG.levelgrid(DGG.HEALPixSystem(), 2)
values = ones(DGG.ncells(source))
regridplan = DGG.plan_regrid(values; from=DGG.DGGSpace(source; chunklevel=0),
to=DGG.DGGSpace(target; chunklevel=1), method=DGG.NearestCell(), lazy=true)
assignment = DGG.partition(regridplan, 4)
destination_chunks = DGG.partchunks(assignment, 1)
source_chunks = DGG.partsources(assignment, 1)
@assert sort(vcat([DGG.partindices(assignment, p) for p in 1:4]...)) ==
collect(1:length(DGG.partitionproblem(regridplan).ids))
(; destination_chunks, source_chunks)(destination_chunks = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12], source_chunks = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11])This reads the dependency relation the plan already owns. It builds neither regridding weights nor a replacement dependency graph. The relation is conservative: a listed source may be ruled out when the actual interpolation or overlap weights are built.
For a run whose graph rows correspond to application-specific chunk IDs, provide the mapping explicitly. For example, a Copernicus DEM run can use its destination store chunk numbers and source tile numbers. The following is a schematic for an application that already has these objects:
| Placeholder | Required meaning |
|---|---|
dag.graph | A GlobalRegridding chunk dependency graph |
dag.order | A permutation of its destination row positions |
todochunks / tiles | Application IDs aligned with destination / source graph rows |
estimated_work / tile_bytes | Nonnegative work / source weights in those row orders |
worker_ids / worker_capacities | Worker labels and relative positive compute capacities |
problem = DGG.partitionproblem(dag.graph;
ids=todochunks, sourceids=tiles, order=dag.order,
weights=estimated_work, sourceweights=tile_bytes)
assignment = DGG.partition(problem, length(worker_ids);
capacities=worker_capacities)
jobs = [(chunks=DGG.partchunks(assignment, p),
tiles=DGG.partsources(assignment, p))
for p in 1:DGG.npartitions(assignment)]Here estimated_work follows graph row order, tile_bytes follows source order, and dag.order is the existing traversal permutation. On a restricted dependency graph, default destination IDs retain their original chunk numbers; partindices still refers to the restricted graph's local rows.
Send each job and the input configuration to its worker. The worker opens its own data handles and schedules the assigned destination chunks. Tasks within that worker can pull work dynamically while sharing a tile cache. Recompute cache consumer counts for the worker's assigned destinations; global consumer counts include work that other workers will finish.
For resumed runs, describe only the pending destinations and retain their application IDs. An assignment describes ownership; the executor still owns task completion, retries, data exchange, and output writes.
Choose an algorithm
| Algorithm | Use it for | Load on the coordinator |
|---|---|---|
WeightedContiguous | A quick split that keeps traversal order | Available by default |
MetisPartition | Balancing work while keeping strongly connected chunks together | using Metis |
KaHyParPartition | Reducing shared source replication, especially when one source serves many chunks | using KaHyPar_jll |
ScotchPartition | A graph mapping alternative with weighted worker capacities | using Scotch |
All four return the same ChunkPartition. Only the coordinator that builds an assignment needs its optional library. The workers receive chunk and source IDs and open their own data handles.
The algorithm types are always available from DiscreteGlobalGrids. Calling an optional algorithm before loading its library raises PartitionBackendUnavailable with an installation and loading hint.
Reduce shared source replication with KaHyPar
KaHyPar is a useful starting point for a CopDEM run whose workers share a tile cache. It balances destination work while grouping chunks that read the same tiles:
using KaHyPar_jll
assignment = DGG.partition(problem, length(worker_ids);
algorithm=DGG.KaHyParPartition(seed=42, imbalance=0.03),
capacities=worker_capacities)Each source forms one hyperedge: the set of destination chunks that read it. The connectivity objective charges sourceweight * (partitions_using_source - 1) for each used source. With tile byte sizes as source weights, this estimates extra bytes read when each worker loads a tile once. Actual reuse still depends on the worker's cache and eviction policy.
This representation stores one entry per destination-to-source dependency. A tile shared by a thousand destinations remains a single source set. The extension retains large source sets during optimization and disables the upstream preset's incidence sparsification.
The extension calls the KaHyPar C API through KaHyPar_jll 1.3.5 or later. This supports the current native build; KaHyPar.jl 0.3.1 pins an older build with a conflicting Boost dependency in this workspace. KaHyPar's native library and bundled configuration use GPL-3.0. The 1.3.5 artifacts support Linux, macOS, and FreeBSD, with no Windows build.
KaHyPar uses integer weights. The extension scales work and source costs into bounded integer ranges and gives zero-cost work the minimum unit required by the library. Worker capacities become integer upper bounds that include imbalance. Indivisible chunks can exceed those bounds; inspect partweights when planning a run.
Use METIS
Load Metis.jl to enable MetisPartition:
using Metis
assignment = DGG.partition(problem, 4;
algorithm=DGG.MetisPartition(seed=42, imbalance=0.03))The extension builds a graph whose vertices are work chunks. Two chunks are connected when they read a common source. A source with d consumers adds sourceweight / (d - 1) to each pair's connection weight. METIS balances work weights while reducing connections that cross partitions.
This graph approximates shared-data costs. It does not count exact replicated bytes or enforce a cache budget. A source with many consumers creates many pairwise connections; maxedges bounds the graph expansion. Native weights are quantized to METIS integers. Use consistent inputs and a fixed seed for repeatable calls; partition labels are not a persistent identity across library versions.
Map work with Scotch
Scotch.jl offers another graph partitioner through the same API:
using Scotch
assignment = DGG.partition(problem, 4;
algorithm=DGG.ScotchPartition(seed=42), capacities=[1, 2, 1, 2])The extension uses the same bounded shared-source graph as METIS and maps it onto a complete worker graph weighted by capacities. Work, connections, and capacities become integer weights; zero-cost work receives the minimum unit. maxedges limits graph expansion. Scotch has a fresh random context for each assignment, so calls with the same seed and inputs can be repeated.
All native backends use heuristics, and their results can change with library versions. A fixed seed supports repeatable calls within a build. Empty work, one partition, and problems with at least as many partitions as chunks use the contiguous splitter after the requested backend is loaded. Problems without positive shared-source costs also use that splitter.
Add a partitioner
Implement partitionproblem to adapt another kind of plan to work rows and source dependencies. The problem retains the full dependency relation so an algorithm can use it directly or build its own graph representation.
For another algorithm, subtype AbstractPartitioningAlgorithm and implement partitionlabels:
function DGG.partitionlabels(algorithm::MyPartitioner,
problem::DGG.PartitionProblem, nparts::Int,
capacities::Vector{Float64})
# Return one logical partition number per original problem row.
endThe caller supplies validated, normalized capacities. Return integer labels in 1:nparts, aligned with the original rows even when problem.order is a permutation. The common API validates those labels and constructs the result's ordered work lists and unique source lists.
DiscreteGlobalGrids.PartitionProblem Type
PartitionProblem(reads; weights=nothing, sourceweights=nothing, ids=nothing,
sourceids=nothing, order=nothing)A bipartite chunk-partitioning problem. Each row of reads contains source resource positions used by one indivisible work item.
The fields reads, weights, sourceweights, ids, sourceids, and order are owned canonical vectors. Treat them as read-only. order contains work row positions, while ids and sourceids carry stable external identities.
DiscreteGlobalGrids.partitionproblem Function
partitionproblem(input; kwargs...) -> PartitionProblemAdapt a chunk plan or dependency graph to a partitioning problem. Packages may add methods for their own work descriptions.
sourceDiscreteGlobalGrids.partition Function
partition(problem, nparts; algorithm=WeightedContiguous(), capacities=ones(nparts))
-> ChunkPartitionAssign indivisible work rows to exactly nparts logical partitions. Empty partitions are retained.
DiscreteGlobalGrids.ChunkPartition Type
ChunkPartitionA serializable logical chunk assignment. Its vector fields are owned by the result and exposed read-only through the accessors below.
sourceDiscreteGlobalGrids.npartitions Function
Return the exact logical partition count, including empty partitions.
sourceDiscreteGlobalGrids.partindices Function
Return input row positions in traversal order. Treat the vector as read-only.
sourceDiscreteGlobalGrids.partchunks Function
Return stable work IDs in traversal order. Treat the vector as read-only.
sourceDiscreteGlobalGrids.partsources Function
Return unique source IDs in source-axis order. Treat the vector as read-only.
sourceDiscreteGlobalGrids.partweights Function
Return original work-weight sums by partition. Treat the vector as read-only.
sourceDiscreteGlobalGrids.AbstractPartitioningAlgorithm Type
AbstractPartitioningAlgorithmSupertype for chunk partitioning algorithms. Implement partitionlabels to add an algorithm.
DiscreteGlobalGrids.WeightedContiguous Type
WeightedContiguous()Cut work into capacity-weighted contiguous segments in the problem's traversal order.
sourceDiscreteGlobalGrids.MetisPartition Type
MetisPartition(; seed=0, imbalance=0.03, maxedges=1_000_000)Partition work by shared-source affinity through the optional Metis backend. maxedges bounds the projected graph's unique undirected edges. METIS represents imbalance in thousandths and uses 0.001 for a requested zero.
DiscreteGlobalGrids.KaHyParPartition Type
KaHyParPartition(; seed=0, imbalance=0.03)Group chunks by shared source data with KaHyPar's connectivity objective. Each source costs its weight for every additional partition that reads it. Load KaHyPar_jll to enable this optional backend. Capacity targets become native integer upper bounds, including imbalance; indivisible chunks may exceed them.
DiscreteGlobalGrids.ScotchPartition Type
ScotchPartition(; seed=0, imbalance=0.03, maxedges=1_000_000)Map the shared-source graph onto workers with relative compute capacities. Load Scotch to enable this optional backend. maxedges bounds the projected graph's unique undirected edges.
DiscreteGlobalGrids.partitionlabels Function
partitionlabels(algorithm, problem, nparts, capacities) -> Vector{Int}Return one label in 1:nparts for each problem row. capacities contains validated normalized targets. Extensions define this hook for custom AbstractPartitioningAlgorithm subtypes.
DiscreteGlobalGrids.PartitionBackendUnavailable Type
PartitionBackendUnavailableThe requested optional partitioning backend has not been loaded.
source