feat(storage): worker-local partition loads with coordinator-owned evidence (#1456 item 1) - #1459
Conversation
…cal evidence Work item 1 of #1456. Workers load and sort independent partitions and return their own results plus local counters; the coordinator consumes them in canonical partition order and owns the shared evidence, the allocation ledger and publication. The merge of worker counters into the evidence is explicit and total (`PartitionLoadCounters::merge_into`). `consume_in_partition_order` bounds running loads plus loaded-but-unconsumed results plus the partition being written to `workers`, the bound #1448's reorder buffer claimed and did not deliver. Cancellation keeps the caller's `FnMut() -> bool` on the coordinator; workers poll a stop flag. The allocation ledger is untouched by workers, so the exact coexistence peak is unchanged. A test-only `materialized_records_bound` names the order-independent reservation figure for the admission gate that does not exist yet. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…d a held in-flight load Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
|
Important Review skippedAuto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: CHILL Plan: Advanced Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
Warning Billing warning: we have not been able to collect payment for this subscription for more than 72 hours. Please update the payment method or pay any pending invoices in Billing to avoid service interruption. Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
) The previous commit showed the surrogate chain is ~9% of shaping. This scopes the rest, so the other ~91% is attributed rather than assumed. 131,072 identities, debug build, quiet host: | region | wall ms | cpu ms | eff cores | % shaping | |-----------------------------------|--------:|-------:|----------:|----------:| | shape.chunk_loop | 9264.1 | 9080.0 | 0.98 | 9.3% | | **shape.loop_to_chain** | 44034.4 |26570.0 | **0.60** | **44.2%** | | chain.validate_staged_details | 947.7 | 950.0 | 1.00 | 1.0% | | chain.assign_surrogates | 687.2 | 660.0 | 0.96 | 0.7% | | chain.resolve_endpoint_surrogates | 7720.7 | 5220.0 | 0.68 | 7.8% | | shape.after_chain | 26116.5 |17570.0 | 0.67 | 26.2% | | shaping | 99543.1 |70840.0 | 0.71 | 100.0% | 89.2% of shaping wall is accounted. **`shape.loop_to_chain` is the dominant region: 44.2% of shaping at 0.60 effective cores.** It is the four `finish_optional` calls -- identities, node details, edge details, endpoints -- each consuming, sorting and concatenating that family's partitions. It is 4.6x the whole surrogate chain. That re-points the critical path, and the mechanism for it already exists. #1459 shipped `consume_in_partition_order(partitions, workers, load, consume)` with a bounded window and schedule-independent evidence, and #1456 names `RowRangePartitioner::finish` as the next candidate for that mechanism, "same shape, `RecordBatch` is `Send`". This measurement says that candidate is the largest single region in shaping rather than a follow-up. `shape.after_chain` is second at 26.2% and 0.67 effective cores, so the two regions outside the chain are 70.4% of shaping between them. Caveat unchanged from the previous commit: 131,072 identities in a debug build. The ratios are measurements, but the regions scale differently, so this does not establish the split at S20/S22. It does establish that the chain is not where the serial fraction lives at this scale, and that `finish_optional` is. Refs #1464, #1456 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Summary
Work item 1 of #1456: separate worker results from shared evidence, on the narrowest slice that exercises the real shape — the fixed-width
finish_optionalloop ingraph_construction/partition_shaping.rs, whoseevidence.clone()→ mutate →*evidence = committedacross an fsync is the structural obstacle #1429 and #1448 died on.load_fixed_partitionnever seesGraphConstructionEvidence; it returns the sorted partition and aPartitionLoadCounters(plain sums plus per-partition maxima). The merge is explicit and total:PartitionLoadCounters::merge_intowrites exactly nine evidence fields, pinned bymerge_is_total_and_touches_nothing_else.consume_in_partition_order(newgraph_construction/partition_load.rs) runs loads on a bounded pool and hands results to the calling thread in canonical partition index order. The ordered critical section is now: merge counters → write the partition. Decoding, sorting and bulk reads run outside it. No worker touches the allocation ledger, so the exact coexistence peak and ordered transition log are byte-for-byte what a sequential pass records.workerspartitions are materialized at any instant, counting loads in flight, loaded-but-unconsumed results, and the one being written — thethreadsbound perf(storage): make one complete ingest scale with available cores #1448's reorder buffer claimed and did not deliver. A slow first partition holds the pool at the bound (test:slow_first_partition_is_consumed_first_and_holds_the_pool_at_the_bound).FnMut() -> boolstays on the coordinator; workers poll anAtomicBoolevery 4096 records.materialized_records_bound(workers, partitions, peak_partition_records)names the order-independent reservation figure. It is test-only for now: no production budget governs materialized partition records (max_run_recordsbounds staged runs), so there is no admission gate to hand it to yet.Default is
PARTITION_LOAD_WORKERS = 2(one partition loading and sorting while the coordinator writes the previous one). Row partitions are unchanged.Proof
fixed_partition_finish_is_schedule_independent_across_worker_countsroutes 4,096 keys into 16 partitions under forced worker counts 1, 2, 3, 5 and asserts equal output SHA-256, equal relabelled evidence (including the in-memory transition log andstorage_transient_peak_total_allocated_bytes), and that replaying the ordered ledger reaches the recorded exact peak with ≥ 2 identities coexisting.peak_partition_recordsfrom the merge fails 2; a worker forgetting its read counters fails the schedule-independence test.partition_load::tests9,tests::determinism8,tests::lifecycle_budget22,partition::tests12,recovery::tests5,tests::57,shape::tests10; unfilteredgraphforge-storage1169 passed / 0 failed / 2 ignored.cargo test -p graphforge-api --test bdd: TCK 3897 of 3897.cargo clippy --workspace -- -D warningsclean;cargo fmt --all -- --checkclean;make pre-push-fastpassed.No timings were taken: the host was running the S24–S26 ladder.
Related: #1456, #1429, #1448.
🤖 Generated with Claude Code
Need help on this PR? Tag
@codesmith-botwith what you need. Autofix is disabled.