Skip to content

Staged parallel execution in the current SPARQL evaluator #1682

Description

@patrykorwat

Is your feature request related to a problem? Please describe.

Oxigraph today runs single threaded in every engine layer. The only concurrent work in the tree is on the ingestion path: the RocksDB bulk loader in lib/oxigraph/src/storage/rocksdb.rs uses std::thread::spawn to build SST files from multiple batches at once (capped by max_num_threads), lib/oxttl exposes Send chunk iterators so a caller can drive parallel parsing with rayon, and the CLI uses rayon_core::ThreadPoolBuilder to fan out across multiple input files during load. SPARQL evaluation, storage range scans, query planning, the transaction commit path, and reasoning all use one core.

On multi-million triple stores this leaves most of the machine idle during long running queries and reasoning runs. On a 16 core workstation a single BGP scan over a large store pins one core at 100 percent while the other 15 sit at 0. Large aggregate queries, property path closures, and OWL 2 RL forward chaining over GeoSPARQL data show the same pattern.

This ticket proposes a staged path to parallelism inside the current evaluator as the forward direction for oxigraph. It is offered as an alternative to the DataFusion investigation in #780 and #1527, which has surfaced useful design questions but starts from a model (relational OLAP with SQL semantics) that cannot hold the semantics of RDF and OWL without forking DataFusion far enough that the reuse argument disappears. Open world assumption, entailment as query semantics, the open schema over URIs, and global URI identity are not peripheral features to patch in; they are the defining properties of an RDF store. The detailed argument is in Additional Context. The short version: parallelism inside sparopt and the current evaluator serves the full semantics of RDF and OWL at scale, lands in four independently useful stages, and keeps every existing SPARQL-specific rewrite intact.

Prior attempts to push rayon into the engine were made and reverted before landing. The blocker is not lack of interest but the absence of a design that addresses whatever broke those attempts.

Describe the solution you'd like

A staged plan gated behind a Cargo feature, default off, with four independently landable stages. Rule evaluation inside oxreason stays single threaded throughout. RocksDB writes stay serialised through a single writer thread.

Stage 0, investigation. Document in this issue which subsystems the prior rayon attempts touched, why each was reverted, and what benchmarks were run. Every later stage references which blocker it solves. Without this the work just relearns the same lessons.

Stage 1, dedicated writer thread with SPSC channel. Move RocksDB batch commits off the thread that produces them, using the same pattern the bulk loader already uses. Applies to Store::load, Store::bulk_load callers that stream in, and the sink path for Store::reason once it lands. Not a thread pool. One writer, one channel, clean RocksDB semantics.

Stage 2, parallel range scans inside a single query. Split one quads_for_pattern scan into N key prefix partitions driven by worker threads against one RocksDB snapshot. Preserve ordering for operators that need it, merge freely for hash join build sides, aggregates, and DISTINCT. This is where the Send and iterator lifetime issues that likely killed the prior attempts have to be solved properly. If this stage does not land, the rest of the stack does not land.

Stage 3, parallel hash join build side. Build the hash table from partitioned build input on N workers, probe single threaded on the main thread. The probe side stays sequential because that is where result ordering determinism is cheapest to preserve.

Stage 4, cost threshold in the existing optimizer. lib/sparopt/src/optimizer.rs already picks physical variants (see JoinAlgorithm, LeftJoinAlgorithm). Add one more branch that selects between sequential and parallel variants based on an estimated row count threshold. This is a small extension to the existing optimizer, not a new planner. Cardinality estimates come from the same per predicate counts the join reorderer already uses, and extending them with characteristic-set style summaries (RDF-3X, Virtuoso) is the standard path for triple stores and has no DataFusion equivalent worth reusing.

Each stage must:

  1. Preserve determinism of query results. Every parallel variant has a deterministic merge step. Tests run with --test-threads=1 in CI on the parallel feature flag.
  2. Stay behind a Cargo feature so minimal builds do not pull rayon into the engine.
  3. Add no regression on small query microbenchmarks on the default path.
  4. Carry its own crash test where a worker panic has to drain cleanly.

Describe alternatives you've considered

The DataFusion evaluator in #780 and #1527. Rejected as a direction for oxigraph. DataFusion is a relational OLAP engine with SQL semantics. Wrapping it in SPARQL syntax does not make it an RDF store. Open world assumption, entailment as query semantics, open schema over URIs, and global URI identity are not operator-level details to be patched in; they are the defining properties of an RDF and OWL system. A DataFusion backed evaluator either sheds those properties (at which point it is a fast triple-to-SQL bridge, not an RDF store) or reimplements them on top (at which point it is a bespoke planner wearing DataFusion as a coat, with all the Arrow decoding tax and none of the upstream leverage). The framing of #1527 as "lives alongside" is honest about this: it is a separate system for a different semantic contract, not the future of the evaluator. For workloads that legitimately want SPARQL-over-Parquet OLAP with pre-materialised closure and no reasoning, a separate project on top of DataFusion is fine. For the oxigraph evaluator proper, parallelism in sparopt and the current engine is the forward path.

Rayon inside the oxreason rule evaluation loop. Rejected for this ticket. Rule body parallelism has its own determinism questions, the maintainer has already rejected the naive version, and it needs its own design discussion after Stage 2 lands.

Multi-threaded RocksDB writes (multiple concurrent write_batch calls). Rejected. RocksDB write semantics favour a single writer; the Stage 1 SPSC pattern is the safe version of the same idea.

Global rayon opt in with no planner threshold. Rejected. Prior experience shows pool overhead hurts point queries enough that a global feature flag either gets left off by everyone (no benefit) or hurts the default path (regression).

A background query pool for unrelated queries running concurrently. Out of scope. That is a transactional concurrency question, not a single query parallelism question, and it needs its own ticket.

Threads managed by hand without rayon. Considered. Rayon's work stealing pool is well understood and the CLI already pulls it in through rayon_core, so adding it to one more subsystem behind a feature flag is cheaper than rolling a bespoke pool. The bulk loader's std::thread::spawn pattern is the right template for Stage 1 specifically (single writer thread, no stealing needed). If the investigation phase surfaces a reason rayon specifically was the problem (rather than parallelism in general), this choice should be revisited.

Additional context

Why this ticket is a replacement for the DataFusion direction in #780 and #1527, not a complement to it:

  1. Open world vs closed world assumption. SQL operates under the closed world assumption; RDF and OWL operate under the open world assumption. This changes the semantics of NOT EXISTS, MINUS, FILTER(!BOUND(?x)), counting over negation, and anything that reasons about what is not asserted. DataFusion cannot be taught this by adding operators because the assumption is baked into how the query is interpreted end to end.
  2. Entailment is part of query semantics, not a pre-pass. Even without an explicit reasoner, RDF and OWL queries are evaluated under an entailment regime: rdfs:subClassOf propagation, rdfs:domain and rdfs:range inferences, owl:sameAs equivalence. A DataFusion evaluator that materialises closure up front works only for finite static closures. Streaming data, OWL 2 DL, owl:sameAs smushing, and any setup that cannot afford to double storage by materialising everything are all dead ends on that path.
  3. The schema is open. SQL assumes a fixed column set per relation; RDF is effectively EAV where every subject can carry any predicate, predicates are first class URIs, and the predicate set is not known up front. The Arrow schema in Figure out an RDF representation in Apache Arrow #1463 handles this by grouping by predicate into dense buffers, which works at read time but makes transactional inserts of new predicates expensive.
  4. URIs are global identifiers, not primary keys. Two triples from two different sources with the same IRI refer to the same resource. This is how linked data, owl:sameAs, and federated queries work. Lowering to DataFusion either loses globality or rebuilds it on top, and rebuilding it is most of what an RDF store does.
  5. Reasoning changes what a query means. Under OWL 2 RL forward chaining, queries are evaluated against the closure. The evaluator has to know about closure materialisation, semi-naive delta state, and the T-Box to evaluate correctly.
  6. Property paths with full SPARQL 1.1 syntax (union, inverse, negation, sequence, Kleene star) are iterative closures over the triple store. Lowering them to DataFusion recursive CTEs is awkward for the non-trivial cases and already flagged as incomplete in New query evaluator using DataFusion #1527.
  7. Blank node scoping is document local and blank nodes are generated during CONSTRUCT. Arrow can carry a blank node ID but the scoping rules are evaluator owned.
  8. SPARQL 1.1 operators that depend on evaluator-owned semantics (SERVICE, EXISTS/NOT EXISTS, MINUS, GRAPH, FROM, FROM NAMED) either need custom DataFusion physical operators or pre-rewriting. New query evaluator using DataFusion #1527 explicitly calls out EXISTS as spotty today.
  9. Cardinality estimation for triple stores is a distinct research lineage (RDF-3X characteristic sets, Virtuoso's cost model, Jena TDB's per-predicate stats). Relational histogram-per-column models are a poor fit because every query reads the same (s, p, o) table and selectivity is dominated by predicate distribution. sparopt can evolve along this lineage; DataFusion's optimizer cannot without forking it.
  10. Write semantics. Transactional single-triple insert/delete over RocksDB is a different workload from Parquet batch writes. Tpt's own "read-only workloads first" framing concedes this.

All ten of these apply whether reasoning is on or off, whether the workload is analytical or transactional, and whether the closure is pre-materialised or not. They are properties of RDF and OWL as a data and knowledge model, not of a particular query workload. An evaluator that does not honour them is not an RDF evaluator in any useful sense. sparopt can evolve to handle them at scale with parallelism, characteristic-set cardinality, and the existing RocksDB snapshot model. A DataFusion backed evaluator cannot, without forking DataFusion far enough that the reuse argument disappears.

Existing concurrent code in the tree, for reference:

  1. lib/oxigraph/src/storage/rocksdb.rs around line 1392, RocksDbStorageBulkLoader::load_batch uses std::thread::spawn to run FileBulkLoader::new(...).load(...) concurrently, with a FIFO queue of handles capped by max_num_threads. Comment on line 1419 says // TODO: better spawn. This is the template for the Stage 1 SPSC writer thread pattern.
  2. lib/oxttl exposes Send chunk iterators and documents rayon usage in its parser doc examples, but the crate itself spawns no threads; rayon is a dev dependency only.
  3. cli/src/main.rs line 159 uses rayon_core::ThreadPoolBuilder to fan out across multiple input files during load.

None of these touch the query engine, storage scans, or reasoning. This ticket would be the first place rayon enters an engine layer.

Comparable systems that went the parallel native triple store route:

  1. RDFox parallelises rule evaluation as well as scans, with a published determinism model. https://www.oxfordsemantic.tech/rdfox
  2. Virtuoso runs parallel scans on its partitioned column store and exposes ThreadsPerQuery in the INI. https://docs.openlinksw.com/virtuoso/rthrcfg/
  3. Apache Jena TDB2 evaluates BGPs with a thread pool behind a config flag (tdb2:threadPoolSize). https://jena.apache.org/documentation/tdb2/

Relevant RocksDB documentation:

  1. Snapshot semantics and thread safety. https://github.com/facebook/rocksdb/wiki/Snapshot
  2. Iterator thread model. https://github.com/facebook/rocksdb/wiki/Iterator

Related oxigraph issues:

  1. Investigate using DataFusion #780 DataFusion investigation. This ticket argues for a different direction; see Additional Context above.
  2. Figure out an RDF representation in Apache Arrow #1463 RDF in Arrow schema design. Depends on Investigate using DataFusion #780 landing; orthogonal to this ticket.
  3. New query evaluator using DataFusion #1527 DataFusion evaluator draft PR. This ticket argues for a different direction; see Additional Context above.

Benchmarks that should gate the work: the LUBM synthetic bench, the spatial reasoning bench in bench/geosparql, a property path microbench over transitive rdfs:subClassOf chains, and a point query microbench added as part of Stage 2 to catch regressions on the default path.

Activity

  1. Tpt commented on Apr 22, 2026

    @Tpt
    Collaborator

    Thank you for opening this topic.

    I quite disagree that DataFusion is not fit for SPARQL evaluation. #1527 passes most of the SPARQL testsuite outside of things around functions (not implemented yet) or EXISTS (the decorrelation implementation in DataFusion is only partial but doing it for general SQL is possible, see "Unnesting Arbitrary Queries" by Thomas Neumann). There have been research papers on SPARQL -> SQL rewriting. Similarly, I don't see what prevents implementing OWL 2 QL (by SPARQL rewriting) or OWL 2 RL (one can rewrite rule languages like datalog into relational algebra + the least fixed point operator provided by WITH RECURSIVE) on top of DataFusion. Also, SPARQL operates under the closed world assumption.

    However, I agree it's not the best fit for RocksDB. RocksDB is optimized for mixed read/write OLTP workload and DataFusion assumes that large columnar reads are fast. I picked RocksDB for Oxigraph mostly to get something stable and ready to use (it was the start of RDF in Rust, I had to write everything from scratch) and as a reaction to Blazegraph. This makes Oxigraph one of fastest store for writes, But, retrospectively I am not sure it was the good choice: complex SPARQL queries, especially with OWL reasoning, look way more like OLAP SQL queries than OLTP. So, I am considering writing an other storage backend more OLAP-like, hence the stalled MR. It would likely rely on Arrow and Parquet and be a LSM tree like RocksDB. In my mind "spareval" would become the "reference" simple implementation, easy to compile/embed and also very useful for things like differential fuzzing (amazing to have when testing complex optimizations and for running coding agents).

    Also, I am a bit scared to put a lot of effort on writing a full query evaluator and optimizer from scratch, this is a lot of work and the DataFusion community has made a lot of effort on it. Reusing that looked like a good idea to me (Oxigraph is still a one guy hobby project). So, my main concern is "what is the highest yield path for the project to keep being alive in 10 years while not requiring too much work".

    But my mind is not set on that at all, I would love convincing arguments. I am not sure that what I believe are sometime LLM hallucinations (please correct me if I assume wrongly, I am deeply sorry if I am wrong) is helping there.

  2. patrykorwat commented on Apr 22, 2026

    @patrykorwat
    ContributorAuthor

    Thanks for the thorough explanation. I acknowledge the reality: a fine tuned bespoke planner is not something the project can carry, DataFusion first makes sense.

    I'd suggest leaving this ticket open rather than closing it. Further advancements in LLM assisted development could make this kind of work a one evening job rather than a quarters long commitment. If that inflection arrives, the calculus flips and a bespoke parallel path becomes viable again. Worth keeping the thread as a pointer for that revisit.

    On the concrete ask for #1560: I'll keep the WKB storage aligned with GeoArrow's wkb extension type, uncompressed little endian bytes in the wkbs side CF with no oxigraph specific framing inside the blob, so a future GeoArrow column can borrow the blob directly.

  3. Tpt commented on Apr 23, 2026

    @Tpt
    Collaborator

    Further advancements in LLM assisted development could make this kind of work a one evening job rather than a quarters long commitment.

    Yes! This is why I am very keen of keeping a small "reference executor" implementation even if we move to e.g. DataFusion. This makes correctness testing much easier.

    I'll keep the WKB storage aligned with GeoArrow's wkb extension type

    Amazing!

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

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions