Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 21 additions & 4 deletions crates/graphforge-storage/src/graph_construction/shape.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1098,9 +1098,11 @@ pub(super) fn route_fixed_run<const N: usize>(
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).
/// Maximum wire bytes retained by one routing accumulator.
const PARTITION_RUN_BYTES: usize = 64 * 1024;

/// Accumulates bounded same-partition slices of wire bytes, amortizing writes
/// without retaining a whole partition run (#1445).
///
/// The staged runs routed through [`route_fixed_run`] and
/// [`route_identity_run`] are UUID-sorted and the partition function is
Expand All @@ -1109,14 +1111,20 @@ pub(super) fn route_fixed_run<const N: usize>(
/// [`RowRangePartitioner::route_batch`] already exploits for Arrow rows;
/// this is its fixed-width-record equivalent.
struct PartitionRun {
bound: usize,
partition: Option<usize>,
bytes: Vec<u8>,
records: u64,
}

impl PartitionRun {
fn new() -> Self {
Self::with_bound(PARTITION_RUN_BYTES)
}

fn with_bound(bound: usize) -> Self {
Self {
bound,
partition: None,
bytes: Vec::new(),
records: 0,
Expand All @@ -1130,9 +1138,18 @@ impl PartitionRun {
target: &mut FixedRangePartitioner<N>,
evidence: &mut GraphConstructionEvidence,
) -> Result<(), GfError> {
if self.partition.is_some_and(|current| current != partition) {
if self.partition.is_some_and(|current| current != partition)
|| wire.len() > self.bound - self.bytes.len()
{
self.flush(target, evidence)?;
}
// Keep an oversized record whole without growing the accumulator.
if wire.len() >= self.bound {
return target.route_slice(partition, wire, 1, evidence);
}
if self.bytes.capacity() == 0 {
self.bytes.reserve_exact(self.bound);
}
self.partition = Some(partition);
self.bytes.extend_from_slice(wire);
self.records = self
Expand Down
79 changes: 79 additions & 0 deletions crates/graphforge-storage/src/graph_construction/shape/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -710,3 +710,82 @@ fn resolved_endpoints_never_cost_one_write_submission_each() {
"shaping submitted {submissions} writes for {resolved_endpoints} resolved endpoints"
);
}

/// A single partition larger than every candidate buffer must stream without
/// growing the accumulator or changing its authenticated bytes/row counts.
#[test]
fn partition_run_bounds_fixed_records_without_changing_output() {
let records = (1_u128..=131_073)
.map(u128::to_be_bytes)
.collect::<Vec<_>>();
for bound in [64 * 1024, 256 * 1024, 1024 * 1024] {
assert_bounded_partition_run(&records, None, bound);
}
}

#[test]
fn partition_run_keeps_compact_records_whole_across_flushes() {
let records = (1_u128..=512)
.map(|key| {
let mut record = [0; NODE_DETAIL_WIDTH];
record[..16].copy_from_slice(&key.to_be_bytes());
let length = if key % 2 == 0 { 255 } else { 1 };
record[16] = length;
record[17..17 + usize::from(length)].fill(b'x');
record
})
.collect::<Vec<_>>();
// The 64-byte case also exercises a record larger than the buffer.
for bound in [64, 1024, 64 * 1024] {
assert_bounded_partition_run(&records, Some(DetailCodec::Compact), bound);
}
}

fn assert_bounded_partition_run<const N: usize>(
records: &[[u8; N]],
codec: Option<DetailCodec>,
bound: usize,
) {
let root = TempDir::new().unwrap();
crate::open_or_initialize_project(root.path()).unwrap();
let mut session = open(&root, 0x1445);
let mut target =
FixedRangePartitioner::<N>::new(&session.root, PartitionFamily::Identities, 2, codec, true)
.unwrap();
let evidence = &mut session.checkpoint.evidence;
let mut run = PartitionRun::with_bound(bound);
let mut expected = Vec::new();
// Two equally populated partitions exercise partition changes as well as
// byte-triggered flushes. Partition changes must retain partial buffers.
for (index, record) in records.iter().enumerate() {
let wire = codec
.map_or(Ok(record.as_slice()), |codec| codec.bytes(record))
.unwrap();
expected.extend_from_slice(wire);
run.push(
index / records.len().div_ceil(2),
wire,
&mut target,
evidence,
)
.unwrap();
assert!(run.bytes.len() <= bound);
assert!(run.bytes.capacity() <= bound);
}
run.flush(&mut target, evidence).unwrap();
run.flush(&mut target, evidence).unwrap(); // Empty flush cannot duplicate rows.
assert_eq!(target.balance().total(), records.len() as u64);
let output = target
.finish_optional("staged-identities.run", &mut || false, evidence)
.unwrap()
.unwrap();
let mut actual = Vec::new();
session
.root
.open_child_file(OsStr::new(&output))
.unwrap()
.read_to_end(&mut actual)
.unwrap();
assert_eq!(actual, expected, "bound={bound}");
assert_eq!(evidence.partition_rows, records.len() as u64);
}
182 changes: 182 additions & 0 deletions docs/development/evidence/partition-run-bound-1445.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,182 @@
{
"base_sha": "51a0d6b984b3fab960a9c1116b7e1b4940e681d8",
"sha256": {
"/home/ubuntu/gf-1445-evidence/bin/bound-256k": "d5df2ddda8731795f626b6b7d09467c3ebd0e125b98e4bc375eb8a33d140fd1b",
"/home/ubuntu/gf-1445-evidence/bin/baseline": "dc10f5717d425611fc01742192384eaa8321b2551669b02fca9c86371e489bde",
"/home/ubuntu/gf-1445-evidence/bin/bound-64k": "00f681220c83e92a694957ff058baeb3128bc89f8f9531bf77349aa639a400f1",
"/home/ubuntu/gf-1445-evidence/bin/bound-1m": "2b2094da75352ea97ac0130cfadb9bf77848e93fb6b61178f60fa3c0f921977e",
"/home/ubuntu/gf-1445-evidence/candidate.patch": "7c1977d0ae00a8236103b4e27c91279e26acf20e2d4c2ddcdbec329d64c85f1e",
"/home/ubuntu/gf-1445-evidence/histogram.patch": "bd9ab301acc6672e36bc144bae7b68bd968fb745c012047604afb4ab4a00ec4c",
"/home/ubuntu/gf-1445-evidence/ingest.py": "485e1b968ecfda7fc9693e88712dcfdbf2a318b2ec47f5d0aab56d4bccde0a0a",
"/home/ubuntu/gf-1445-evidence/run-measured.sh": "6ca8c521b5956682c7f9abf8ab1113c81a82232b4de611ce63458d620e0aef56",
"/home/ubuntu/graphforge-ladder/cargo-bench/release/graphforge-benchmark-graph500-generator": "b1ddbd8bbcac15e50d530d7599f706603ab5a2e357a64730fb5ffa969d00e968",
"/home/ubuntu/gf-1445-evidence/quiet.py": "a6e86d2461a6bac9626d225d30c6ac17c8b2f73def00539117d9494aa04a54c5",
"/home/ubuntu/gf-1445-evidence/resume-curve.py": "9a8037abc4245652e52c3f13157373ac062c74b9abec4ea9a6e7fd62a5237bf8",
"/home/ubuntu/gf-1445-evidence/summarize-curve.py": "78baf16a2564a12db25e26fd77b8c72b6292b67ea59c3a54b8d3b48ec511278f",
"/home/ubuntu/gf-1445-evidence/measure-floors.py": "cac4be0cba0e780a4ad7262c2fa2269a9c2c0fabbb08ebcdca61f6c56c8a2a9d",
"/home/ubuntu/gf-1445-evidence/measurement-plan.txt": "69c3f83b6b471eaa2b9c712c3774fc94189a877ae1909ba9fedca7a8a7351ff3"
},
"seed": 13907095936298285200,
"edge_factor": 16,
"input_sha256": {
"/home/ubuntu/gf-1445-evidence/input/edges.parquet": "30f197e5813adf8e0352d71d4451061686dfb96d267e681524c30344a96a52f4",
"/home/ubuntu/gf-1445-evidence/input/nodes.parquet": "44c9dfd9325013d0f6ea2f03bd86b00d2bea01265254cb28f70f0881a6478075",
"/home/ubuntu/gf-1445-evidence/input-s19/edges.parquet": "bd71780610259b1e6c721097926d930e20c434963ed1b3aa658802bdd2be91a0",
"/home/ubuntu/gf-1445-evidence/input-s19/nodes.parquet": "86491d7aa07331f25bbc788dc8aade06f97e7979f013e496e1970eee23cc1d98",
"/home/ubuntu/gf-1445-evidence/input-s20/edges.parquet": "c8e6b66f3d56b31f284a1288621d5536b7e5e347b6050e0ed3b0c8fed60e8953",
"/home/ubuntu/gf-1445-evidence/input-s20/nodes.parquet": "5792da943d39a3ec0cfe48c375fef1b078ae31f2c086af5d9bcfd417c74f24aa"
},
"curve": {
"base_sha": "51a0d6b984b3fab960a9c1116b7e1b4940e681d8",
"selection_limit_wall_s": 30.78607031319989,
"selected": "bound-64k",
"curve": {
"baseline": {
"samples": [
{
"label": "baseline-s18-guarded-round1",
"wall_s": 29.823490793118253,
"cpu_s": 25.896252,
"cgroup_peak_bytes": 792477696
},
{
"label": "baseline-s18-guarded-round2-attempt1",
"wall_s": 29.143522070953622,
"cpu_s": 25.505372,
"cgroup_peak_bytes": 821202944
},
{
"label": "baseline-s18-guarded-round3-attempt1",
"wall_s": 29.320066964952275,
"cpu_s": 25.723618,
"cgroup_peak_bytes": 763043840
}
],
"wall_s": 29.320066964952275,
"cpu_s": 25.723618,
"cgroup_peak_bytes": 792477696
},
"bound-64k": {
"samples": [
{
"label": "bound-64k-s18-guarded-round1-attempt2",
"wall_s": 29.374653718201444,
"cpu_s": 25.724124,
"cgroup_peak_bytes": 799113216
},
{
"label": "bound-64k-s18-guarded-round2-attempt1",
"wall_s": 28.60761404898949,
"cpu_s": 25.612138,
"cgroup_peak_bytes": 790102016
},
{
"label": "bound-64k-s18-guarded-round3-attempt1",
"wall_s": 28.79827523091808,
"cpu_s": 25.088846,
"cgroup_peak_bytes": 791957504
}
],
"wall_s": 28.79827523091808,
"cpu_s": 25.612138,
"cgroup_peak_bytes": 791957504
},
"bound-256k": {
"samples": [
{
"label": "bound-256k-s18-guarded-round1-attempt1",
"wall_s": 29.13570951204747,
"cpu_s": 25.580735,
"cgroup_peak_bytes": 845651968
},
{
"label": "bound-256k-s18-guarded-round2-attempt1",
"wall_s": 28.282239037798718,
"cpu_s": 25.22293,
"cgroup_peak_bytes": 784314368
},
{
"label": "bound-256k-s18-guarded-round3-attempt1",
"wall_s": 29.86603315989487,
"cpu_s": 25.689847,
"cgroup_peak_bytes": 803442688
}
],
"wall_s": 29.13570951204747,
"cpu_s": 25.580735,
"cgroup_peak_bytes": 803442688
},
"bound-1m": {
"samples": [
{
"label": "bound-1m-s18-guarded-round1-attempt1",
"wall_s": 29.29350886703469,
"cpu_s": 25.370921,
"cgroup_peak_bytes": 840220672
},
{
"label": "bound-1m-s18-guarded-round2-attempt1",
"wall_s": 28.540352625073865,
"cpu_s": 24.991708,
"cgroup_peak_bytes": 791252992
},
{
"label": "bound-1m-s18-guarded-round3-attempt1",
"wall_s": 29.30297133116983,
"cpu_s": 25.272261,
"cgroup_peak_bytes": 791572480
}
],
"wall_s": 29.29350886703469,
"cpu_s": 25.272261,
"cgroup_peak_bytes": 791572480
}
}
},
"floors": {
"selected": "bound-64k",
"results": [
{
"scale": 18,
"wall_s": 28.79827523091808,
"edges_per_second": 145644.2778731748,
"floor": 101175,
"basis": "three-run S18 median"
},
{
"scale": 19,
"label": "bound-64k-s19-floor-attempt1",
"wall_s": 57.463328373851255,
"cpu_s": 50.350206,
"cgroup_peak_bytes": 1402183680,
"edges_per_second": 145981.93730485762,
"floor": 105415,
"basis": "one qualified complete-ingest observation"
},
{
"scale": 20,
"label": "bound-64k-s20-floor-attempt1",
"wall_s": 113.88196149794385,
"cpu_s": 99.884978,
"cgroup_peak_bytes": 2673209344,
"edges_per_second": 147321.10142222053,
"floor": 105802,
"basis": "one qualified complete-ingest observation"
}
]
},
"excluded_attempts": [
"baseline-s18-round1",
"baseline-s18-round2",
"bound-1m-s18-round1",
"bound-1m-s18-round2",
"bound-1m-s18-round3",
"bound-256k-s18-round1",
"bound-256k-s18-round2",
"bound-256k-s18-round3",
"bound-64k-s18-guarded-round1-attempt1",
"bound-64k-s18-guarded-round1",
"bound-64k-s18-round1",
"bound-64k-s18-round2"
]
}
Loading