Skip to content

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.

julia
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.

julia
unequal = DGG.partition(problem, 2; capacities=[1, 2])
DGG.partweights(unequal)
2-element Vector{Float64}:
 1.0
 3.0

The API keeps row positions separate from application IDs:

valuemeaning
ordera 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.

julia
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)
end

Each 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.

julia
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:

PlaceholderRequired meaning
dag.graphA GlobalRegridding chunk dependency graph
dag.orderA permutation of its destination row positions
todochunks / tilesApplication IDs aligned with destination / source graph rows
estimated_work / tile_bytesNonnegative work / source weights in those row orders
worker_ids / worker_capacitiesWorker labels and relative positive compute capacities
text
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 ​

AlgorithmUse it forLoad on the coordinator
WeightedContiguousA quick split that keeps traversal orderAvailable by default
MetisPartitionBalancing work while keeping strongly connected chunks togetherusing Metis
KaHyParPartitionReducing shared source replication, especially when one source serves many chunksusing KaHyPar_jll
ScotchPartitionA graph mapping alternative with weighted worker capacitiesusing 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:

julia
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:

julia
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:

julia
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:

julia
function DGG.partitionlabels(algorithm::MyPartitioner,
        problem::DGG.PartitionProblem, nparts::Int,
        capacities::Vector{Float64})
    # Return one logical partition number per original problem row.
end

The 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
julia
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.

source
DiscreteGlobalGrids.partitionproblem Function
julia
partitionproblem(input; kwargs...) -> PartitionProblem

Adapt a chunk plan or dependency graph to a partitioning problem. Packages may add methods for their own work descriptions.

source
DiscreteGlobalGrids.partition Function
julia
partition(problem, nparts; algorithm=WeightedContiguous(), capacities=ones(nparts))
    -> ChunkPartition

Assign indivisible work rows to exactly nparts logical partitions. Empty partitions are retained.

source
DiscreteGlobalGrids.ChunkPartition Type
julia
ChunkPartition

A serializable logical chunk assignment. Its vector fields are owned by the result and exposed read-only through the accessors below.

source
DiscreteGlobalGrids.npartitions Function

Return the exact logical partition count, including empty partitions.

source
DiscreteGlobalGrids.partindices Function

Return input row positions in traversal order. Treat the vector as read-only.

source
DiscreteGlobalGrids.partchunks Function

Return stable work IDs in traversal order. Treat the vector as read-only.

source
DiscreteGlobalGrids.partsources Function

Return unique source IDs in source-axis order. Treat the vector as read-only.

source
DiscreteGlobalGrids.partweights Function

Return original work-weight sums by partition. Treat the vector as read-only.

source
DiscreteGlobalGrids.AbstractPartitioningAlgorithm Type
julia
AbstractPartitioningAlgorithm

Supertype for chunk partitioning algorithms. Implement partitionlabels to add an algorithm.

source
DiscreteGlobalGrids.WeightedContiguous Type
julia
WeightedContiguous()

Cut work into capacity-weighted contiguous segments in the problem's traversal order.

source
DiscreteGlobalGrids.MetisPartition Type
julia
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.

source
DiscreteGlobalGrids.KaHyParPartition Type
julia
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.

source
DiscreteGlobalGrids.ScotchPartition Type
julia
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.

source
DiscreteGlobalGrids.partitionlabels Function
julia
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.

source
DiscreteGlobalGrids.PartitionBackendUnavailable Type
julia
PartitionBackendUnavailable

The requested optional partitioning backend has not been loaded.

source