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
Original file line number Diff line number Diff line change
Expand Up @@ -508,6 +508,7 @@ mod determinism {
max_batch_rows: 8_192,
max_run_records: 262_144,
partition_count,
target_partition_records: 16,
..GraphConstructionBudgets::default()
};
let (fingerprint, layout) =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1552,6 +1552,7 @@ mod lifecycle_budget {
Some("endpoints.run"),
None,
window_rows,
partition::default_materialization_bytes(),
&mut || false,
&mut evidence,
)
Expand Down
37 changes: 36 additions & 1 deletion crates/graphforge-storage/src/graph_construction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@ use catalog::{
load_parent_runtime_catalog_from_compact,
};
mod partition_load;
mod partition_memory;
mod partition_records;
pub(crate) mod partition_shaping;
mod supersession;
Expand Down Expand Up @@ -473,12 +474,26 @@ 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.
/// Maximum 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.
pub partition_count: u32,
/// Desired identities per sampled range; a planning target, not a byte bound.
#[serde(
default = "partition::default_target_records",
skip_serializing_if = "partition::is_default_target_records"
)]
pub target_partition_records: u64,
/// Maximum accounted materialization bytes per partition load. Fixed-family
/// workers admit at most two loads; Arrow rows are loaded serially. This is
/// separate from staging windows, writer buffers and whole-process memory.
#[serde(
default = "partition::default_materialization_bytes",
skip_serializing_if = "partition::is_default_materialization_bytes"
)]
pub max_partition_bytes: u64,
}

impl Default for GraphConstructionBudgets {
Expand All @@ -495,6 +510,8 @@ impl Default for GraphConstructionBudgets {
max_catalog_decoded_bytes: 256 << 20,
max_catalog_identifier_bytes: 64 << 20,
partition_count: partition::DEFAULT_PARTITION_COUNT,
target_partition_records: partition::default_target_records(),
max_partition_bytes: partition::default_materialization_bytes(),
}
}
}
Expand All @@ -511,6 +528,8 @@ impl GraphConstructionBudgets {
|| self.max_catalog_entries == 0
|| self.max_catalog_decoded_bytes == 0
|| self.max_catalog_identifier_bytes == 0
|| self.target_partition_records == 0
|| self.max_partition_bytes == 0
|| self.partition_count == 0
|| self.partition_count > partition::MAX_PARTITION_COUNT
{
Expand Down Expand Up @@ -1288,6 +1307,22 @@ impl GraphConstructionSession {
Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
Err(error) => return Err(storage(error)),
};
// The former public default recorded a 256-partition maximum. Opening
// that exact legacy default with today's default retains its authority;
// explicit non-default changes still fail validate_checkpoint below.
let legacy_defaults = GraphConstructionBudgets {
partition_count: 256,
..GraphConstructionBudgets::default()
};
let budgets = if budgets == GraphConstructionBudgets::default()
&& recovered_checkpoint
.as_ref()
.is_some_and(|checkpoint| checkpoint.budgets == legacy_defaults)
{
legacy_defaults
} else {
budgets
};
if let Some(checkpoint) = recovered_checkpoint.as_mut()
&& (DetailCodec::from_version(checkpoint.format_version).is_err()
|| checkpoint.operation_uuid != operation_uuid
Expand Down
64 changes: 56 additions & 8 deletions crates/graphforge-storage/src/graph_construction/partition.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,16 +33,50 @@
use super::storage;
use graphforge_core::GfError;

/// Number of range partitions requested by a construction session.
/// Maximum range partitions requested by a construction session.
///
/// 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.
pub(crate) const DEFAULT_PARTITION_COUNT: u32 = 256;
pub(crate) const DEFAULT_PARTITION_COUNT: u32 = 4_096;

/// Upper bound on the recorded partition count.
pub(crate) const MAX_PARTITION_COUNT: u32 = 4_096;

/// Deterministic planning target; not a byte bound or an optimality claim.
/// Admission checks actual bytes independently; skew can exceed this target.
pub(super) const fn default_target_records() -> u64 {
16_384
}
#[expect(
clippy::trivially_copy_pass_by_ref,
reason = "Serde skip_serializing_if requires a reference"
)]
pub(super) fn is_default_target_records(value: &u64) -> bool {
*value == default_target_records()
}
pub(super) const fn default_materialization_bytes() -> u64 {
256 << 20
}
#[expect(
clippy::trivially_copy_pass_by_ref,
reason = "Serde skip_serializing_if requires a reference"
)]
pub(super) fn is_default_materialization_bytes(value: &u64) -> bool {
*value == default_materialization_bytes()
}

/// Fail before allocating the retained partition or its sort/reorder buffers.
pub(super) fn admit_materialization(bytes: Option<u64>, limit: u64) -> Result<u64, GfError> {
let bytes = bytes.ok_or_else(|| storage("partition materialization byte count overflows"))?;
if bytes > limit {
return Err(storage(format!(
"partition materialization requires {bytes} bytes, exceeds recorded budget {limit}"
)));
}
Ok(bytes)
}

/// 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.
Expand All @@ -62,12 +96,9 @@ const BALANCE_MIN_MEAN_ROWS: u64 = 16;

/// Never cut more partitions than the balance check can meaningfully assert on.
///
/// Each partition costs a durable spill artifact, with its own barrier and
/// writer receipt, in every family. Cutting the full recorded count regardless
/// of input size makes that cost a constant: a two-thousand-row graph would pay
/// the same durability price as a sixty-seven-million-row one, and shaping's
/// barrier count would stop tracking the work it protects. Bounding the cut by
/// the recorded record count keeps the cost proportional.
/// Each partition retains spill names and writer state in every family.
/// Bounding the cut by the recorded record count avoids charging tiny inputs
/// for the maximum cut. Directory synchronization is batched across spills.
///
/// 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
Expand Down Expand Up @@ -216,6 +247,23 @@ pub(crate) struct IdentitySampler {
}

impl IdentitySampler {
/// Adapt within the recorded maximum, retaining the small-input balance floor.
pub(crate) fn with_target(
partition_count: u32,
total_records: u64,
target: u64,
) -> Result<Self, GfError> {
validate_partition_count(partition_count)?;
if target == 0 {
return Err(storage("partition target must be positive"));
}
let desired = total_records
.div_ceil(target)
.max(256)
.min(u64::from(partition_count));
Self::new(u32::try_from(desired).map_err(storage)?, total_records)
}

/// Build a sampler for `total_records` staged identities.
///
/// # Errors
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -236,3 +236,65 @@ fn every_partition_is_reachable_for_a_well_sampled_ingest() {
balance.rows()
);
}

#[test]
fn adaptive_cut_grows_from_recorded_data_and_stops_at_recorded_ceiling() {
for (records, expected) in [
(0, 1),
(32, 2),
(1 << 20, 256),
(1 << 23, 512),
(1 << 24, 1024),
(1 << 26, 4096),
(1 << 30, 4096),
] {
let sampler =
IdentitySampler::with_target(4096, records, default_target_records()).unwrap();
assert_eq!(sampler.cut(), expected);
assert_eq!(
sampler.positions().collect::<Vec<_>>(),
IdentitySampler::with_target(4096, records, default_target_records())
.unwrap()
.positions()
.collect::<Vec<_>>()
);
}
assert_eq!(
IdentitySampler::with_target(64, 1 << 30, 16).unwrap().cut(),
64
);
assert!(IdentitySampler::with_target(4096, 100, 0).is_err());
}

#[test]
fn materialization_budget_refuses_overflow_and_one_byte_over_limit() {
assert_eq!(admit_materialization(Some(1024), 1024).unwrap(), 1024);
assert!(admit_materialization(Some(1025), 1024).is_err());
assert!(admit_materialization(None, u64::MAX).is_err());
}

#[test]
fn legacy_budget_serialization_is_unchanged_and_new_limits_round_trip() {
let legacy = super::super::GraphConstructionBudgets {
partition_count: 256,
..Default::default()
};
let encoded = serde_json::to_value(legacy).unwrap();
assert!(encoded.get("max_partition_bytes").is_none());
assert!(encoded.get("target_partition_records").is_none());
assert_eq!(
serde_json::from_value::<super::super::GraphConstructionBudgets>(encoded).unwrap(),
legacy
);
let custom = super::super::GraphConstructionBudgets {
max_partition_bytes: 12345,
target_partition_records: 123,
..legacy
};
let encoded = serde_json::to_value(custom).unwrap();
assert_eq!(encoded["max_partition_bytes"], 12345);
assert_eq!(
serde_json::from_value::<super::super::GraphConstructionBudgets>(encoded).unwrap(),
custom
);
}
Loading