Skip to content

epic(storage): evaluate DataFusion, Tokio, and Arrow replacements for custom construction machinery #1504

Description

@DecisionNerd

Purpose

Comprehensively evaluate replacing GraphForge's custom construction sorting, partitioning, spill management, and scheduling with DataFusion, Tokio, and Arrow facilities. Deliver evidence-backed decisions for each mechanism, including performance, correctness, and maintenance tradeoffs. Neither migration nor retention is the predetermined answer; partial reuse and hybrid designs are valid outcomes.

Why this is new work

#1456 improved the existing construction implementation and closed with bounded experiments and documented dispositions. It did not establish that the custom engine is best practice or that off-the-shelf facilities generally harm performance. #1465 evaluated a particular encode integration; its slower result cannot reject reuse in shaping, sorting, partitioning, spilling, or scheduling. Arrow was already used; the work did not migrate construction to DataFusion or Tokio.

The unresolved architectural question deserves its own close gate. This epic addresses architecture and proof debt, not just throughput tuning.

Scope and work packages

1. Inventory and decision protocol

  • Map every custom sorting, partitioning, spilling, scheduling, buffering, cancellation, and memory-accounting mechanism in the construction path to its callers, invariants, tests, and maintenance burden. Identify existing library reuse separately from possible new reuse.
  • Establish the current-main baseline and exact dependency versions. Consult version-matched upstream documentation and implementation; link primary sources for capability and limitation claims.
  • Build a mechanism-by-candidate matrix: current implementation, applicable DataFusion facilities, applicable Arrow kernels/batch facilities, applicable Tokio coordination facilities, and meaningful hybrid designs. These libraries have different roles; do not assume each replaces every mechanism. Include the existing Rayon execution arrangement as a scheduling baseline.
  • Separate execution partitioning from durable layout, generic mechanisms from graph semantics, transient spill from recovery authority, and library-accounted memory from total admitted memory. Challenge existing implementation boundaries where useful without silently changing product contracts.
  • Before measuring, record workloads, resource envelopes, repetitions, uncertainty treatment, regression thresholds, correctness gates, and how maintenance benefits will be weighed against performance costs. Do not invent a universal speedup threshold or infer one from historical experiments. Explicitly record maintainer acceptance of any proposed production tradeoff that changes an existing requirement.

2. Sorting and partitioning evaluation

  • Prototype representative construction/shaping operations using suitable library sorting and partitioning facilities, including selective kernel reuse and integrated execution alternatives.
  • Compare allocation/copy costs, ordering semantics, global merges, repartitioning, skew behavior, and memory admission against the current implementation.
  • Cover power-law hubs and partitions exceeding the current fail-closed budget; distinguish safe refusal from successfully processing a larger input. Evaluate bounded processing/spill alternatives without assuming adaptive partitioning alone solves skew.
  • Explain how recorded UUID splitters, surrogate assignment, endpoint resolution, and canonical publication interact with each alternative. Do not impose byte-identical transient intermediates where the publication contract does not require them.

3. Spill, buffering, and memory evaluation

  • Prototype applicable external-sort/spill and memory-pool facilities under forced memory pressure, not merely in-memory inputs.
  • Account for input/output batches, sort state, queues, completed results awaiting consumption, compression and writer buffers, adapters, spill files, and unregistered allocations. Distinguish reservation bounds, observed RSS, and page cache.
  • Evaluate transient-file ownership, cleanup, disk limits, corruption detection, descriptor authentication, crash/restart behavior, and coexistence with GraphForge receipts and recovery. A library spill file is not automatically a durable checkpoint.
  • Measure bytes copied/read/written, spill amplification, file counts, durability barriers, simultaneous scratch, and refusal behavior.

4. Scheduling and cancellation evaluation

  • Compare current coordination with applicable Tokio task/stream/channel facilities and DataFusion scheduling, including hybrid CPU execution where appropriate.
  • Exercise shared process CPU admission, bounded queues by bytes, bounded completed-result retention, backpressure, fairness, cancellation latency, error propagation, shutdown, and concurrent imports/queries.
  • Verify actual behavior of blocking work and nested runtimes/pools from code and tests. An async API alone is not evidence of parallel CPU execution or cheaper I/O.

5. Integrated experiments and tradeoff report

  • Integrate the most credible compatible alternatives into the real Rust construction/import path in isolated experimental worktrees. Measure interactions between mechanisms; isolated microbenchmarks are supporting evidence only.
  • Compare complete import through publication and reopen/query verification against the same baseline under equal CPU, memory, storage, durability, and semantic constraints.
  • Produce a maintenance assessment for every alternative: custom code and tests deleted versus adapter code added, ownership boundaries, invariant duplication, dependency/API stability, upgrade effort, diagnostics, debugging, portability, build/binary cost, and required upstream changes. Code-line reduction alone is insufficient; support judgments with prototype diffs and concrete change/upgrade examples, marking estimates explicitly.
  • Write an ADR and a concise decision matrix: adopt / retain / hybrid, evidence, drawbacks, confidence, limitations, and revisit triggers for each mechanism. Include a staged migration/rollback plan and bounded implementation issues for accepted replacements. Production migration is a separately tracked outcome, not something to imply from completing this evaluation.

Evidence and correctness requirements

Use small deterministic fixtures plus representative Graph500 and property-bearing inputs: empty/tiny, normal, duplicate-heavy, skewed/hub, variable-width/nested properties, and inputs large enough to force spill under constrained memory. Include repeated quiet-host paired measurements at S18–S20 and at least one larger or tighter-memory stress case; state the exact supported scale, and record failures and environmental limits. Do not claim S26 qualification from these experiments.

Measure end-to-end and phase wall time, CPU-seconds, effective concurrency, peak RSS and reservation bounds, copies, I/O, scratch, fsyncs, cancellation, and recovery cost. Preserve commands, input identities, revisions, configuration, raw observations, and variability. Use natural quiet windows around other sessions; do not interrupt their builds. Compare semantically equivalent work, including successful published row counts.

Correctness is a hard gate: graph identities, duplicate policy, surrogate/endpoint mappings, properties, canonical published artifacts where required by ADR 0038, query answers, CAS/receipt integrity, structured errors, and atomic publication must remain valid. Test forced worker counts and schedules; cancellation at distinct stages; malformed/truncated/corrupt spills; disk exhaustion; admission refusal; and interrupted/resumed versus uninterrupted imports. Demonstrate that failures publish no partial graph and that reopen/recovery preserves the contract. Review temporary-file and descriptor handling without weakening authentication for benchmark convenience.

Likely surfaces: crates/graphforge-storage/src/graph_construction*, construction_directory.rs, construction determinism/lifecycle tests, and API import/resource-policy code. Rust continues to own behavior; bindings remain thin. Use repository-required validation and exact-head CI for merged experiment infrastructure or production changes. Strategically request CodeRabbit review for invariant boundaries, recovery adapters, and the final experimental implementation; independently verify findings.

Acceptance criteria / close gate

  • Inventory and version-grounded candidate matrix cover all four principal mechanisms: sorting, partitioning, spilling, and scheduling, plus their memory/cancellation interactions.
  • The comparison protocol is recorded before the benchmark results and any later changes are explained rather than retrofitted to a favored result.
  • Every principal mechanism has a runnable representative replacement/hybrid experiment, or a specific independently reviewed correctness/API incompatibility proof for an inapplicable candidate. Poor encode performance, speculation, or a time-box expiry is not such a proof.
  • At least one credible combined design is exercised through real construction, publication, reopen, and queries; if no combination is viable, concrete incompatibility evidence accounts for every otherwise credible route.
  • Correctness, failure, recovery, and resource-bound tests have direct evidence; unsupported cases and regressions are explicit.
  • Fair repeated performance measurements and maintenance/correctness tradeoff assessments are available per mechanism and for the integrated candidate(s).
  • A reviewed ADR records bounded conclusions, accepted/rejected tradeoffs, uncertainty, and revisit triggers. No result is generalized to all workloads or all library reuse.
  • Accepted migration work has scoped follow-up issues and sequencing; unresolved evaluation questions remain open under this epic rather than being relabeled as implementation follow-ups to permit closure.

A negative conclusion may close this evaluation only with the evidence above. Retaining custom code is a decision to justify, not the default proof of success. Completing this epic does not itself prove a migration shipped, the #1387 ingest floor was met, or the implementation is universally best practice.

Completion scenarios

  1. Given a candidate passes a sorting microbenchmark, when its construction integration is evaluated, then adoption advice includes complete-ingest performance, semantic equivalence, memory pressure, and recovery evidence.
  2. Given a skewed property-bearing import that forces spill, when workers are rescheduled, cancelled, or interrupted, then output or safe refusal follows the declared contract, partial state is not published, and recovery matches an uninterrupted result.
  3. Given an alternative reduces custom code but costs resources or adds adapter complexity, when the ADR recommends it, then the tradeoff is quantified where possible, uncertainty is explicit, and changes to acceptance requirements require an explicit maintainer decision.
  4. Given one library integration is slower, when the evaluation closes, then rejection is limited to the tested mechanism/design/envelope and other credible candidates have their own evidence.

Tracking and boundaries

This epic is the canonical close gate for this evaluation. Split work into native sub-issues only for independently reviewable concerns needed by these criteria; make them block this epic, avoid duplicated implementation ownership, and follow the repository WIP limit and current-main worktree policy. Inventory/protocol precedes prototypes; integrated evidence and the ADR follow them.

Related, not duplicates:

Non-goals: mandatory wholesale engine replacement; rewriting graph semantics or public APIs; distributed execution; changing durability guarantees to win a benchmark; implementing every recommended production migration within this evaluation; unrelated M11 work. No new release-blocking milestone assignment is implied.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    coreCore source code changestestingTest coverage and testing infrastructure

    Type

    No type

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions