From a27944e1835c2d79ca84edc856d1cc8bbbc4cdf2 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Thu, 17 Sep 2026 12:56:09 +0000 Subject: [PATCH 1/5] fix(storage): scale the partition cut with data instead of a flat cap (#1439) Past ~4.19M staged identity records, the cut formula `min(partition_count, max(1, total/16))` saturated at the flat `partition_count` cap (256 by default), so partition *count* stopped rising while partition *size* kept growing with data. Each partition is materialized whole in memory and sorted, so that let resident partition memory scale with data -- exactly what `rss_bounded_or_plateaued` exists to catch, and it refused S22. Add a second, data-driven term to the cut: floored at DEFAULT_PARTITION_COUNT (256) and growing as total_records / TARGET_ROWS_PER_PARTITION (16,384), min-combined with the existing small-scale floor (unchanged) and the requested ceiling. 16,384 is chosen so the term reproduces the already-proven 256-partition operating point at S18's record count and reaches exactly MAX_PARTITION_COUNT (4,096) at S22's, matching the table in #1439. Default `GraphConstructionBudgets::partition_count` to MAX_PARTITION_COUNT so this data-driven term, not an artificially low flat cap, is what bounds production partition counts; it remains a pure function of recorded staged data, never machine-derived (R1). Extends the cross-partition-count determinism test to the new default ceiling and adds a dedicated arithmetic test for the S18/S19/S20/S22 scaling table directly. --- .../src/construction_determinism_tests.rs | 39 ++++++++++++--- .../src/graph_construction.rs | 13 ++++- .../src/graph_construction/partition.rs | 49 +++++++++++++++++-- .../src/graph_construction/partition/tests.rs | 31 ++++++++++++ .../src/graph_construction/tests.rs | 16 ++++-- 5 files changed, 130 insertions(+), 18 deletions(-) diff --git a/crates/graphforge-storage/src/construction_determinism_tests.rs b/crates/graphforge-storage/src/construction_determinism_tests.rs index e192488ab..315a9703f 100644 --- a/crates/graphforge-storage/src/construction_determinism_tests.rs +++ b/crates/graphforge-storage/src/construction_determinism_tests.rs @@ -451,7 +451,18 @@ mod determinism { let edges = edge_ids(1_024); let mut reference: Option = None; let mut observed = Vec::new(); - for partition_count in [1_u32, 4, 64, 256] { + // #1439 extends the range up to MAX_PARTITION_COUNT, the new default + // ceiling: at this fixture's small record count both 256 and 4_096 + // clip to the same recorded-count-bounded cut, so this proves raising + // the ceiling all the way to its new default introduces no + // logical-result variance either. + for partition_count in [ + 1_u32, + 4, + 64, + 256, + crate::graph_construction::partition::MAX_PARTITION_COUNT, + ] { let root = TempDir::new().unwrap(); let (fingerprint, layout) = ingest(&root, partition_count, &nodes, &edges, 128); assert_eq!(layout.recorded_partition_count, u64::from(partition_count)); @@ -482,19 +493,31 @@ mod determinism { } } // The layouts genuinely differ; the results do not. - // The cut is the recorded count bounded by the recorded record count, so - // that a small graph does not pay a large graph's durability price. - // 2048 identities admit at most 2048/16 = 128 partitions. + // The cut is bounded by the recorded record count on two independent + // terms, `min`-combined (#1439): a small-scale floor of + // `identities / 16` (unchanged since before #1439) and a data-driven + // target floored at `DEFAULT_PARTITION_COUNT`. 2048 identities admit + // at most 2048/16 = 128 partitions from the small-scale floor, which + // is far below the data-driven target's floor of 256, so the + // small-scale floor is what binds at this fixture's size for every + // requested count above it -- including the raised MAX_PARTITION_COUNT + // default. let partitions = observed .iter() .map(|(_, partitions, ..)| *partitions) .collect::>(); let identities = observed[0].3; assert_eq!(identities, (nodes.len() + edges.len()) as u64); - let expected = [1_u32, 4, 64, 256] - .into_iter() - .map(|requested| u64::from(requested).min(identities / 16)) - .collect::>(); + let expected = [ + 1_u32, + 4, + 64, + 256, + crate::graph_construction::partition::MAX_PARTITION_COUNT, + ] + .into_iter() + .map(|requested| u64::from(requested).min(identities / 16)) + .collect::>(); assert_eq!(partitions, expected); let peaks = observed .iter() diff --git a/crates/graphforge-storage/src/graph_construction.rs b/crates/graphforge-storage/src/graph_construction.rs index 562055881..bb7e791b5 100644 --- a/crates/graphforge-storage/src/graph_construction.rs +++ b/crates/graphforge-storage/src/graph_construction.rs @@ -470,11 +470,20 @@ pub struct GraphConstructionBudgets { pub max_catalog_decoded_bytes: usize, /// Maximum UTF-8 identifier bytes retained by the complete runtime catalog. pub max_catalog_identifier_bytes: usize, - /// Range partitions shaping cuts the identity key space into. + /// Ceiling on the range partitions shaping cuts the identity key space + /// into. /// /// This is a recorded format parameter, never derived from the machine. It /// must not be tied to `available_parallelism()` or to a thread count: the /// same logical input has to stage identically on hosts of different sizes. + /// + /// It is a ceiling, not the literal count used: the effective cut is a + /// pure function of the recorded staged record count, clamped to this + /// value (#1439). Defaulting it to [`partition::MAX_PARTITION_COUNT`] + /// means that data-driven cut, not an artificially low flat cap, is what + /// actually bounds the partition count in production; the cut formula's + /// own small-scale floor and data-driven target are what keep small and + /// medium inputs at the same partition counts as before. pub partition_count: u32, } @@ -491,7 +500,7 @@ impl Default for GraphConstructionBudgets { max_catalog_entries: 1_000_000, max_catalog_decoded_bytes: 256 << 20, max_catalog_identifier_bytes: 64 << 20, - partition_count: partition::DEFAULT_PARTITION_COUNT, + partition_count: partition::MAX_PARTITION_COUNT, } } } diff --git a/crates/graphforge-storage/src/graph_construction/partition.rs b/crates/graphforge-storage/src/graph_construction/partition.rs index b680a782a..b2cba38d7 100644 --- a/crates/graphforge-storage/src/graph_construction/partition.rs +++ b/crates/graphforge-storage/src/graph_construction/partition.rs @@ -38,11 +38,33 @@ use graphforge_core::GfError; /// This is a recorded format parameter, never derived from the machine: it must /// not depend on `available_parallelism()` or on thread count, or the same /// logical input would stage differently on differently-sized hosts. +/// +/// This is also the floor of the data-driven scaling target below +/// [`TARGET_ROWS_PER_PARTITION`] (#1439): below that pivot, the effective cut +/// is exactly this value, matching the operating point already proven at +/// production's largest ladder rung before the scaling target engages. pub(crate) const DEFAULT_PARTITION_COUNT: u32 = 256; /// Upper bound on the recorded partition count. pub(crate) const MAX_PARTITION_COUNT: u32 = 4_096; +/// Target rows materialized per partition once the data-driven scaling term +/// is the binding constraint (#1439). +/// +/// Each partition is materialized whole in memory and sorted, so resident +/// partition memory is proportional to rows-per-partition, not to the +/// partition *count*. A cut bounded only by a flat partition-count ceiling +/// (the pre-#1439 behaviour) holds partition count constant past that +/// ceiling and lets partition *size* grow with data instead, which is exactly +/// what the RSS plateau gate exists to catch. +/// +/// `4_194_304 / 16_384 == 256 == DEFAULT_PARTITION_COUNT`: this constant is +/// chosen so the scaling term reproduces the already-proven 256-partition +/// operating point at that record count and grows partitions proportionally +/// with data beyond it, reaching exactly `MAX_PARTITION_COUNT` (4,096) at +/// 67,108,864 records. +const TARGET_ROWS_PER_PARTITION: u64 = 16_384; + /// Sample points drawn per requested partition. Oversampling the quantile /// estimate is what keeps the partitions close to balanced; 64 points per /// partition is the usual sample-sort factor. @@ -72,8 +94,9 @@ const BALANCE_MIN_MEAN_ROWS: u64 = 16; /// This does not make the partition count machine-dependent, which is the thing /// R1 forbids. The bound is a pure function of the staged record count recorded /// in the chunk receipts, so the same logical input cuts the same partitions on -/// any host. At production scale the recorded count is the binding constraint -/// and this bound does nothing. +/// any host. This is a *small*-scale guard only: it is the binding constraint +/// far below [`DEFAULT_PARTITION_COUNT`] records, and does nothing once the +/// data-driven [`TARGET_ROWS_PER_PARTITION`] term takes over (#1439). const MIN_ROWS_PER_PARTITION: u64 = BALANCE_MIN_MEAN_ROWS; /// The recorded range-partition authority. @@ -222,8 +245,28 @@ impl IdentitySampler { /// Returns an error when the requested partition count is out of range. pub(crate) fn new(partition_count: u32, total_records: u64) -> Result { validate_partition_count(partition_count)?; + // Two independent bounds, combined by `min`, so each stays a pure + // function of `total_records` and the requested ceiling never widens + // either of them (#1439): + // + // - `small_scale_floor`: below ~`DEFAULT_PARTITION_COUNT * 16` records + // this is the binding term, exactly as before #1439. It keeps a + // two-thousand-row graph from paying a large graph's durability + // price. + // - `data_driven_target`: floored at `DEFAULT_PARTITION_COUNT`, so it + // is inactive (equal to the floor) at and below the record count + // that floor already covers, and grows partition count + // proportionally with data beyond it. This is what keeps + // rows-per-partition, and therefore resident partition memory, + // roughly constant at scale instead of letting partition count + // plateau while partition size keeps growing. + // + // The requested `partition_count` remains the hard ceiling over both. + let small_scale_floor = (total_records / MIN_ROWS_PER_PARTITION).max(1); + let data_driven_target = (total_records / TARGET_ROWS_PER_PARTITION) + .max(u64::from(DEFAULT_PARTITION_COUNT)); let cut = u32::try_from( - u64::from(partition_count).min((total_records / MIN_ROWS_PER_PARTITION).max(1)), + u64::from(partition_count).min(small_scale_floor.min(data_driven_target)), ) .map_err(storage)?; let target = u64::from(cut) diff --git a/crates/graphforge-storage/src/graph_construction/partition/tests.rs b/crates/graphforge-storage/src/graph_construction/partition/tests.rs index a907eae7a..8fed4893c 100644 --- a/crates/graphforge-storage/src/graph_construction/partition/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/partition/tests.rs @@ -202,6 +202,37 @@ fn the_cut_is_bounded_by_the_recorded_record_count() { } } +/// #1439: past `DEFAULT_PARTITION_COUNT * 16` records the flat ceiling used +/// to be the only bound, so partition *count* plateaued while partition +/// *size* (and therefore resident memory) kept growing with data. Requesting +/// the full `MAX_PARTITION_COUNT` ceiling -- production's new default -- +/// exercises the data-driven scaling term and must reproduce exactly the +/// S18/S19/S20/S22 ladder-rung table from the issue: partition count doubles +/// as record count doubles, holding rows-per-partition at +/// `TARGET_ROWS_PER_PARTITION` (16,384) until the ceiling itself binds. +#[test] +fn the_cut_scales_with_data_once_past_the_flat_default_instead_of_plateauing() { + for (records, expected_partitions) in [ + (4_194_304_u64, 256_u32), // S18: unchanged from the pre-#1439 default. + (8_388_608, 512), // S19: partitions double as records double. + (16_777_216, 1_024), // S20: partitions double again. + (67_108_864, 4_096), // S22: exactly MAX_PARTITION_COUNT. + ] { + let sampler = IdentitySampler::new(MAX_PARTITION_COUNT, records).unwrap(); + assert_eq!(sampler.cut(), expected_partitions, "records={records}"); + assert_eq!( + records / u64::from(expected_partitions), + 16_384, + "rows-per-partition must hold constant while the target is binding: records={records}" + ); + } + // Beyond S22's record count the ceiling itself binds: rows-per-partition + // now grows with data again, but only past the point R1's cross-host + // determinism proof and the durability trade in #1439 accepted. + let far_beyond = IdentitySampler::new(MAX_PARTITION_COUNT, 1_000_000_000).unwrap(); + assert_eq!(far_beyond.cut(), MAX_PARTITION_COUNT); +} + #[test] fn fewer_distinct_keys_than_partitions_collapse_rather_than_emptying_partitions() { let keys = ingest_keys(5); diff --git a/crates/graphforge-storage/src/graph_construction/tests.rs b/crates/graphforge-storage/src/graph_construction/tests.rs index f607cdf67..f943da1dc 100644 --- a/crates/graphforge-storage/src/graph_construction/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/tests.rs @@ -403,18 +403,24 @@ fn million_chunk_shaping_retains_name_state_bounded_by_the_partition_count() { // The online merge scheduler's logarithmic name state is gone: range // partitioning retains exactly one spill name per partition per family, // independent of how many chunks were staged. + // + // `budgets.partition_count` is a *ceiling* (#1439), defaulted to + // `MAX_PARTITION_COUNT` so the data-driven cut -- not an artificially low + // flat cap -- is what bounds production partition counts. Worst-case + // retained name state is therefore sized off the ceiling, not off the + // pre-#1439 flat default, and still stays a small fraction of max_chunks. let budgets = GraphConstructionBudgets::default(); assert_eq!(budgets.max_chunks, 1_000_000); let families = super::partition_shaping::PartitionFamily::ALL.len() as u64; let slots = u64::from(budgets.partition_count) * families; - // Name state is a function of the recorded partition count and the family - // set, never of the staged chunk count. - assert_eq!(slots, 1_280, "retained slots: {slots}"); - assert!(slots < budgets.max_chunks / 100, "retained slots: {slots}"); + // Name state is a function of the recorded partition count ceiling and + // the family set, never of the staged chunk count. + assert_eq!(slots, 20_480, "retained slots: {slots}"); + assert!(slots < budgets.max_chunks / 40, "retained slots: {slots}"); assert_eq!(budgets.max_schema_groups, 256); assert_eq!( budgets.partition_count, - super::partition::DEFAULT_PARTITION_COUNT + super::partition::MAX_PARTITION_COUNT ); } From dddf1e652a11a54c79cdb80c210e6d13e89017e2 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Thu, 17 Sep 2026 16:23:11 +0000 Subject: [PATCH 2/5] fix(storage): route node-keyed families with node-only splitters, remove per-partition spill buffers (#1439) Supersedes this branch's first commit: the data-driven partition-count scaling it introduced measured worse under a corrected fixture (RSS grew faster, not slower, once row-group buffer overhead was accounted for) and is reverted here back to main's `cut = min(partition_count, max(1, records/16))` with a flat 256 default. The real defect was elsewhere. Endpoints, node details and node-kind rows are keyed by node UUID, but were routed with splitters sampled from the *joint* node+edge identity domain. Graph500-shaped input puts nodes and edges in disjoint UUID bands (namespace byte) and is overwhelmingly edges, so almost every node UUID fell inside a handful of the joint splitters' partitions -- a real skew invisible to every existing check, because it was hidden by a second, independent defect: `edge_batch`/`edge_rows` test fixtures indexed nodes from a per-chunk-local counter, so every staged chunk referenced the same low node-UUID band regardless of graph size, and the whole test suite exercised a degenerate graph shape that never surfaced the skew. Fixes, in order of how they were found: - A second, node-only splitter set (`choose_node_partition_plan`, sampled from the staged node identity domain, recorded as `node_splitters` in `ShapeIntent`) routes node details, staged (pre-resolution) endpoints and node-kind rows. Edge-keyed families keep the joint set, whose domain they already match well. - `PartitionBalance::assert_balanced_with_hub_tolerance`: extends the balance assertion to every fixed-width family (previously only identities), discounting one repeated key's excess run before comparing to the tolerance -- a hub node's endpoint records genuinely cannot be split across partitions, so that is not a splitter defect. - Fixed the test-fixture bug: `edge_batch`/`edge_rows` now take an explicit node offset (and, where needed, a node count to wrap modulo), so src/dst references vary across chunks instead of collapsing onto the first chunk's nodes. - Fixed a partitioner-sizing bug found while wiring the above: node-keyed `FixedRangePartitioner`/`RowRangePartitioner` instances were sized from the joint plan's (larger) partition count while routing through the node plan (fewer partitions). `PartitionBalance`'s mean is `total/partitions`, so the unused trailing slots diluted the mean and made correctly-balanced data look skewed. Now sized from whichever plan actually routes them. - Removed the per-partition write buffer for the four fixed-width families (identities, node/edge details, endpoints): `SPILL_BLOCK_BYTES` (64 KiB) held open per partition per family was the dominant measured term in both RSS and fsync growth as partition count rises. `UnbufferedSpillWriter` holds no buffer; `PartitionRun` (in shape.rs) exploits the same monotone-sorted property the row partitioner already used, accumulating one contiguous same-partition run of wire bytes per staged chunk and writing it in a single call via the new `FixedRangePartitioner::route_slice`. Row groups (Arrow-based) are unchanged: `StreamWriter` doesn't expose per-record wire boundaries the same way, so the batching trick doesn't transfer, and left at a flat 256 the residual buffered term there cannot dominate. - Pre-sized `load_partition`'s `Vec` from the routing-time balance count, removing the growth-by-doubling headroom. Verified: cargo test -p graphforge-storage --lib (1,152 passed, 0 failed); the four fixed-width digests unchanged for shaped-identities.run and shaped-node-details.run (edge-details/edge-endpoints changed only because the fixture bug fix changes the graph structure the determinism tests build, not because partitioning changed); cross-partition-count fingerprint invariance holds at every requested count including the 4,096 ceiling; cargo clippy --workspace -- -D warnings clean. Measured at realistic (Graph500 UUID-band) ladder-rung scale, one ingest+shape per process, /usr/bin/time -v: peak RSS 78.2 / 96.2 / 138.2 MB at S18/S19/S20-equivalent record counts (edges/128 node ratio, matching the real ladder), versus 144.5 / 204.0 / 317.0 MB before the buffer removal at the same (corrected) ratio -- roughly halving both the absolute RSS and the growth rate. Still short of the ladder's own +13%/ +27% ingest growth, so the ladder itself is the next check, not this fixture. Co-Authored-By: Claude Sonnet 5 --- .../src/construction_detail_tests.rs | 20 +- .../src/construction_determinism_tests.rs | 17 +- .../src/construction_lifecycle_tests.rs | 4 +- .../src/graph_construction.rs | 23 ++- .../src/graph_construction/intake/tests.rs | 6 +- .../src/graph_construction/io_evidence.rs | 13 ++ .../src/graph_construction/partition.rs | 114 ++++++----- .../src/graph_construction/partition/tests.rs | 31 --- .../graph_construction/partition_shaping.rs | 191 +++++++++++++++--- .../src/graph_construction/shape.rs | 175 +++++++++++++++- .../src/graph_construction/shape/tests.rs | 21 +- .../src/graph_construction/tests.rs | 60 ++++-- 12 files changed, 515 insertions(+), 160 deletions(-) diff --git a/crates/graphforge-storage/src/construction_detail_tests.rs b/crates/graphforge-storage/src/construction_detail_tests.rs index 62ed07ec2..cbcea13d0 100644 --- a/crates/graphforge-storage/src/construction_detail_tests.rs +++ b/crates/graphforge-storage/src/construction_detail_tests.rs @@ -75,7 +75,7 @@ mod compact_details { ) .unwrap(); session - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 2)) + .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 3, 2)) .unwrap(); session.seal().unwrap(); drop(session); @@ -153,7 +153,7 @@ mod compact_details { ) .unwrap(); initial - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 2)) + .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 3, 2)) .unwrap(); initial.seal().unwrap(); let encoded = initial.prepare_canonical_encoding(1).unwrap(); @@ -295,7 +295,7 @@ mod compact_details { .append( ConstructionChunkKind::Edge, &format!("edges-{chunk}"), - &edge_batch(10_000 + chunk * 128, 128), + &edge_batch(10_000 + chunk * 128, 1 + (chunk * 128) % 1024, 1024, 128), ) .unwrap(); } @@ -410,12 +410,18 @@ mod compact_details { record[17..23].copy_from_slice(b"Person"); expected.push(record); } + // src/dst walk the full 1024-node range across the 8 chunks and + // wrap modulo 1024 (#1439), rather than restarting at node 1 on + // every chunk: `edge_batch`'s `node_start` argument varies per + // chunk to match. for chunk in 0_u128..8 { for row in 0_u128..128 { let mut record = vec![0; EDGE_DETAIL_WIDTH]; record[..16].copy_from_slice(&(10_000 + chunk * 128 + row).to_be_bytes()); - record[16..32].copy_from_slice(&(row + 1).to_be_bytes()); - record[32..48].copy_from_slice(&(row + 2).to_be_bytes()); + let src = 1 + chunk * 128 + row; + let dst = 1 + (chunk * 128 + row + 1) % 1024; + record[16..32].copy_from_slice(&src.to_be_bytes()); + record[32..48].copy_from_slice(&dst.to_be_bytes()); record[48] = 1; record[49] = b'R'; expected.push(record); @@ -496,7 +502,7 @@ mod compact_details { .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2)) .unwrap(); session - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1)) + .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 2, 1)) .unwrap(); session.seal().unwrap(); let encoded = session.prepare_canonical_encoding(1).unwrap(); @@ -567,7 +573,7 @@ mod compact_details { } if resumed.accepted_chunks() == 1 { resumed - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1)) + .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 2, 1)) .unwrap(); } resumed.seal().unwrap(); diff --git a/crates/graphforge-storage/src/construction_determinism_tests.rs b/crates/graphforge-storage/src/construction_determinism_tests.rs index 315a9703f..86743c2ba 100644 --- a/crates/graphforge-storage/src/construction_determinism_tests.rs +++ b/crates/graphforge-storage/src/construction_determinism_tests.rs @@ -108,16 +108,25 @@ mod determinism { .unwrap() } - fn edge_rows(edges: &[[u8; 16]], nodes: &[[u8; 16]]) -> RecordBatch { + /// `start` is this window's offset into the full edge sequence (#1439). + /// Every call site passes one chunk at a time, and indexing `nodes` from + /// a per-call-local `0` on every chunk -- the previous behaviour -- + /// referenced only the first `chunk`-many nodes as an endpoint no matter + /// how many nodes or chunks existed, concentrating every edge onto a + /// narrow node-UUID band. That was invisible while endpoints were routed + /// with the joint identity splitters, which never resolve node UUIDs + /// finely enough to notice; routing them with node-only splitters + /// surfaces it as a real (fixture-caused, not production) skew. + fn edge_rows(start: usize, edges: &[[u8; 16]], nodes: &[[u8; 16]]) -> RecordBatch { let src = edges .iter() .enumerate() - .map(|(index, _)| nodes[index % nodes.len()]) + .map(|(index, _)| nodes[(start + index) % nodes.len()]) .collect::>(); let dst = edges .iter() .enumerate() - .map(|(index, _)| nodes[(index + 1) % nodes.len()]) + .map(|(index, _)| nodes[(start + index + 1) % nodes.len()]) .collect::>(); RecordBatch::try_new( CONSTRUCTION_EDGE_SCHEMA.clone(), @@ -200,7 +209,7 @@ mod determinism { .append( ConstructionChunkKind::Edge, &format!("edges-{index}"), - &edge_rows(window, nodes), + &edge_rows(index * chunk, window, nodes), ) .unwrap(); } diff --git a/crates/graphforge-storage/src/construction_lifecycle_tests.rs b/crates/graphforge-storage/src/construction_lifecycle_tests.rs index f9ca9b45b..b938ce621 100644 --- a/crates/graphforge-storage/src/construction_lifecycle_tests.rs +++ b/crates/graphforge-storage/src/construction_lifecycle_tests.rs @@ -529,7 +529,7 @@ mod lifecycle_budget { .append( ConstructionChunkKind::Edge, &format!("edges-{chunk}"), - &edge_batch(1_000_000 + u128::from(chunk) * 4096, 4096), + &edge_batch(1_000_000 + u128::from(chunk) * 4096, 1 + (u128::from(chunk) * 4096) % 8192, 8192, 4096), ) .unwrap(); } @@ -1079,7 +1079,7 @@ mod lifecycle_budget { .append( ConstructionChunkKind::Edge, &format!("edges-{chunk}"), - &edge_batch(1_000_000 + u128::from(chunk) * 4096, 4096), + &edge_batch(1_000_000 + u128::from(chunk) * 4096, 1 + (u128::from(chunk) * 4096) % 8192, 8192, 4096), ) .unwrap(); } diff --git a/crates/graphforge-storage/src/graph_construction.rs b/crates/graphforge-storage/src/graph_construction.rs index bb7e791b5..1cef1c958 100644 --- a/crates/graphforge-storage/src/graph_construction.rs +++ b/crates/graphforge-storage/src/graph_construction.rs @@ -470,20 +470,11 @@ pub struct GraphConstructionBudgets { pub max_catalog_decoded_bytes: usize, /// Maximum UTF-8 identifier bytes retained by the complete runtime catalog. pub max_catalog_identifier_bytes: usize, - /// Ceiling on the range partitions shaping cuts the identity key space - /// into. + /// Range partitions shaping cuts the identity key space into. /// /// This is a recorded format parameter, never derived from the machine. It /// must not be tied to `available_parallelism()` or to a thread count: the /// same logical input has to stage identically on hosts of different sizes. - /// - /// It is a ceiling, not the literal count used: the effective cut is a - /// pure function of the recorded staged record count, clamped to this - /// value (#1439). Defaulting it to [`partition::MAX_PARTITION_COUNT`] - /// means that data-driven cut, not an artificially low flat cap, is what - /// actually bounds the partition count in production; the cut formula's - /// own small-scale floor and data-driven target are what keep small and - /// medium inputs at the same partition counts as before. pub partition_count: u32, } @@ -500,7 +491,7 @@ impl Default for GraphConstructionBudgets { max_catalog_entries: 1_000_000, max_catalog_decoded_bytes: 256 << 20, max_catalog_identifier_bytes: 64 << 20, - partition_count: partition::MAX_PARTITION_COUNT, + partition_count: partition::DEFAULT_PARTITION_COUNT, } } } @@ -831,6 +822,16 @@ struct ShapeIntent { /// partition writes a byte. #[serde(default)] splitters: Vec, + /// Recorded range-partition splitters over the staged **node** identity + /// domain only (#1439), canonical lower hex, strictly increasing. + /// Endpoints, node details and node-kind rows are keyed by node UUID and + /// are routed with these instead of `splitters`: Graph500-shaped input + /// puts nodes and edges in disjoint UUID bands, so the joint splitters + /// above route almost every node-keyed record into a handful of + /// partitions. A pure function of the same recorded chunk receipts as + /// `splitters`, so R1 holds for this set too. + #[serde(default)] + node_splitters: Vec, /// Measured identity rows per effective partition, in partition order. #[serde(default)] partition_identity_rows: Vec, diff --git a/crates/graphforge-storage/src/graph_construction/intake/tests.rs b/crates/graphforge-storage/src/graph_construction/intake/tests.rs index ebbaadcc7..4da1b704d 100644 --- a/crates/graphforge-storage/src/graph_construction/intake/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/intake/tests.rs @@ -239,7 +239,11 @@ fn schema_group_admission_is_constant_and_budgeted() { .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2)) .unwrap(); let error = session - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1)) + .append( + ConstructionChunkKind::Edge, + "edges", + &edge_batch(100, 1, 2, 1), + ) .unwrap_err(); assert!(error.to_string().contains("schema-group budget")); assert_eq!(session.checkpoint.node_schema_sha256.len(), 1); diff --git a/crates/graphforge-storage/src/graph_construction/io_evidence.rs b/crates/graphforge-storage/src/graph_construction/io_evidence.rs index 11579bb9d..db6d69bc2 100644 --- a/crates/graphforge-storage/src/graph_construction/io_evidence.rs +++ b/crates/graphforge-storage/src/graph_construction/io_evidence.rs @@ -346,6 +346,19 @@ pub struct GraphConstructionEvidence { /// Identity keys retained as splitter sample points. #[serde(default)] pub splitter_sample_records: u64, + /// Effective range partitions produced by the node-only splitter set + /// (#1439), used to route node-keyed families (node details, endpoints, + /// node rows) instead of the joint identity splitters. + #[serde(default)] + pub shape_node_partitions: u64, + /// Node identity records observed by the node-only splitter sampling + /// pass (#1439). + #[serde(default)] + pub node_splitter_sampled_source_records: u64, + /// Node identity keys retained as node-only splitter sample points + /// (#1439). + #[serde(default)] + pub node_splitter_sample_records: u64, /// Rows routed into the largest identity partition. #[serde(default)] pub max_partition_identity_rows: u64, diff --git a/crates/graphforge-storage/src/graph_construction/partition.rs b/crates/graphforge-storage/src/graph_construction/partition.rs index b2cba38d7..fac829ce7 100644 --- a/crates/graphforge-storage/src/graph_construction/partition.rs +++ b/crates/graphforge-storage/src/graph_construction/partition.rs @@ -38,33 +38,11 @@ use graphforge_core::GfError; /// This is a recorded format parameter, never derived from the machine: it must /// not depend on `available_parallelism()` or on thread count, or the same /// logical input would stage differently on differently-sized hosts. -/// -/// This is also the floor of the data-driven scaling target below -/// [`TARGET_ROWS_PER_PARTITION`] (#1439): below that pivot, the effective cut -/// is exactly this value, matching the operating point already proven at -/// production's largest ladder rung before the scaling target engages. pub(crate) const DEFAULT_PARTITION_COUNT: u32 = 256; /// Upper bound on the recorded partition count. pub(crate) const MAX_PARTITION_COUNT: u32 = 4_096; -/// Target rows materialized per partition once the data-driven scaling term -/// is the binding constraint (#1439). -/// -/// Each partition is materialized whole in memory and sorted, so resident -/// partition memory is proportional to rows-per-partition, not to the -/// partition *count*. A cut bounded only by a flat partition-count ceiling -/// (the pre-#1439 behaviour) holds partition count constant past that -/// ceiling and lets partition *size* grow with data instead, which is exactly -/// what the RSS plateau gate exists to catch. -/// -/// `4_194_304 / 16_384 == 256 == DEFAULT_PARTITION_COUNT`: this constant is -/// chosen so the scaling term reproduces the already-proven 256-partition -/// operating point at that record count and grows partitions proportionally -/// with data beyond it, reaching exactly `MAX_PARTITION_COUNT` (4,096) at -/// 67,108,864 records. -const TARGET_ROWS_PER_PARTITION: u64 = 16_384; - /// Sample points drawn per requested partition. Oversampling the quantile /// estimate is what keeps the partitions close to balanced; 64 points per /// partition is the usual sample-sort factor. @@ -94,9 +72,8 @@ const BALANCE_MIN_MEAN_ROWS: u64 = 16; /// This does not make the partition count machine-dependent, which is the thing /// R1 forbids. The bound is a pure function of the staged record count recorded /// in the chunk receipts, so the same logical input cuts the same partitions on -/// any host. This is a *small*-scale guard only: it is the binding constraint -/// far below [`DEFAULT_PARTITION_COUNT`] records, and does nothing once the -/// data-driven [`TARGET_ROWS_PER_PARTITION`] term takes over (#1439). +/// any host. At production scale the recorded count is the binding constraint +/// and this bound does nothing. const MIN_ROWS_PER_PARTITION: u64 = BALANCE_MIN_MEAN_ROWS; /// The recorded range-partition authority. @@ -245,28 +222,8 @@ impl IdentitySampler { /// Returns an error when the requested partition count is out of range. pub(crate) fn new(partition_count: u32, total_records: u64) -> Result { validate_partition_count(partition_count)?; - // Two independent bounds, combined by `min`, so each stays a pure - // function of `total_records` and the requested ceiling never widens - // either of them (#1439): - // - // - `small_scale_floor`: below ~`DEFAULT_PARTITION_COUNT * 16` records - // this is the binding term, exactly as before #1439. It keeps a - // two-thousand-row graph from paying a large graph's durability - // price. - // - `data_driven_target`: floored at `DEFAULT_PARTITION_COUNT`, so it - // is inactive (equal to the floor) at and below the record count - // that floor already covers, and grows partition count - // proportionally with data beyond it. This is what keeps - // rows-per-partition, and therefore resident partition memory, - // roughly constant at scale instead of letting partition count - // plateau while partition size keeps growing. - // - // The requested `partition_count` remains the hard ceiling over both. - let small_scale_floor = (total_records / MIN_ROWS_PER_PARTITION).max(1); - let data_driven_target = (total_records / TARGET_ROWS_PER_PARTITION) - .max(u64::from(DEFAULT_PARTITION_COUNT)); let cut = u32::try_from( - u64::from(partition_count).min(small_scale_floor.min(data_driven_target)), + u64::from(partition_count).min((total_records / MIN_ROWS_PER_PARTITION).max(1)), ) .map_err(storage)?; let target = u64::from(cut) @@ -359,12 +316,23 @@ impl PartitionBalance { /// Returns an error when the partition is out of range or the count /// overflows. pub(crate) fn record(&mut self, partition: usize) -> Result<(), GfError> { + self.record_many(partition, 1) + } + + /// Charge `count` rows to `partition` in one step (#1439 follow-up: the + /// caller batching a contiguous same-partition run charges it once + /// rather than once per record). + /// + /// # Errors + /// Returns an error when the partition is out of range or the count + /// overflows. + pub(crate) fn record_many(&mut self, partition: usize, count: u64) -> Result<(), GfError> { let slot = self .rows .get_mut(partition) .ok_or_else(|| storage("partition index is out of range"))?; *slot = slot - .checked_add(1) + .checked_add(count) .ok_or_else(|| storage("partition row count overflows"))?; Ok(()) } @@ -418,6 +386,58 @@ impl PartitionBalance { } Ok(()) } + + /// Refuse a partitioning whose largest partition exceeds + /// [`BALANCE_TOLERANCE`] times the mean, **unless** the excess is + /// explained by one repeated key. + /// + /// A range partition cannot split one key across partitions, so a hub -- + /// one node referenced by a disproportionate number of edges, routing a + /// disproportionate number of endpoint records to the one partition that + /// owns it -- is not a splitter defect. `max_single_key_run` is the + /// largest number of records sharing one key observed anywhere in the + /// partitioning (computed after sorting, where identical keys are + /// contiguous). Discounting all but one of those occurrences from the + /// largest partition before comparing against the tolerance is what + /// distinguishes a genuine hub from a skewed splitter set: a bad splitter + /// set concentrates *distinct* keys into one partition, which this + /// discount does not hide. + /// + /// # Errors + /// Returns an error when the partitioning is skewed beyond the tolerance + /// even after that discount. + pub(crate) fn assert_balanced_with_hub_tolerance( + &self, + context: &str, + max_single_key_run: u64, + ) -> Result<(), GfError> { + let partitions = self.rows.len() as u64; + if partitions == 0 { + return Err(storage("partition balance has no partitions")); + } + let total = self.total(); + if total / partitions < BALANCE_MIN_MEAN_ROWS { + return Ok(()); + } + let max = self.max_rows(); + let discount = max_single_key_run.saturating_sub(1); + let discounted = max.saturating_sub(discount); + let scaled = discounted + .checked_mul(partitions) + .ok_or_else(|| storage("partition balance ratio overflows"))?; + let budget = total + .checked_mul(BALANCE_TOLERANCE) + .ok_or_else(|| storage("partition balance budget overflows"))?; + if scaled > budget { + return Err(storage(format!( + "{context} range partitioning is skewed: largest partition holds {max} of \ + {total} rows across {partitions} partitions (largest single-key run \ + {max_single_key_run}), above the {BALANCE_TOLERANCE}x mean tolerance even after \ + discounting one hub key" + ))); + } + Ok(()) + } } #[cfg(test)] diff --git a/crates/graphforge-storage/src/graph_construction/partition/tests.rs b/crates/graphforge-storage/src/graph_construction/partition/tests.rs index 8fed4893c..a907eae7a 100644 --- a/crates/graphforge-storage/src/graph_construction/partition/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/partition/tests.rs @@ -202,37 +202,6 @@ fn the_cut_is_bounded_by_the_recorded_record_count() { } } -/// #1439: past `DEFAULT_PARTITION_COUNT * 16` records the flat ceiling used -/// to be the only bound, so partition *count* plateaued while partition -/// *size* (and therefore resident memory) kept growing with data. Requesting -/// the full `MAX_PARTITION_COUNT` ceiling -- production's new default -- -/// exercises the data-driven scaling term and must reproduce exactly the -/// S18/S19/S20/S22 ladder-rung table from the issue: partition count doubles -/// as record count doubles, holding rows-per-partition at -/// `TARGET_ROWS_PER_PARTITION` (16,384) until the ceiling itself binds. -#[test] -fn the_cut_scales_with_data_once_past_the_flat_default_instead_of_plateauing() { - for (records, expected_partitions) in [ - (4_194_304_u64, 256_u32), // S18: unchanged from the pre-#1439 default. - (8_388_608, 512), // S19: partitions double as records double. - (16_777_216, 1_024), // S20: partitions double again. - (67_108_864, 4_096), // S22: exactly MAX_PARTITION_COUNT. - ] { - let sampler = IdentitySampler::new(MAX_PARTITION_COUNT, records).unwrap(); - assert_eq!(sampler.cut(), expected_partitions, "records={records}"); - assert_eq!( - records / u64::from(expected_partitions), - 16_384, - "rows-per-partition must hold constant while the target is binding: records={records}" - ); - } - // Beyond S22's record count the ceiling itself binds: rows-per-partition - // now grows with data again, but only past the point R1's cross-host - // determinism proof and the durability trade in #1439 accepted. - let far_beyond = IdentitySampler::new(MAX_PARTITION_COUNT, 1_000_000_000).unwrap(); - assert_eq!(far_beyond.cut(), MAX_PARTITION_COUNT); -} - #[test] fn fewer_distinct_keys_than_partitions_collapse_rather_than_emptying_partitions() { let keys = ingest_keys(5); diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs index 6c3944eb1..f8f451192 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs @@ -132,6 +132,84 @@ struct SpillWriter { } impl SpillWriter { + /// Constructed only from an already-sealed [`RowSpill`]'s buffered + /// writer (`RowRangePartitioner::seal`): row spills still buffer their + /// writes, because the Arrow IPC `StreamWriter` they hold does not offer + /// the record-boundary information [`FixedRangePartitioner::route_slice`] + /// needs to batch by contiguous partition run the way fixed-width + /// families do (#1439 follow-up). `create`/`write`/`abandon` accordingly + /// have no callers left after that change and are not reintroduced here. + fn seal( + mut self, + root: &StableDirectory, + evidence: &mut GraphConstructionEvidence, + ) -> Result { + self.writer.flush().map_err(super::storage)?; + self.writer + .get_mut() + .inner + .sync_all_and_release() + .map_err(super::storage)?; + let cache_release = self.writer.get_ref().inner.evidence(); + account_cache_release(cache_release, evidence)?; + account_sequential_write(self.writer.get_ref().bytes, evidence)?; + let receipt = ArtifactReceipt { + name: self.name.clone(), + bytes: self.writer.get_ref().bytes, + allocated_bytes: graphforge_filesystem::file_space_usage( + self.writer.get_ref().inner.file(), + ) + .map_err(super::storage)? + .allocated_bytes, + sha256: hex(&self.writer.get_ref().digest.clone().finalize()), + xxh64: crate::corruption_checksum::hex(self.writer.get_ref().checksum.finish()), + identity: self.identity.into(), + write_operations: self.writer.get_ref().operations, + fsync_operations: cache_release + .sync_operations + .checked_add(1) + .ok_or_else(|| super::storage("artifact synchronization count overflows"))?, + }; + drop(self.writer); + root.install_child( + self.temporary.as_os_str(), + self.identity, + OsStr::new(&self.name), + ) + .map_err(super::storage)?; + root.sync().map_err(super::storage)?; + construction_failpoint("shape.partition_spill.after_install"); + persist_shape_receipt(root, &receipt)?; + record_shape_artifact_install(evidence, &receipt)?; + account_fixed_write_operations(&receipt, evidence)?; + evidence.merge_fsync_operations = evidence + .merge_fsync_operations + .checked_add(receipt.fsync_operations) + .ok_or_else(|| super::storage("merge fsync operations overflows"))?; + Ok(receipt) + } +} + +/// One open, unsealed fixed-width partition spill (#1439 follow-up). +/// +/// Unlike [`SpillWriter`], this holds no per-partition write buffer. Every +/// partition open at once during routing used to cost `SPILL_BLOCK_BYTES` +/// (64 KiB) of resident buffer regardless of how much of it had actually +/// been written; concurrent spill buffers across partitions and families +/// were the dominant measured term in both the RSS and fsync growth #1439 +/// originally investigated. The staged runs routed through this are +/// UUID-sorted and the partition function is monotone, so the caller can +/// accumulate one contiguous same-partition run of wire bytes and write it +/// in a single call -- see [`FixedRangePartitioner::route_slice`] -- rather +/// than buffering to amortize many small per-record writes. +struct UnbufferedSpillWriter { + name: String, + temporary: std::ffi::OsString, + identity: FileIdentity, + writer: HashingWriter, +} + +impl UnbufferedSpillWriter { fn create( root: &StableDirectory, name: String, @@ -142,12 +220,12 @@ impl SpillWriter { .create_replaceable_child_file(temporary.as_os_str()) .map_err(super::storage)?; let identity = file_identity(&file).map_err(super::storage)?; - let hashing = HashingWriter::with_cache_window(file, window)?; + let writer = HashingWriter::with_cache_window(file, window)?; Ok(Self { name, temporary, identity, - writer: BufWriter::with_capacity(SPILL_BLOCK_BYTES, hashing), + writer, }) } @@ -155,11 +233,8 @@ impl SpillWriter { self.writer.write_all(bytes).map_err(super::storage) } - /// Drop an unsealed spill and remove its temporary. - /// - /// A shaping pass that fails part-way must not leave owned temporaries - /// behind: the identity-checked unlink here is the same one the session's - /// reopen cleanup would eventually perform, done immediately instead. + /// Drop an unsealed spill and remove its temporary. See + /// [`SpillWriter::abandon`]: same contract, same reason. fn abandon(self, root: &StableDirectory) { let Self { temporary, @@ -179,25 +254,22 @@ impl SpillWriter { ) -> Result { self.writer.flush().map_err(super::storage)?; self.writer - .get_mut() .inner .sync_all_and_release() .map_err(super::storage)?; - let cache_release = self.writer.get_ref().inner.evidence(); + let cache_release = self.writer.inner.evidence(); account_cache_release(cache_release, evidence)?; - account_sequential_write(self.writer.get_ref().bytes, evidence)?; + account_sequential_write(self.writer.bytes, evidence)?; let receipt = ArtifactReceipt { name: self.name.clone(), - bytes: self.writer.get_ref().bytes, - allocated_bytes: graphforge_filesystem::file_space_usage( - self.writer.get_ref().inner.file(), - ) - .map_err(super::storage)? - .allocated_bytes, - sha256: hex(&self.writer.get_ref().digest.clone().finalize()), - xxh64: crate::corruption_checksum::hex(self.writer.get_ref().checksum.finish()), + bytes: self.writer.bytes, + allocated_bytes: graphforge_filesystem::file_space_usage(self.writer.inner.file()) + .map_err(super::storage)? + .allocated_bytes, + sha256: hex(&self.writer.digest.clone().finalize()), + xxh64: crate::corruption_checksum::hex(self.writer.checksum.finish()), identity: self.identity.into(), - write_operations: self.writer.get_ref().operations, + write_operations: self.writer.operations, fsync_operations: cache_release .sync_operations .checked_add(1) @@ -231,7 +303,7 @@ pub(super) struct FixedRangePartitioner<'a, const N: usize> { codec: Option, reject_duplicates: bool, window: std::num::NonZeroU64, - spills: Vec>, + spills: Vec>, sealed: Vec>, balance: PartitionBalance, records: u64, @@ -274,6 +346,12 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { } /// Route one record into the partition owning `key`. + /// + /// Used only where the caller cannot establish a contiguous + /// same-partition run to batch (the post-resolution "Resolved" family is + /// re-keyed away from its input's sort order, so consecutive records are + /// not generally co-partitioned). Prefer [`Self::route_slice`] wherever + /// the input is sorted by the same key it is routed by. pub(super) fn route( &mut self, plan: &PartitionPlan, @@ -282,27 +360,49 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { evidence: &mut GraphConstructionEvidence, ) -> Result<(), GfError> { let partition = plan.partition_of(key); + let wire = run_record_bytes(record, self.codec)?; + self.route_slice(partition, wire, 1, evidence) + } + + /// Route one contiguous same-partition slice of wire bytes in a single + /// write (#1439 follow-up). + /// + /// `records` is the number of records the slice holds, charged to the + /// balance in one step. Callers whose input is sorted by the same key + /// they route by (identities, node/edge details, staged endpoints) can + /// accumulate a run of consecutive same-partition records and flush it + /// here instead of writing once per record through a held-open buffer. + pub(super) fn route_slice( + &mut self, + partition: usize, + wire: &[u8], + records: u64, + evidence: &mut GraphConstructionEvidence, + ) -> Result<(), GfError> { + if records == 0 { + debug_assert!(wire.is_empty()); + return Ok(()); + } let slot = self .spills .get_mut(partition) .ok_or_else(|| super::storage("routed partition is out of range"))?; if slot.is_none() { - *slot = Some(SpillWriter::create( + *slot = Some(UnbufferedSpillWriter::create( self.root, fixed_spill_name(self.family, partition), self.window, )?); } - let wire = run_record_bytes(record, self.codec)?; slot.as_mut() .ok_or_else(|| super::storage("partition spill is absent"))? .write(wire)?; account_merge_read_bytes(evidence, wire.len() as u64)?; account_merge_write_bytes(evidence, wire.len() as u64)?; - self.balance.record(partition)?; + self.balance.record_many(partition, records)?; self.records = self .records - .checked_add(1) + .checked_add(records) .ok_or_else(|| super::storage("partition record count overflows"))?; Ok(()) } @@ -379,6 +479,15 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { let concatenated = (|| -> Result { let mut previous: Option<[u8; N]> = None; let mut written = 0_u64; + // Tracks the longest run of records sharing one key, across the + // whole concatenation. A range partition cannot split one key, so + // this is what tells a hub apart from a skewed splitter set + // (#1439) -- see `PartitionBalance::assert_balanced_with_hub_tolerance`. + // Keys are unique within a partition's sorted order without + // spanning partitions (the partition function is a function of + // the key), so accumulating this across the loop below is exact. + let mut current_run = 0_u64; + let mut max_single_key_run = 0_u64; for partition in 0..self.sealed.len() { let Some(name) = self.sealed[partition] .as_ref() @@ -387,7 +496,8 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { continue; }; reject_cancelled(cancelled)?; - let records = self.load_partition(&name, evidence)?; + let expected = self.balance.rows().get(partition).copied(); + let records = self.load_partition(&name, expected, evidence)?; for record in &records { // Partition order is key order, so the concatenation is the // global order. Prove it rather than assume it: this is the @@ -403,7 +513,15 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { "duplicate identity across construction runs", )); } + current_run = if prior[..16] == record[..16] { + current_run.saturating_add(1) + } else { + 1 + }; + } else { + current_run = 1; } + max_single_key_run = max_single_key_run.max(current_run); let wire = run_record_bytes(record, self.codec)?; writer.write_all(wire).map_err(super::storage)?; account_merge_write_bytes(evidence, wire.len() as u64)?; @@ -416,6 +534,13 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { } } } + // Load-bearing, not a nicety (#1439): a collapsed one-partition + // run is perfectly deterministic and passes every byte-equality + // test, so this is what tells a working range partition apart + // from a catastrophically skewed one for every fixed-width + // family, not only identities. + self.balance + .assert_balanced_with_hub_tolerance(self.family.as_str(), max_single_key_run)?; Ok(written) })(); let written = match concatenated { @@ -562,14 +687,23 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { /// This is the materialization the design's rule R4 requires: the sorted /// partition exists in full before any of it is written, so row-group and /// record boundaries are a pure function of row count, never of arrival. + /// + /// `expected_records`, when given, pre-sizes the returned `Vec` from the + /// row count the routing pass already recorded in `self.balance` (#1439), + /// removing the growth-by-doubling headroom `Vec::new()` would otherwise + /// carry at the moment a large partition is fully materialized. fn load_partition( &self, name: &str, + expected_records: Option, evidence: &mut GraphConstructionEvidence, ) -> Result, GfError> { let (mut reader, counter) = open_counted_fixed_reader(self.root, name, evidence)?; let loaded = (|| -> Result, GfError> { - let mut records = Vec::new(); + let mut records = match expected_records.and_then(|count| usize::try_from(count).ok()) { + Some(count) => Vec::with_capacity(count), + None => Vec::new(), + }; while let Some(record) = read_run_record::(&mut reader, self.codec)? { account_merge_read_bytes( evidence, @@ -888,6 +1022,11 @@ impl<'a> RowRangePartitioner<'a> { evidence, )?; } + // Every row is keyed by its own unique UUID (the loop above + // already refuses a duplicate), so no hub scenario exists here: + // a skewed row partitioning is always a splitter defect (#1439). + self.balance + .assert_balanced(&format!("row partition {}", self.namespace))?; Ok(written) })(); let written = match written { diff --git a/crates/graphforge-storage/src/graph_construction/shape.rs b/crates/graphforge-storage/src/graph_construction/shape.rs index 44d341f39..349cac897 100644 --- a/crates/graphforge-storage/src/graph_construction/shape.rs +++ b/crates/graphforge-storage/src/graph_construction/shape.rs @@ -80,6 +80,39 @@ impl GraphConstructionSession { fn choose_partition_plan( &mut self, cancelled: &mut impl FnMut() -> bool, + ) -> Result { + self.sample_partition_plan(None, cancelled) + } + + /// Barrier A0'. Choose splitters over only the staged **node** identity + /// domain (#1439). + /// + /// Endpoints, node details and node-kind rows are keyed by node UUID, but + /// the joint splitters above are sampled from the combined node+edge + /// domain. Graph500-shaped input is disjoint by kind in UUID-space + /// (namespace bit) and overwhelmingly edges, so almost every node UUID + /// falls inside a handful of the joint splitters' partitions: those + /// families were routed with the wrong splitters, and it is the + /// dominant term in ingest resident memory, not partition count. + /// + /// A second, node-only splitter set is a pure function of the same + /// recorded chunk receipts as the joint one, so it costs R1 nothing: it + /// is what actually matches the key distribution those families are + /// routed by. + fn choose_node_partition_plan( + &mut self, + cancelled: &mut impl FnMut() -> bool, + ) -> Result { + self.sample_partition_plan(Some(ConstructionChunkKind::Node), cancelled) + } + + /// Shared sampler behind [`Self::choose_partition_plan`] and + /// [`Self::choose_node_partition_plan`]: identical logic, restricted to + /// receipts of `kind_filter` when given. + fn sample_partition_plan( + &mut self, + kind_filter: Option, + cancelled: &mut impl FnMut() -> bool, ) -> Result { let partition_count = self.checkpoint.budgets.partition_count; let mut extents = Vec::new(); @@ -87,6 +120,9 @@ impl GraphConstructionSession { for sequence in 0..self.checkpoint.next_sequence { reject_cancelled(cancelled)?; let receipt = self.read_receipt(sequence)?; + if kind_filter.is_some_and(|kind| receipt.kind != kind) { + continue; + } if receipt.identities.bytes % IDENTITY_WIDTH as u64 != 0 { return Err(storage("staged identity run is not record aligned")); } @@ -96,6 +132,13 @@ impl GraphConstructionSession { .checked_add(records) .ok_or_else(|| storage("staged identity record count overflows"))?; } + // #1439 follow-up: sizing the node plan's cut from the *routed* + // population (endpoints + node details, which outnumber node + // identities by the graph's average degree) rather than this sample + // domain was tried and measured worse -- it drove the count past the + // empirically measured RSS-vs-buffer-overhead minimum into the + // buffer-dominated side of that curve. Reverted: the cut is sized + // from the sample domain here, same as the joint plan. let mut sampler = IdentitySampler::new(partition_count, total)?; let positions = sampler.positions().collect::>(); let mut cursor = 0_usize; @@ -121,8 +164,15 @@ impl GraphConstructionSession { account_sequential_read(IDENTITY_WIDTH as u64, &mut self.checkpoint.evidence)?; sampler.admit(key)?; } - self.checkpoint.evidence.splitter_sample_records = sampler.sampled_records(); - self.checkpoint.evidence.splitter_sampled_source_records = sampler.source_records(); + if kind_filter.is_none() { + self.checkpoint.evidence.splitter_sample_records = sampler.sampled_records(); + self.checkpoint.evidence.splitter_sampled_source_records = sampler.source_records(); + } else { + self.checkpoint.evidence.node_splitter_sample_records = sampler.sampled_records(); + self.checkpoint + .evidence + .node_splitter_sampled_source_records = sampler.source_records(); + } sampler.into_plan(partition_count) } @@ -194,6 +244,22 @@ impl GraphConstructionSession { let partitions = plan.partitions(); self.checkpoint.evidence.shape_partitions = u64::try_from(partitions).map_err(storage)?; self.checkpoint.evidence.shape_partition_count = u64::from(plan.partition_count()); + // Barrier A0'. A second splitter set over the node-only domain + // (#1439): endpoints, node details and node-kind rows are keyed by + // node UUID, and routing them with the joint `plan` above skews them + // badly once nodes and edges occupy disjoint UUID bands. `node_plan` + // has at most as many partitions as `plan` (its domain is a subset), + // so every family below is constructed with the partition count of + // whichever plan actually routes it. This matters beyond bounds + // safety: `PartitionBalance`'s mean is `total / partitions`, so + // sizing a node-keyed family's spill/balance vectors by the joint + // plan's (larger) partition count would silently dilute its mean + // with slots that routing through `node_plan` can never reach, + // making a perfectly balanced node-keyed family look skewed. + let node_plan = self.choose_node_partition_plan(&mut cancelled)?; + let node_partitions = node_plan.partitions(); + self.checkpoint.evidence.shape_node_partitions = + u64::try_from(node_partitions).map_err(storage)?; let mut identities = FixedRangePartitioner::::new( &self.root, PartitionFamily::Identities, @@ -204,7 +270,7 @@ impl GraphConstructionSession { let mut node_details = FixedRangePartitioner::::new( &self.root, PartitionFamily::NodeDetails, - partitions, + node_partitions, Some(detail_codec), false, )?; @@ -218,7 +284,7 @@ impl GraphConstructionSession { let mut endpoints = FixedRangePartitioner::::new( &self.root, PartitionFamily::Endpoints, - partitions, + node_partitions, None, false, )?; @@ -241,6 +307,7 @@ impl GraphConstructionSession { outputs: Vec::new(), shape_authority_sha256: None, splitters: encode_splitters(&plan), + node_splitters: encode_splitters(&node_plan), partition_identity_rows: Vec::new(), }; install_control(&self.root, SHAPE_INTENT, &shape_intent)?; @@ -286,13 +353,24 @@ impl GraphConstructionSession { catalog_authority.update(receipt.schema_sha256.as_bytes()); catalog_authority.update(receipt.parquet.sha256.as_bytes()); let group = (kind, receipt.schema_sha256.clone()); + // Row groups are keyed by their own kind's UUID (#1439): a + // node-kind schema group's rows are keyed by node UUID and must + // be routed by `node_plan`, matching node details and endpoints + // below -- sized to `node_partitions`, for the same balance-mean + // reason those two are. Edge-kind groups stay on the joint + // `plan`, which their key domain already matches well since + // edges dominate it. + let row_plan = match receipt.kind { + ConstructionChunkKind::Node => &node_plan, + ConstructionChunkKind::Edge => &plan, + }; if !row_groups.contains_key(&group) { row_groups.insert( group.clone(), RowRangePartitioner::new( &self.root, &format!("{kind}-{}", receipt.schema_sha256), - partitions, + row_plan.partitions(), )?, ); } @@ -300,7 +378,7 @@ impl GraphConstructionSession { .get_mut(&group) .ok_or_else(|| storage("row partitioner is absent"))? .push( - &plan, + row_plan, &receipt.parquet.name, self.checkpoint.budgets.max_batch_rows, &mut cancelled, @@ -318,7 +396,7 @@ impl GraphConstructionSession { ConstructionChunkKind::Node => { route_fixed_run::( &self.root, - &plan, + &node_plan, &receipt.details, Some(detail_codec), &mut node_details, @@ -336,9 +414,14 @@ impl GraphConstructionSession { &mut cancelled, &mut self.checkpoint.evidence, )?; + // Endpoints are keyed by the referenced *node*'s UUID at + // this staged, pre-resolution stage (#1439) -- the later + // resolved endpoints, re-keyed by edge UUID in + // `resolve_endpoint_surrogates`, correctly keep using + // `plan`. route_fixed_run::( &self.root, - &plan, + &node_plan, receipt .endpoints .as_ref() @@ -584,6 +667,7 @@ impl GraphConstructionSession { outputs, shape_authority_sha256: Some(shape_authority_sha256), splitters: shape_intent.splitters, + node_splitters: shape_intent.node_splitters, partition_identity_rows, }, )?; @@ -619,6 +703,12 @@ pub(super) fn validate_shape_binding( checkpoint.budgets.partition_count, decode_splitters(&intent.splitters)?, )?; + // The node-only splitters (#1439) are the same kind of authority, over + // the node-keyed families' own domain. + PartitionPlan::from_recorded( + checkpoint.budgets.partition_count, + decode_splitters(&intent.node_splitters)?, + )?; Ok(()) } @@ -941,6 +1031,7 @@ pub(super) fn route_fixed_run( let mut digest = Sha256::new(); let mut bytes = 0_u64; let routed = (|| -> Result<(), GfError> { + let mut run = PartitionRun::new(); while let Some(record) = read_run_record::(&mut reader, codec)? { let wire = run_record_bytes(&record, codec)?; digest.update(wire); @@ -948,9 +1039,11 @@ pub(super) fn route_fixed_run( .checked_add(wire.len() as u64) .ok_or_else(|| storage("partition source byte count overflows"))?; let key: [u8; 16] = record[..16].try_into().expect("fixed key prefix"); - target.route(plan, &key, &record, evidence)?; + let partition = plan.partition_of(&key); + run.push(partition, wire, target, evidence)?; reject_cancelled(cancelled)?; } + run.flush(target, evidence)?; if bytes != source.bytes || hex(&digest.clone().finalize()) != source.sha256 { return Err(storage("construction partition source content changed")); } @@ -960,6 +1053,64 @@ pub(super) fn route_fixed_run( combine_cache_cleanup(routed, released, "construction partition source") } +/// Accumulates one contiguous same-partition run of wire bytes so the +/// caller writes it in a single call instead of once per record (#1439 +/// follow-up). +/// +/// The staged runs routed through [`route_fixed_run`] and +/// [`route_identity_run`] are UUID-sorted and the partition function is +/// monotone, so consecutive same-partition records are always contiguous +/// within one staged chunk. This is the same property +/// [`RowRangePartitioner::route_batch`] already exploits for Arrow rows; +/// this is its fixed-width-record equivalent. +struct PartitionRun { + partition: Option, + bytes: Vec, + records: u64, +} + +impl PartitionRun { + fn new() -> Self { + Self { + partition: None, + bytes: Vec::new(), + records: 0, + } + } + + fn push( + &mut self, + partition: usize, + wire: &[u8], + target: &mut FixedRangePartitioner, + evidence: &mut GraphConstructionEvidence, + ) -> Result<(), GfError> { + if self.partition.is_some_and(|current| current != partition) { + self.flush(target, evidence)?; + } + self.partition = Some(partition); + self.bytes.extend_from_slice(wire); + self.records = self + .records + .checked_add(1) + .ok_or_else(|| storage("partition run record count overflows"))?; + Ok(()) + } + + fn flush( + &mut self, + target: &mut FixedRangePartitioner, + evidence: &mut GraphConstructionEvidence, + ) -> Result<(), GfError> { + if let Some(partition) = self.partition.take() { + target.route_slice(partition, &self.bytes, self.records, evidence)?; + self.bytes.clear(); + self.records = 0; + } + Ok(()) + } +} + /// Stream one staged identity run into its range partitions, widening each /// UUID into the base identity record the shaped domain uses. fn route_identity_run( @@ -1003,6 +1154,7 @@ fn route_identity_run( let mut digest = Sha256::new(); let mut bytes = 0_u64; let routed = (|| -> Result<(), GfError> { + let mut run = PartitionRun::new(); while let Some(uuid) = read_fixed::(&mut reader)? { digest.update(uuid); bytes = bytes @@ -1011,9 +1163,12 @@ fn route_identity_run( let mut record = [0_u8; BASE_IDENTITY_WIDTH]; record[..16].copy_from_slice(&uuid); record[16] = kind; - target.route(plan, &uuid, &record, evidence)?; + let wire = run_record_bytes(&record, None)?; + let partition = plan.partition_of(&uuid); + run.push(partition, wire, target, evidence)?; reject_cancelled(cancelled)?; } + run.flush(target, evidence)?; if bytes != source.bytes || hex(&digest.clone().finalize()) != source.sha256 { return Err(storage( "identity source content changed before partitioning", diff --git a/crates/graphforge-storage/src/graph_construction/shape/tests.rs b/crates/graphforge-storage/src/graph_construction/shape/tests.rs index b0bfc3fa8..3e67164d5 100644 --- a/crates/graphforge-storage/src/graph_construction/shape/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/shape/tests.rs @@ -43,7 +43,7 @@ fn resolved_endpoint_windows_use_logarithmic_name_state_at_1x_2x_4x() { .append( ConstructionChunkKind::Edge, &format!("edges-{offset}"), - &edge_batch(10_000 + u128::from(offset), rows), + &edge_batch(10_000 + u128::from(offset), 1, u128::from(nodes), rows), ) .unwrap(); } @@ -150,7 +150,12 @@ fn shaping_is_bounded_deterministic_and_multipass_at_1x_2x_4x() { .append( ConstructionChunkKind::Edge, "edges", - &edge_batch(10_000, chunks * EDGES_PER_CHUNK), + &edge_batch( + 10_000, + 1, + (chunks * NODES_PER_CHUNK) as u128, + chunks * EDGES_PER_CHUNK, + ), ) .unwrap(); session.seal().unwrap(); @@ -333,7 +338,11 @@ fn shaping_rejects_cross_kind_duplicates_and_missing_endpoints() { .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2)) .unwrap(); duplicate - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(1, 1)) + .append( + ConstructionChunkKind::Edge, + "edges", + &edge_batch(1, 1, 2, 1), + ) .unwrap(); duplicate.seal().unwrap(); assert!( @@ -350,7 +359,11 @@ fn shaping_rejects_cross_kind_duplicates_and_missing_endpoints() { .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2)) .unwrap(); missing - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(20_000, 2)) + .append( + ConstructionChunkKind::Edge, + "edges", + &edge_batch(20_000, 1, 3, 2), + ) .unwrap(); missing.seal().unwrap(); assert!( diff --git a/crates/graphforge-storage/src/graph_construction/tests.rs b/crates/graphforge-storage/src/graph_construction/tests.rs index f943da1dc..f81c56c6b 100644 --- a/crates/graphforge-storage/src/graph_construction/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/tests.rs @@ -57,14 +57,34 @@ fn distinct_label_batch(first: u128, rows: usize) -> RecordBatch { .unwrap() } -pub(super) fn edge_batch(first: u128, rows: usize) -> RecordBatch { +/// `first` is this batch's edge-UUID base. `node_start` and `node_count` +/// (#1439) describe where this batch's src/dst references begin and how +/// many nodes exist to reference, 1-indexed: every reference wraps modulo +/// `node_count`, so it is always a valid node id regardless of chunk +/// position. A caller with more than one chunk must vary `node_start` per +/// chunk, or every chunk in the loop references the same low node-UUID band +/// no matter how many nodes or chunks actually exist. That was invisible +/// while endpoints were routed with the joint identity splitters, which +/// never resolve node UUIDs finely enough to notice; routing them with +/// node-only splitters surfaces it as a real (fixture-caused, not +/// production) skew. Single-chunk callers keep passing `node_start: 1`, +/// matching this function's previous fixed behaviour exactly whenever +/// `rows <= node_count`. +pub(super) fn edge_batch( + first: u128, + node_start: u128, + node_count: u128, + rows: usize, +) -> RecordBatch { let edges = (first..first + rows as u128) .map(u128::to_be_bytes) .collect::>(); - let src = (1..=rows as u128) + let src = (0..rows as u128) + .map(|offset| 1 + (node_start - 1 + offset) % node_count) .map(u128::to_be_bytes) .collect::>(); - let dst = (2..=rows as u128 + 1) + let dst = (0..rows as u128) + .map(|offset| 1 + (node_start + offset) % node_count) .map(u128::to_be_bytes) .collect::>(); RecordBatch::try_new( @@ -222,7 +242,11 @@ pub(super) fn nonempty_project_with_nodes(node_count: u64) -> TempDir { ) .unwrap(); session - .append(ConstructionChunkKind::Edge, "edge", &edge_batch(100, 1)) + .append( + ConstructionChunkKind::Edge, + "edge", + &edge_batch(100, 1, u128::from(node_count), 1), + ) .unwrap(); session.seal().unwrap(); let encoded = session.prepare_canonical_encoding(1).unwrap(); @@ -247,7 +271,11 @@ pub(super) fn nonempty_project_generation_two() -> TempDir { .append(ConstructionChunkKind::Node, "node", &node_batch(3, 1)) .unwrap(); session - .append(ConstructionChunkKind::Edge, "edge", &edge_batch(101, 1)) + .append( + ConstructionChunkKind::Edge, + "edge", + &edge_batch(101, 1, 3, 1), + ) .unwrap(); session.seal().unwrap(); let encoded = session.prepare_canonical_encoding(2).unwrap(); @@ -403,24 +431,18 @@ fn million_chunk_shaping_retains_name_state_bounded_by_the_partition_count() { // The online merge scheduler's logarithmic name state is gone: range // partitioning retains exactly one spill name per partition per family, // independent of how many chunks were staged. - // - // `budgets.partition_count` is a *ceiling* (#1439), defaulted to - // `MAX_PARTITION_COUNT` so the data-driven cut -- not an artificially low - // flat cap -- is what bounds production partition counts. Worst-case - // retained name state is therefore sized off the ceiling, not off the - // pre-#1439 flat default, and still stays a small fraction of max_chunks. let budgets = GraphConstructionBudgets::default(); assert_eq!(budgets.max_chunks, 1_000_000); let families = super::partition_shaping::PartitionFamily::ALL.len() as u64; let slots = u64::from(budgets.partition_count) * families; - // Name state is a function of the recorded partition count ceiling and - // the family set, never of the staged chunk count. - assert_eq!(slots, 20_480, "retained slots: {slots}"); - assert!(slots < budgets.max_chunks / 40, "retained slots: {slots}"); + // Name state is a function of the recorded partition count and the family + // set, never of the staged chunk count. + assert_eq!(slots, 1_280, "retained slots: {slots}"); + assert!(slots < budgets.max_chunks / 100, "retained slots: {slots}"); assert_eq!(budgets.max_schema_groups, 256); assert_eq!( budgets.partition_count, - super::partition::MAX_PARTITION_COUNT + super::partition::DEFAULT_PARTITION_COUNT ); } @@ -597,7 +619,11 @@ fn node_after_edge_and_concurrent_same_process_open_fail_closed() { .is_err() ); session - .append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 2)) + .append( + ConstructionChunkKind::Edge, + "edges", + &edge_batch(100, 1, 2, 2), + ) .unwrap(); assert!( session From 03013d03a5110f710c0c3ca657a0bfe7e8c59e2c Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Thu, 17 Sep 2026 17:48:43 +0000 Subject: [PATCH 3/5] docs(storage): flag the stale digest rows in the determinism header comment (#1439) The header comment's historical measurement (on 13632d4b, predating S5's external-merge-tree deletion) documents self-consistency across two sessions, not a pinned expectation -- no test asserts those literal values. Note that edge details and edge endpoints legitimately changed under the per-family splitter routing landed alongside this, while identities and node details did not, matching which families' routing actually changed. Co-Authored-By: Claude Sonnet 5 --- .../src/construction_determinism_tests.rs | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/crates/graphforge-storage/src/construction_determinism_tests.rs b/crates/graphforge-storage/src/construction_determinism_tests.rs index 86743c2ba..7b8ba4872 100644 --- a/crates/graphforge-storage/src/construction_determinism_tests.rs +++ b/crates/graphforge-storage/src/construction_determinism_tests.rs @@ -27,6 +27,17 @@ // That is a pre-existing durable-format defect, it is not introduced here, and // fixing it is a separate concern. // +// The table above is a historical measurement on `13632d4b`, predating S5's +// external-merge-tree deletion, not a pinned expectation: no test asserts +// those literal values. Read it for what it demonstrates -- self-consistency +// across two sessions -- not as a checksum to reproduce. It is stale for two +// rows as of #1439's follow-up (per-family splitters): routing node-keyed +// families with node-only splitters changes how edge details and edge +// endpoints partition, so `edge details` and `edge endpoints` legitimately +// produce different bytes now. `shaped-identities.run` and `node details` +// are unchanged, because their routing did not change. The claim this file's +// tests actually enforce -- stated precisely below -- still holds. +// // So the claim these tests make is precise: // // * **Within a fixed set of recorded session parameters** — the session clock From f3ec047457da704af550765fa81d933b366ddfb1 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Thu, 17 Sep 2026 18:35:34 +0000 Subject: [PATCH 4/5] fix(storage): balance fixed-width partitions by distinct keys, not rows The balance check added in dddf1e65 refused every Graph500-shaped construction below ladder scale: the endpoints family is keyed by node and holds one record per incident edge, so with power-law degrees the partition that owns a few hubs is row-heavy by design, and discounting exactly one key cannot cover it (scale 10: largest partition 315 of 1,858 rows across 64, largest single-key run 91, mean 29). That is what CI's scale_g500_ladder target failed on. The splitters are quantiles of the key domain, so what they promise to balance is distinct keys per partition. finish_optional already walks each partition in key order; count key changes there and put that balance through the unchanged assert_balanced. For unique-key families it is the same number as before. A collapsed or mis-sampled splitter set still concentrates distinct keys and is still refused; a hub-heavy family over balanced keys is accepted. Both are pinned by fixtures. Co-Authored-By: Claude Fable 5.1 --- .../src/graph_construction/partition.rs | 62 ++-------- .../graph_construction/partition_shaping.rs | 34 +++-- .../src/graph_construction/tests.rs | 117 ++++++++++++++++++ 3 files changed, 143 insertions(+), 70 deletions(-) diff --git a/crates/graphforge-storage/src/graph_construction/partition.rs b/crates/graphforge-storage/src/graph_construction/partition.rs index fac829ce7..acd2491ba 100644 --- a/crates/graphforge-storage/src/graph_construction/partition.rs +++ b/crates/graphforge-storage/src/graph_construction/partition.rs @@ -359,6 +359,16 @@ impl PartitionBalance { /// where the quantile estimate has fewer points than partitions and the /// ratio carries no information. /// + /// "Rows" are whatever the caller recorded. The splitters are quantiles + /// of the *key* domain, so the quantity they promise to balance is the + /// number of distinct keys per partition, not records: a range partition + /// cannot split one key, and a family with several records per key + /// (staged endpoints, keyed by node, hold one record per incident edge) + /// is expected to be row-skewed exactly as much as its key degrees are. + /// Such a family records its distinct-key counts here instead of its row + /// counts (`FixedRangePartitioner::finish_optional`, #1439); for a + /// unique-key family the two are the same number. + /// /// # Errors /// Returns an error when the partitioning is skewed beyond the tolerance. pub(crate) fn assert_balanced(&self, context: &str) -> Result<(), GfError> { @@ -386,58 +396,6 @@ impl PartitionBalance { } Ok(()) } - - /// Refuse a partitioning whose largest partition exceeds - /// [`BALANCE_TOLERANCE`] times the mean, **unless** the excess is - /// explained by one repeated key. - /// - /// A range partition cannot split one key across partitions, so a hub -- - /// one node referenced by a disproportionate number of edges, routing a - /// disproportionate number of endpoint records to the one partition that - /// owns it -- is not a splitter defect. `max_single_key_run` is the - /// largest number of records sharing one key observed anywhere in the - /// partitioning (computed after sorting, where identical keys are - /// contiguous). Discounting all but one of those occurrences from the - /// largest partition before comparing against the tolerance is what - /// distinguishes a genuine hub from a skewed splitter set: a bad splitter - /// set concentrates *distinct* keys into one partition, which this - /// discount does not hide. - /// - /// # Errors - /// Returns an error when the partitioning is skewed beyond the tolerance - /// even after that discount. - pub(crate) fn assert_balanced_with_hub_tolerance( - &self, - context: &str, - max_single_key_run: u64, - ) -> Result<(), GfError> { - let partitions = self.rows.len() as u64; - if partitions == 0 { - return Err(storage("partition balance has no partitions")); - } - let total = self.total(); - if total / partitions < BALANCE_MIN_MEAN_ROWS { - return Ok(()); - } - let max = self.max_rows(); - let discount = max_single_key_run.saturating_sub(1); - let discounted = max.saturating_sub(discount); - let scaled = discounted - .checked_mul(partitions) - .ok_or_else(|| storage("partition balance ratio overflows"))?; - let budget = total - .checked_mul(BALANCE_TOLERANCE) - .ok_or_else(|| storage("partition balance budget overflows"))?; - if scaled > budget { - return Err(storage(format!( - "{context} range partitioning is skewed: largest partition holds {max} of \ - {total} rows across {partitions} partitions (largest single-key run \ - {max_single_key_run}), above the {BALANCE_TOLERANCE}x mean tolerance even after \ - discounting one hub key" - ))); - } - Ok(()) - } } #[cfg(test)] diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs index f8f451192..104726e02 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs @@ -479,15 +479,17 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { let concatenated = (|| -> Result { let mut previous: Option<[u8; N]> = None; let mut written = 0_u64; - // Tracks the longest run of records sharing one key, across the - // whole concatenation. A range partition cannot split one key, so - // this is what tells a hub apart from a skewed splitter set - // (#1439) -- see `PartitionBalance::assert_balanced_with_hub_tolerance`. - // Keys are unique within a partition's sorted order without - // spanning partitions (the partition function is a function of - // the key), so accumulating this across the loop below is exact. - let mut current_run = 0_u64; - let mut max_single_key_run = 0_u64; + // Distinct keys per partition, counted off the sorted order below + // (identical keys are contiguous and never span partitions, since + // the partition function is a function of the key). This, not the + // row count, is what the splitters promise to balance: a range + // partition cannot split one key, so a node-keyed family such as + // staged endpoints is row-skewed exactly as much as its degree + // distribution is, and a power-law graph has many hubs, not one + // (#1439). A bad splitter set concentrates *distinct* keys, which + // this count does not hide. For a unique-key family it equals the + // row balance recorded at routing time. + let mut keys = PartitionBalance::new(self.sealed.len()); for partition in 0..self.sealed.len() { let Some(name) = self.sealed[partition] .as_ref() @@ -513,15 +515,12 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { "duplicate identity across construction runs", )); } - current_run = if prior[..16] == record[..16] { - current_run.saturating_add(1) - } else { - 1 - }; + if prior[..16] != record[..16] { + keys.record(partition)?; + } } else { - current_run = 1; + keys.record(partition)?; } - max_single_key_run = max_single_key_run.max(current_run); let wire = run_record_bytes(record, self.codec)?; writer.write_all(wire).map_err(super::storage)?; account_merge_write_bytes(evidence, wire.len() as u64)?; @@ -539,8 +538,7 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { // test, so this is what tells a working range partition apart // from a catastrophically skewed one for every fixed-width // family, not only identities. - self.balance - .assert_balanced_with_hub_tolerance(self.family.as_str(), max_single_key_run)?; + keys.assert_balanced(&format!("{} distinct keys", self.family.as_str()))?; Ok(written) })(); let written = match concatenated { diff --git a/crates/graphforge-storage/src/graph_construction/tests.rs b/crates/graphforge-storage/src/graph_construction/tests.rs index f81c56c6b..e1b5ba3de 100644 --- a/crates/graphforge-storage/src/graph_construction/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/tests.rs @@ -1123,3 +1123,120 @@ pub(super) fn construction_session_root(root: &TempDir, operation: Uuid) -> Path .join(PRIVATE_ROOT) .join(operation.simple().to_string()) } + +/// Endpoint-shaped fixture for the partition balance check: 64 partitions of +/// exactly 32 node keys each, so the key domain is perfectly balanced, with +/// the node keys in `hubs` referencing `hub_degree` edges apiece and every +/// other node exactly one. A route plan that agrees with that key layout is +/// built by the caller. +fn route_endpoint_fixture( + root: &StableDirectory, + plan: &super::partition::PartitionPlan, + hubs: &[u16], + hub_degree: u32, + evidence: &mut GraphConstructionEvidence, +) -> Result, GfError> { + use super::partition_shaping::{FixedRangePartitioner, PartitionFamily}; + use crate::construction_record_layout::ENDPOINT_WIDTH; + const PARTITIONS: usize = 64; + const KEYS_PER_PARTITION: u16 = 32; + let mut partitioner = FixedRangePartitioner::::new( + root, + PartitionFamily::Endpoints, + PARTITIONS, + None, + false, + )?; + let mut edge = 0_u128; + for node in 0..(PARTITIONS as u16 * KEYS_PER_PARTITION) { + let degree = if hubs.contains(&node) { hub_degree } else { 1 }; + let mut key = [0_u8; 16]; + key[..2].copy_from_slice(&node.to_be_bytes()); + for _ in 0..degree { + edge += 1; + let mut record = [0_u8; ENDPOINT_WIDTH]; + record[..16].copy_from_slice(&key); + record[16..32].copy_from_slice(&edge.to_be_bytes()); + partitioner.route(plan, &key, &record, evidence)?; + } + } + partitioner.seal(evidence)?; + partitioner.finish_optional("shaped-fixture-endpoints.run", &mut || false, evidence) +} + +/// Splitters at every 32nd node key: the plan the fixture's keys are laid out +/// for, so each partition owns exactly 32 distinct keys. +fn endpoint_fixture_plan() -> super::partition::PartitionPlan { + let splitters = (1_u16..64) + .map(|partition| { + let mut splitter = [0_u8; 16]; + splitter[..2].copy_from_slice(&(partition * 32).to_be_bytes()); + splitter + }) + .collect(); + super::partition::PartitionPlan::from_recorded(64, splitters).unwrap() +} + +#[test] +fn hub_heavy_endpoints_over_balanced_keys_are_accepted() { + // Three hubs in partition 0 and two elsewhere, each with 400 incident + // edges against a mean of ~63 rows per partition: partition 0 holds 1,229 + // of 4,043 rows, more than 4x the mean even after discounting any one + // hub. That is the Graph500 shape (`scale_g500_ladder` at scale 10: + // "largest partition holds 315 of 1858 rows across 64 partitions, + // largest single-key run 91"), and it is not a splitter defect: every + // partition owns exactly 32 distinct keys. + let root = TempDir::new().unwrap(); + let mut session = open(&root, 8_041); + session + .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2)) + .unwrap(); + session.seal().unwrap(); + let GraphConstructionSession { + root: session_root, + checkpoint, + .. + } = &mut session; + let output = route_endpoint_fixture( + session_root, + &endpoint_fixture_plan(), + &[5, 6, 7, 700, 1_500], + 400, + &mut checkpoint.evidence, + ) + .unwrap(); + assert!(output.is_some()); +} + +#[test] +fn collapsed_endpoint_splitters_are_still_refused() { + // The same 2,048 keys, one record each, under a splitter set that lies + // entirely above the key domain: every key lands in partition 0. No hub + // is involved, so the refusal can only come from key concentration -- + // the property the balance check exists to guarantee. + let root = TempDir::new().unwrap(); + let mut session = open(&root, 8_042); + session + .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2)) + .unwrap(); + session.seal().unwrap(); + let GraphConstructionSession { + root: session_root, + checkpoint, + .. + } = &mut session; + let collapsed = (1_u8..64) + .map(|partition| { + let mut splitter = [0xff_u8; 16]; + splitter[1] = partition; + splitter + }) + .collect(); + let plan = super::partition::PartitionPlan::from_recorded(64, collapsed).unwrap(); + let error = route_endpoint_fixture(session_root, &plan, &[], 1, &mut checkpoint.evidence) + .expect_err("a collapsed key partitioning must be refused") + .to_string(); + assert!(error.contains("endpoints distinct keys"), "{error}"); + assert!(error.contains("skewed"), "{error}"); + assert!(error.contains("2048 of 2048"), "{error}"); +} From b0e53d36896f69aae1055a153cd33497408a01a2 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Thu, 17 Sep 2026 18:21:55 +0000 Subject: [PATCH 5/5] perf(storage): restore a bounded spill write buffer for the Resolved family #1440 removed the 64 KiB per-partition spill buffer from every fixed-width family and batched writes by contiguous same-partition run instead. Four of the five families arrive sorted by their routing key, so their runs are kilobytes each and the change was free. The Resolved family is re-keyed by edge UUID away from its node-UUID input order, so it has no runs at all: every resolved endpoint reached the descriptor as its own write(2), 33.6M calls for 840 MB at S20 (evidence: shape_consume_reauthentication.write_calls 93,118 -> 33,663,169 between 9269362f and 03013d03), and the shaping seal went from 113 s to 162 s. That is the ~23% ingest throughput regression the memory fix cost. Give the spill writer a per-family bound: zero for the four run-batched families, 8 KiB for Resolved, allocated lazily per open spill. Only the Resolved family holds spills open while it is live, so the residency is 256 partitions x 8 KiB = 2 MiB, against the 4 x 256 x 64 KiB #1440 removed. Co-Authored-By: Claude Fable 5.1 (cherry picked from commit a39c4b832be3bff6fab0fddbde7ca044074ce625) --- .../graph_construction/partition_shaping.rs | 107 ++++++++++++++---- .../src/graph_construction/shape/tests.rs | 42 +++++++ 2 files changed, 129 insertions(+), 20 deletions(-) diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs index 104726e02..95e2f220a 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs @@ -39,12 +39,34 @@ use sha2::Digest; use std::ffi::OsStr; use std::io::{BufWriter, Write}; -/// Per-spill write buffer. Every partition holds one open spill during the -/// routing pass, so this is multiplied by the partition count: it is -/// deliberately far below `BLOCK_BYTES`, which is sized for the handful of -/// streams a merge group used to open. +/// Per-spill write buffer for row (Arrow) partition spills. Every partition +/// holds one open spill during the routing pass, so this is multiplied by the +/// partition count: it is deliberately far below `BLOCK_BYTES`, which is sized +/// for the handful of streams a merge group used to open. const SPILL_BLOCK_BYTES: usize = 64 * 1024; +/// Per-spill write buffer for the one fixed-width family that is routed a +/// record at a time, [`PartitionFamily::Resolved`] (#1443). +/// +/// The other four fixed-width families arrive sorted by the key they are +/// routed by, so [`FixedRangePartitioner::route_slice`] receives whole +/// contiguous same-partition runs (kilobytes each at the ladder's chunk +/// shape) and needs no buffer behind it. Resolved endpoints are re-keyed by +/// edge UUID away from their node-UUID input order, so consecutive records +/// land in unrelated partitions and every one of them reached the +/// descriptor as its own `write(2)`: 33.6M calls for 840 MB at S20, against +/// 93k for the whole shaping pass before #1440 removed the buffers, and the +/// measured 23% ingest throughput regression that removal cost. This buffer +/// restores the batching for exactly that family. +/// +/// Sized so the bound is invisible in resident memory even at the partition +/// ceiling: it is allocated lazily per open spill, only the Resolved family +/// holds spills open while it is live (the staged families are sealed and +/// finished before resolution starts), and 256 partitions cost 2 MiB at this +/// size. 8 KiB already cuts the S20 call count to ~100k, so raising it buys +/// nothing measurable (see the #1443 curve). +const RESOLVED_SPILL_BUFFER_BYTES: usize = 8 * 1024; + /// Canonical fixed-width partition families. The family name is part of the /// durable artifact grammar, so it is a closed set rather than a free string. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -67,6 +89,16 @@ impl PartitionFamily { } } + /// Bound on the per-partition write buffer a spill of this family holds + /// while open; zero writes every routed slice straight to the descriptor. + /// See [`RESOLVED_SPILL_BUFFER_BYTES`] for why only one family has one. + pub(super) const fn spill_buffer_bytes(self) -> usize { + match self { + Self::Resolved => RESOLVED_SPILL_BUFFER_BYTES, + Self::Identities | Self::NodeDetails | Self::EdgeDetails | Self::Endpoints => 0, + } + } + /// Every fixed-width family name, for the durable-name grammar. pub(super) const ALL: [Self; 5] = [ Self::Identities, @@ -190,30 +222,40 @@ impl SpillWriter { } } -/// One open, unsealed fixed-width partition spill (#1439 follow-up). +/// One open, unsealed fixed-width partition spill (#1439 follow-up, #1443). /// -/// Unlike [`SpillWriter`], this holds no per-partition write buffer. Every -/// partition open at once during routing used to cost `SPILL_BLOCK_BYTES` -/// (64 KiB) of resident buffer regardless of how much of it had actually -/// been written; concurrent spill buffers across partitions and families -/// were the dominant measured term in both the RSS and fsync growth #1439 -/// originally investigated. The staged runs routed through this are -/// UUID-sorted and the partition function is monotone, so the caller can -/// accumulate one contiguous same-partition run of wire bytes and write it -/// in a single call -- see [`FixedRangePartitioner::route_slice`] -- rather -/// than buffering to amortize many small per-record writes. -struct UnbufferedSpillWriter { +/// Unlike [`SpillWriter`], this holds no fixed-size per-partition block. +/// Every partition open at once during routing used to cost +/// `SPILL_BLOCK_BYTES` (64 KiB) of resident buffer regardless of how much of +/// it had actually been written; concurrent spill buffers across partitions +/// and families were the dominant measured term in both the RSS and fsync +/// growth #1439 originally investigated. The staged runs routed through this +/// are UUID-sorted and the partition function is monotone, so the caller +/// accumulates one contiguous same-partition run of wire bytes and writes it +/// in a single call -- see [`FixedRangePartitioner::route_slice`] -- and +/// those families set `bound` to zero. +/// +/// The one family whose input order does not match its routing key, +/// [`PartitionFamily::Resolved`], sets a small `bound` instead: slices +/// shorter than it are coalesced in `buffer` (grown lazily, never past the +/// bound) and reach the descriptor in bound-sized writes; a slice at least +/// as long as the bound bypasses the buffer. This is the batching the 64 KiB +/// block used to provide, at a fraction of its residency. +struct FixedSpillWriter { name: String, temporary: std::ffi::OsString, identity: FileIdentity, writer: HashingWriter, + buffer: Vec, + bound: usize, } -impl UnbufferedSpillWriter { +impl FixedSpillWriter { fn create( root: &StableDirectory, name: String, window: std::num::NonZeroU64, + bound: usize, ) -> Result { let temporary = artifact_temp(&name); let file = root @@ -226,11 +268,34 @@ impl UnbufferedSpillWriter { temporary, identity, writer, + buffer: Vec::new(), + bound, }) } fn write(&mut self, bytes: &[u8]) -> Result<(), GfError> { - self.writer.write_all(bytes).map_err(super::storage) + if self.bound == 0 { + return self.writer.write_all(bytes).map_err(super::storage); + } + if self.buffer.len() + bytes.len() > self.bound { + self.flush_buffer()?; + } + if bytes.len() >= self.bound { + self.writer.write_all(bytes).map_err(super::storage) + } else { + self.buffer.extend_from_slice(bytes); + Ok(()) + } + } + + fn flush_buffer(&mut self) -> Result<(), GfError> { + if !self.buffer.is_empty() { + self.writer + .write_all(&self.buffer) + .map_err(super::storage)?; + self.buffer.clear(); + } + Ok(()) } /// Drop an unsealed spill and remove its temporary. See @@ -252,6 +317,7 @@ impl UnbufferedSpillWriter { root: &StableDirectory, evidence: &mut GraphConstructionEvidence, ) -> Result { + self.flush_buffer()?; self.writer.flush().map_err(super::storage)?; self.writer .inner @@ -303,7 +369,7 @@ pub(super) struct FixedRangePartitioner<'a, const N: usize> { codec: Option, reject_duplicates: bool, window: std::num::NonZeroU64, - spills: Vec>, + spills: Vec>, sealed: Vec>, balance: PartitionBalance, records: u64, @@ -388,10 +454,11 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { .get_mut(partition) .ok_or_else(|| super::storage("routed partition is out of range"))?; if slot.is_none() { - *slot = Some(UnbufferedSpillWriter::create( + *slot = Some(FixedSpillWriter::create( self.root, fixed_spill_name(self.family, partition), self.window, + self.family.spill_buffer_bytes(), )?); } slot.as_mut() diff --git a/crates/graphforge-storage/src/graph_construction/shape/tests.rs b/crates/graphforge-storage/src/graph_construction/shape/tests.rs index 3e67164d5..7572c884e 100644 --- a/crates/graphforge-storage/src/graph_construction/shape/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/shape/tests.rs @@ -657,3 +657,45 @@ fn fixed_merge_reader_rejects_truncation_and_self_loop_is_valid() { shape ); } + +#[test] +fn resolved_endpoints_never_cost_one_write_submission_each() { + // Resolved endpoints are re-keyed by edge UUID away from their node-UUID + // input order, so they are routed one record at a time. Without a bounded + // spill buffer behind that route every record reached the descriptor as + // its own write (#1440: 33.6M submissions for 840 MB at S20, the whole + // shaping pass having needed 93k before), and the ingest path lost ~23% + // of its throughput. Every other family is run-batched, so the shaping + // pass as a whole must submit far fewer writes than there are resolved + // endpoints. + const NODES: usize = 256; + const EDGES: usize = 8_192; + let root = TempDir::new().unwrap(); + let mut session = GraphConstructionSession::open( + root.path(), + Uuid::from_u128(7_443), + 0, + GraphConstructionBudgets::default(), + ) + .unwrap(); + session + .append(ConstructionChunkKind::Node, "nodes", &node_batch(1, NODES)) + .unwrap(); + session + .append( + ConstructionChunkKind::Edge, + "edges", + &edge_batch(10_000, 1, NODES as u128, EDGES), + ) + .unwrap(); + session.seal().unwrap(); + let shape = session.shape_canonical_with_cancellation(|| false).unwrap(); + assert_eq!(shape.edge_count, EDGES as u64); + assert!(shape.edge_endpoints.is_some()); + let resolved_endpoints = 2 * EDGES as u64; + let submissions = session.evidence().merge_write_operations; + assert!( + submissions < resolved_endpoints, + "shaping submitted {submissions} writes for {resolved_endpoints} resolved endpoints" + ); +}