ADR 0047: Over-budget construction partitions succeed; one CPU budget per instance
ADR 0047: Over-budget construction partitions succeed; one CPU budget per instance
Section titled “ADR 0047: Over-budget construction partitions succeed; one CPU budget per instance”Status: Accepted
Implementation: Decision 1 is implemented by #1585 and decision 2 by #1586; see the implementation updates below.
Build target: v0.6.0
Related: ADR 0046 (construction reuse decisions; this record takes the two decisions it reserved for maintainers), ADR 0038 (determinism at the publication boundary), ADR 0045 (ingest authentication regime); #337 (per-instance execution resource policy); #1448 (shaping parallelism); #1504 (construction reuse epic).
Implementation update: external partitions (#1585)
Section titled “Implementation update: external partitions (#1585)”#1585 compared the #1507 DataFusion adapter with a native bounded merge under
the #1505 protocol and selected the native merge
(docs/development/evidence/external-partition-comparison-1585.md).
A fixed-width partition without a detail codec whose materialization would
exceed max_partition_bytes is sorted on its load worker into runs of at most
max_partition_bytes. Each run is written as a construction artifact
temporary with an XXH64 checksum. The coordinator merges the runs into the
consumer the resident path feeds, so the shaped bytes are unchanged. The
families without a detail codec are identities and the staged and resolved
endpoints; the endpoint families are the ones a hub can push over budget.
Detail-codec partitions and Arrow property-row partitions keep the refusal.
How each obligation is met:
- Recorded parameter.
GraphConstructionBudgets::max_external_partition_bytesbounds the scratch for one partition. It defaults to 64 GiB and must be zero or at leastmax_partition_bytes. Zero keeps the refusal. A checkpoint recorded before the field exists reads it as zero, and a default-budget resume keeps that recorded value. - Owned scratch. Runs are
.artifact-xrun-p<N>-<random>.tmptemporaries in the session directory. Every exit path unlinks them, and session open and recovery reclaim any a crash leaves. Runs are never resume authority; the sealed segments are. - Bounded scratch. A partition whose sealed segments exceed the external
bound is refused before any run is written. The ADR text above says the
limit is “derived from the allocation ledger”. The ledger has no scratch
budget to derive one from, so the limit is a recorded budget instead. Run
counts and bytes are recorded in construction evidence and on the import
receipt (
external_partitions,external_runs,external_run_bytes). Runs are not charged to the allocation ledger’s transient peak, because the load workers write them concurrently and the recorded peak would then depend on their interleaving. - Integrity. The merge verifies each run’s file identity, length, record count and checksum, and refuses before any derived output is installed.
- Tested pool size. No library pool is used. Memory per run is the resident budget itself.
- Thread-based coordinator. The merge runs on the existing plain-thread coordinator.
- Unchanged resident path. Partitions within budget take the resident path unchanged.
Evidence: docs/development/evidence/external-partitions-1585.md.
Implementation update: the instance construction budget (#1586)
Section titled “Implementation update: the instance construction budget (#1586)”Each instance builds one ConstructionCpuAdmission, sized
compute_threads - construction_cpu_reserve, and attaches it to every
construction session it opens. Finish-time partition loads lease their workers
from it. Import normalization leases its compute-pool lanes for each flush on
the calling thread. Concurrent imports on the instance therefore share one
construction limit.
Query kernels are not admission-gated. They keep the instance’s
compute_threads pool, because each kernel splits its work into
compute_threads chunks, and a width that varied with load would change its
chunking and floating-point reduction order. The budget is shared from the
construction side. Construction never uses more than compute_threads - reserve
lanes across all imports. Normalization therefore never occupies more of the
query pool than that, so at least reserve pool threads stay free for queries.
Finish-time loads run on their own threads and count against the same limit.
The default reserve is 1. Measured at 4 and 8 compute threads with two
concurrent imports, a larger reserve did not lower query latency beyond
run-to-run variation, and the limit held at every setting. Evidence:
docs/development/evidence/construction-cpu-budget-1586.md.
Context
Section titled “Context”ADR 0046 kept construction’s own machinery and left two decisions to the maintainers, because each changes an existing requirement:
- Must an over-budget partition succeed rather than refuse?
- Should concurrent imports share one CPU budget?
The refusal. Construction routes endpoint records by node UUID, two per
edge, so every record of one node lands in one partition. A partition whose
materialization exceeds the recorded max_partition_bytes refuses the whole
ingest before allocating; the default budget is 256 MiB. At 33 bytes per
endpoint record, one node with roughly 8.1 million edges fills the budget on
its own. More ranges cannot split it and more RAM does not raise the recorded
budget. The measurements behind this record
(docs/development/evidence/partition-refusal-1584.md) show:
- A 9,000,000-edge star graph, one hub with every other node linked to it, is
refused:
partition materialization requires 297536415 bytes, exceeds recorded budget 268435456. A 4,000,000-edge star publishes. A class node in a knowledge graph, linked to every entity of its type, has this shape. One such hub refuses an ingest of roughly the S19 Graph500 rung’s edge count. - Graph500 inputs grow their largest hub by 1.52× per scale step, measured from S18 to S26, as the generator’s parameters predict. At S26, the certification target, the largest hub has 1,709,763 edges, and the modelled largest partition is 73.7 MB: 27% of the budget. The model fits the measured S18 and S20 partitions within 6%.
- The #1507 DataFusion external-sort adapter, run with its pool equal to the
default budget, failed on the 9M star with DataFusion’s own
Resources exhaustedduring the merge. With 64 MiB and 8 MiB pools it published the same answers. Every earlier run that spilled had used a pool of 2 MiB or less, so a spilling sort under a large pool had never been tested. The cause inside DataFusion is not isolated.
The CPU budget. Each GraphForge instance sizes one private CPU pool from
compute_threads (#337). Query kernels run on it, and import normalization
has run on it since #1472, with nothing reserving a share for queries. The
finish-time partition loads use private threads outside any budget, and #1448
is about to add shaping lanes. #1508 F12 measured the trade: a shared budget of
two held two concurrent imports to two cores, at 35.5 ms against 19.4 ms on
private pools.
Decision
Section titled “Decision”1. An over-budget fixed-width partition is processed externally
Section titled “1. An over-budget fixed-width partition is processed externally”A fixed-width construction partition whose materialization exceeds the
recorded max_partition_bytes is sorted with bounded memory and streamed to
its shaped output, instead of refusing the ingest. The budget keeps bounding
memory; it stops bounding which inputs are accepted.
Refusal remains, as structured errors that publish nothing, for:
- exhausted scratch disk under a recorded limit;
- cancellation;
- integrity failures in scratch or inputs.
The mechanism is not chosen here. #1585 compares the #1507 DataFusion adapter with a native bounded external merge over GraphForge’s sealed segments, under the #1505 protocol, before measuring. Whichever wins must meet every obligation below.
- Recorded parameter. External processing is a recorded construction parameter, so a resume validates it like any budget. A session recorded before it resumes under its recorded refusal contract.
- Owned scratch. Scratch is created through the construction directory and reclaimed at recovery. It is never recovery authority.
- Bounded scratch. A scratch disk limit derived from the allocation ledger refuses cleanly when exhausted, and scratch bytes are accounted in construction evidence.
- Integrity. Scratch is checksummed, or checked by the #1507 record-multiset guard, before anything derived from it is published.
- Tested pool size. Any library pool is a sub-budget whose size is chosen
from tests at the partition sizes it will meet, not set equal to
max_partition_bytes. - Thread-based coordinator. The coordinator is a plain thread, because a streaming merge that blocks on its own runtime cannot run under a Tokio coordinator (#1509).
- Unchanged resident path. Partitions within budget keep today’s resident path, and publication stays byte-stable under ADR 0038.
2. One CPU budget per instance, shared by queries and construction
Section titled “2. One CPU budget per instance, shared by queries and construction”An instance has one CPU budget, sized from compute_threads. All
CPU-parallel work in the instance draws from it:
- query kernels;
- import normalization;
- finish-time partition loads;
- #1448’s shaping lanes.
Construction may hold at most compute_threads - reserve of it, with a reserve
of at least one, so queries always have a share while an import runs. Work
that blocks on file I/O runs on construction-owned threads holding admission,
never on the ComputePool workers that queries depend on. Admission is
scheduling only: published bytes and construction evidence must not depend on
the budget. Waiting for admission is cancellable.
#1586 implements this, chooses the default reserve from measured evidence, and exposes it in the resource policy and its diagnostics. #1448’s lanes draw from it, so #1586 blocks #1448.
Scope limits
Section titled “Scope limits”- Row partitions. Arrow property-row partitions keep their refusal. They were never evaluated under an external path; the question stays open under #1504.
- Instance, not process. The budget is per instance, matching the resource policy’s scope since #337. Two instances in one process each keep their own. A process-wide budget is not decided here.
- Graph500 claim. The Graph500 growth model predicts that the largest partition exceeds the default budget at S29, and that every endpoint partition exceeds it at S30, where the 4,096-range maximum stops growing. Both are model outputs beyond anything measured.
Relation to ADR 0046
Section titled “Relation to ADR 0046”This record resolves the two decisions ADR 0046 reserved. It replaces that record’s “Over-budget partitions: retain refusal as the contract” row with decision 1. Every other ADR 0046 decision stands, including the retained production load pool; decision 2 constrains how that pool and #1448’s lanes are admitted, not which scheduler runs them.
Consequences
Section titled “Consequences”- A graph with a hub of any degree ingests, bounded by scratch disk, once #1585 lands. Until then the refusal stands and the architecture page says so.
- #1582 no longer retires the #1507 adapter until #1585 chooses between it and a native merge.
- #1448’s lanes are built against #1586’s admission, not retrofitted to it.
- Imports can take longer beside interactive queries, by design.
Revisit when
Section titled “Revisit when”- #1585’s comparison finds neither mechanism can meet the obligations above at an acceptable measured cost. The decision then returns to the maintainers with that evidence.
- Multi-instance processes become a supported deployment, which raises the process-wide budget question.
- Row partitions are evaluated under an external path.
Evidence
Section titled “Evidence”| Evidence | What it establishes |
|---|---|
docs/development/evidence/partition-refusal-1584.md |
Graph500 hub growth S18–S26, the partition-size model and its fit, the star-graph refusal, and the #1507 adapter’s failure under a large pool and success under smaller ones |
docs/development/evidence/construction-reuse-integrated-1509.md |
The hybrid’s measured cost at a forced 1 MiB budget, and the Tokio-coordinator incompatibility |
docs/development/evidence/construction-scheduling-spike-1508.md |
F12 shared-admission measurement; F13 on the API runtime’s blocking pool |