java.lang.Object
com.darkcollective.relix.processor.internal.QueryExecutor

public final class QueryExecutor extends Object
Top-level orchestrator that drives query execution against a validated SemanticModel.

The engine is streaming-first. The executeStreaming methods hand each query statement's output to a RowStreamConsumer as a lazy Stream of rows, so a query's result is never fully buffered by the executor itself — a consumer that processes rows incrementally (writing them out, counting them, etc.) runs in constant memory regardless of result size. The convenience execute methods are thin wrappers that collect each stream into a QueryResult for callers that genuinely want the whole result in memory.

For each query statement in SemanticModel.rootQueries() a result is produced in declaration order. Queries are executed one at a time: a query's row stream is fully consumed (and closed) before the next query runs.

Precondition

The model (or semantic result) passed to any execute/ executeStreaming method must be fully valid — i.e. SemanticResult.isFullyValid() must return true. If a SemanticResult with errors is supplied, an IllegalArgumentException is thrown immediately before any execution begins. Passing a raw SemanticModel that was produced from an invalid analysis may cause EvaluationExceptions that would have been caught during semantic analysis.

Usage


 SemanticResult result = analyzer.analyze(inputStream);

 // Streaming: rows are written out incrementally, never all buffered.
 executor.executeStreaming(result, (label, schema, rows) ->
         rows.forEach(row -> writer.write(format(row))));

 // Convenience: collect the whole result into memory.
 List<QueryResult> results = executor.execute(result);
 

Thread safety

This class is stateless and therefore thread-safe.

  • Constructor Details

    • QueryExecutor

      public QueryExecutor()
      Creates an executor reading the system UTC clock.
    • QueryExecutor

      public QueryExecutor(Clock clock)
      Creates an executor whose current-time built-ins read clock.

      Pinning the clock makes a run reproducible: a script with a relative window (NOW() - DURATION 'PT1H') selects the same rows however much later it is replayed. The CLI exposes this as --now / RELIX_NOW.

      Parameters:
      clock - the clock to read; must not be null
  • Method Details

    • executeStreaming

      public void executeStreaming(SemanticResult result, DataSourceConnector connector, RowStreamConsumer consumer)
      Streams the results of all query statements in result to consumer, using the supplied connector for external relations.
      Parameters:
      result - a fully-valid semantic analysis result; must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      Throws:
      IllegalArgumentException - if result contains semantic errors
      EvaluationException - if a runtime data-level error occurs while a stream is read
    • executeStreaming

      public void executeStreaming(SemanticResult result, RowStreamConsumer consumer)
      Streams the results of all query statements in result to consumer using an inline-only connector.
      Parameters:
      result - a fully-valid semantic analysis result; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      Throws:
      IllegalArgumentException - if result contains semantic errors
      EvaluationException - if a runtime error occurs, including external-source access
    • executeStreaming

      public void executeStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer)
      Streams the results of all query statements in model to consumer, using the supplied connector for external relations.

      Each query's row stream is opened, passed to consumer, and closed before the next query begins. The consumer must consume the stream within the callback (see RowStreamConsumer).

      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      Throws:
      EvaluationException - if a runtime data-level error occurs while a stream is read
    • executeStreaming

      public void executeStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds)
      Like executeStreaming(SemanticModel, DataSourceConnector, RowStreamConsumer) but aborts any FIX fixpoint that exceeds maxFixpointRounds iterations with an EvaluationException.
      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      maxFixpointRounds - maximum semi-naïve fixpoint iterations; use ExecutionContext.UNLIMITED_FIXPOINT_ROUNDS for no cap
      Throws:
      EvaluationException - if a runtime data-level error occurs or a fixpoint cap is hit
    • executeStreaming

      public void executeStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds, QueryEventListener listener)
      Like executeStreaming(SemanticModel, DataSourceConnector, RowStreamConsumer, int) but reports the run's PLAN and EXECUTE decisions to listener, each tagged with the owning query's label.

      This is the plain execution path's observability hook: unlike trace(com.darkcollective.relix.semantic.SemanticModel, java.util.List<com.darkcollective.relix.ast.RelNode>, com.darkcollective.relix.events.QueryEventListener), it neither plans a second time nor drains rows on the caller's behalf, so a host can watch an ordinary run without changing what that run does. No OPTIMIZE events reach the feed here — this method runs whatever trees it is given, and rewriting them is the optimizer's own pass, which reports to its own listener.

      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      maxFixpointRounds - maximum semi-naïve fixpoint iterations; use ExecutionContext.UNLIMITED_FIXPOINT_ROUNDS for no cap
      listener - notified on each planning/execution decision; must not be null
      Throws:
      EvaluationException - if a runtime data-level error occurs or a fixpoint cap is hit
    • executeStreaming

      public void executeStreaming(SemanticModel model, RowStreamConsumer consumer)
      Streams the results of all query statements in model to consumer using an inline-only connector.
      Parameters:
      model - the semantic model; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      Throws:
      EvaluationException - if a runtime error occurs, including external-source access
    • execute

      public List<QueryResult> execute(SemanticResult result, DataSourceConnector connector)
      Executes all query statements in result, using the supplied connector, and collects each result into a QueryResult.
      Parameters:
      result - a fully-valid semantic analysis result; must not be null
      connector - provides rows for external data sources; must not be null
      Returns:
      an unmodifiable list of results, one per query statement, in declaration order
      Throws:
      IllegalArgumentException - if result contains semantic errors
      EvaluationException - if a runtime data-level error occurs
    • execute

      public List<QueryResult> execute(SemanticResult result)
      Executes all query statements in result using an inline-only connector and collects each result into a QueryResult.
      Parameters:
      result - a fully-valid semantic analysis result; must not be null
      Returns:
      an unmodifiable list of results, one per query statement
      Throws:
      IllegalArgumentException - if result contains semantic errors
      EvaluationException - if a runtime error occurs, including external-source access
    • execute

      public List<QueryResult> execute(SemanticModel model, DataSourceConnector connector)
      Executes all query statements in model, using the supplied connector, and collects each result into a QueryResult.
      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      Returns:
      an unmodifiable list of results in declaration order
      Throws:
      EvaluationException - if a runtime data-level error occurs
    • execute

      public List<QueryResult> execute(SemanticModel model)
      Executes all query statements in model using an inline-only connector and collects each result into a QueryResult.
      Parameters:
      model - the semantic model; must not be null
      Returns:
      an unmodifiable list of results in declaration order
      Throws:
      EvaluationException - if a runtime error occurs
    • executeProvenance

      public <K> void executeProvenance(SemanticModel model, DataSourceConnector connector, Semiring<K> semiring, ProvenanceConsumer<K> consumer)
      Evaluates each query statement in model as an annotated K-relation over semiring, handing each result to consumer in declaration order.

      Provenance is an explicit, opt-in mode: the chosen semiring threads through the positive-algebra operators (σ/π/×/⋈/∪) while every non-positive operator is read as an opaque base relation — see ProvenanceEvaluator. The query is evaluated from its raw logical tree (no optimizer, no pushdown): annotation tracking is in-engine, above the federation boundary. Because a K-relation is canonical, each result is fully materialised rather than streamed.

      Type Parameters:
      K - the semiring annotation type
      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      semiring - the annotation semiring; must not be null
      consumer - receives each query's label and annotated relation; must not be null
      Throws:
      EvaluationException - if a runtime data-level error occurs during evaluation
    • executeProvenance

      public <K> void executeProvenance(SemanticModel model, DataSourceConnector connector, Semiring<K> semiring, String weightColumn, int maxFixpointRounds, ProvenanceConsumer<K> consumer)
      Evaluates each query as a K-relation over semiring, supplying per-edge weights from weightColumn for any transitive closure in the query — the semiring-weighted-closure path.

      When weightColumn is non-null, each base edge is annotated with the value of that column coerced into the semiring (a double for tropical shortest path, a multiplicity for ℕ; the boolean/security semirings ignore the weight). A null/absent/non-numeric weight, and every non-closure base tuple, lift to one. When weightColumn is null this is the plain one() lift — boolean reachability and ℕ path-counting need no weight column.

      maxFixpointRounds bounds the weighted-closure iteration (see ExecutionContext.maxFixpointRounds()); it matters only for a non-idempotent semiring (ℕ) over a cyclic graph, which otherwise never converges. Use ExecutionContext.UNLIMITED_FIXPOINT_ROUNDS for no cap.

      Type Parameters:
      K - the semiring annotation type
      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      semiring - the annotation semiring; must not be null
      weightColumn - the per-edge weight column for weighted closures, or null
      maxFixpointRounds - the weighted-closure iteration cap (≥ 1)
      consumer - receives each query's label and annotated relation; must not be null
      Throws:
      EvaluationException - if a runtime data-level error occurs during evaluation
    • executeLineage

      public void executeLineage(SemanticModel model, DataSourceConnector connector, ProvenanceConsumer<Polynomial> consumer)
      Evaluates each query statement in model as a full why-provenance K-relation over the polynomial-lineage semiring ℕ[X], handing each result to consumer.

      Unlike the cheap semirings, lineage mints a distinct provenance variable per base-tuple occurrence (<relation>#<ordinal>), so each output tuple's annotation is the polynomial recording exactly which input tuples produced it and how they were combined. The representation is bounded: a polynomial that exceeds PolynomialSemiring.MAX_MONOMIALS terms is truncated and flagged. Like all provenance evaluation this runs on the raw logical tree, in-engine, never pushed down.

      Parameters:
      model - the semantic model; must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label and lineage-annotated relation; must not be null
      Throws:
      EvaluationException - if a runtime data-level error occurs during evaluation
    • executeOptimizedStreaming

      public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer)
      Streams the results of optimizer-rewritten query trees, one per root query in SemanticModel.rootQueries(), using the supplied connector.

      optimizedRoots must be parallel to model.rootQueries(): element i is the rewritten tree to execute in place of root query i. Because the optimizer produces fresh RelNode instances absent from SemanticModel.nodeSchemas() (which is keyed by object identity), this method re-infers schema annotations for each rewritten tree via SchemaInference.annotate(com.darkcollective.relix.symbol.table.SymbolTable, com.darkcollective.relix.ast.RelNode, com.darkcollective.relix.semantic.SchemaAnnotations, com.darkcollective.relix.function.FunctionCatalog), retaining the original annotations so that nested view bodies referenced by — but not rewritten within — the optimized tree remain executable.

      Parameters:
      model - the semantic model the trees were optimized from; must not be null
      optimizedRoots - rewritten trees, parallel to model.rootQueries(); must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      Throws:
      IllegalArgumentException - if optimizedRoots size differs from the root-query count
      EvaluationException - if a runtime data-level error occurs while a stream is read
    • executeOptimizedStreaming

      public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds)
      Like executeOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer) but aborts any FIX fixpoint that exceeds maxFixpointRounds iterations.
      Parameters:
      model - the semantic model; must not be null
      optimizedRoots - rewritten trees, parallel to model.rootQueries(); must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      maxFixpointRounds - maximum semi-naïve fixpoint iterations; use ExecutionContext.UNLIMITED_FIXPOINT_ROUNDS for no cap
      Throws:
      IllegalArgumentException - if optimizedRoots size differs from the root-query count
      EvaluationException - if a runtime data-level error occurs or a fixpoint cap is hit
    • executeOptimizedStreaming

      public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds, QueryEventListener listener)
      Like executeOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer, int) but reports the run's PLAN and EXECUTE decisions to listener, each tagged with the owning query's label — the observing counterpart for a caller that has already optimized. The OPTIMIZE events belong to that earlier pass and are collected there.
      Parameters:
      model - the semantic model the trees were optimized from; must not be null
      optimizedRoots - rewritten trees, parallel to model.rootQueries(); must not be null
      connector - provides rows for external data sources; must not be null
      consumer - receives each query's label, schema, and lazy row stream
      maxFixpointRounds - maximum semi-naïve fixpoint iterations; use ExecutionContext.UNLIMITED_FIXPOINT_ROUNDS for no cap
      listener - notified on each planning/execution decision; must not be null
      Throws:
      IllegalArgumentException - if optimizedRoots size differs from the root-query count
      EvaluationException - if a runtime data-level error occurs or a fixpoint cap is hit
    • executeOptimizedStreaming

      public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, RowStreamConsumer consumer)
      Streams the results of optimizer-rewritten query trees using an inline-only connector.
      Parameters:
      model - the semantic model the trees were optimized from; must not be null
      optimizedRoots - rewritten trees, parallel to model.rootQueries(); must not be null
      consumer - receives each query's label, schema, and lazy row stream
      Throws:
      IllegalArgumentException - if optimizedRoots size differs from the root-query count
      EvaluationException - if a runtime error occurs, including external-source access
    • executeOptimized

      public List<QueryResult> executeOptimized(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector)
      Executes optimizer-rewritten query trees, using the supplied connector, and collects each result into a QueryResult. See executeOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer).
      Parameters:
      model - the semantic model the trees were optimized from; must not be null
      optimizedRoots - rewritten trees, parallel to model.rootQueries(); must not be null
      connector - provides rows for external data sources; must not be null
      Returns:
      an unmodifiable list of results in declaration order
      Throws:
      IllegalArgumentException - if optimizedRoots size differs from the root-query count
      EvaluationException - if a runtime data-level error occurs
    • executeOptimized

      public List<QueryResult> executeOptimized(SemanticModel model, List<RelNode> optimizedRoots)
      Executes optimizer-rewritten query trees using an inline-only connector and collects each result into a QueryResult.
      Parameters:
      model - the semantic model the trees were optimized from; must not be null
      optimizedRoots - rewritten trees, parallel to model.rootQueries(); must not be null
      Returns:
      an unmodifiable list of results in declaration order
      Throws:
      IllegalArgumentException - if optimizedRoots size differs from the root-query count
      EvaluationException - if a runtime error occurs, including external-source access
    • explain

      public void explain(SemanticModel model, PlanConsumer consumer)
      Renders the physical plan of each query statement in model and hands it to consumer, without executing anything. This shows the planner's physical decisions — notably which sub-trees were pushed to the database as a PhysicalNode.PushedScan and each join's algorithm and build side.
      Parameters:
      model - the semantic model; must not be null
      consumer - receives each query's label and rendered plan; must not be null
    • explain

      public void explain(SemanticModel model, List<RelNode> roots, PlanConsumer consumer)
      Renders the physical plan of each (optimizer-rewritten) root tree and hands it to consumer, without executing anything. Use this overload after optimization so the explained plan matches what executeOptimized* would run; pass model.rootQueries()'s resolved nodes (via explain(SemanticModel, PlanConsumer)) for the un-optimized plan.
      Parameters:
      model - the semantic model the trees belong to; must not be null
      roots - the trees to plan, parallel to model.rootQueries(); must not be null
      consumer - receives each query's label and rendered plan; must not be null
      Throws:
      IllegalArgumentException - if roots size differs from the root-query count
    • trace

      public void trace(SemanticModel model, List<RelNode> roots, QueryEventListener listener)
      Plans each (optimizer-rewritten) root tree — without executing it — so the planner emits its physical-decision QueryEvents to listener. Each event is tagged with the owning query's label.

      Pair this with the optimizer's listener-aware optimize overload on the same listener to get one feed of both the optimizer's rule firings and the planner's decisions.

      Parameters:
      model - the semantic model the trees belong to; must not be null
      roots - the trees to plan, parallel to model.rootQueries(); must not be null
      listener - notified on each physical decision; must not be null
      Throws:
      IllegalArgumentException - if roots size differs from the root-query count
    • traceExecute

      public void traceExecute(SemanticModel model, List<RelNode> roots, DataSourceConnector connector, QueryEventListener listener)
      Executes each (optimizer-rewritten) root tree purely for its QueryEvent.Stage.EXECUTE side-effect events — discarding the rows — emitting each to listener tagged with the owning query's label.

      This is the execution-stage companion to trace(com.darkcollective.relix.semantic.SemanticModel, java.util.List<com.darkcollective.relix.ast.RelNode>, com.darkcollective.relix.events.QueryEventListener): where trace surfaces the planner's decisions without running anything, this surfaces decisions made during execution (e.g. a declarative- optimisation group skipped as infeasible). Planning inside execution uses QueryEventListener.NONE, so no PLAN events are double-emitted; only the EXECUTE feed flows here.

      Parameters:
      model - the semantic model the trees belong to; must not be null
      roots - the trees to run, parallel to model.rootQueries(); must not be null
      connector - provides rows for external data sources; must not be null
      listener - notified on each execution-stage decision; must not be null
      Throws:
      IllegalArgumentException - if roots size differs from the root-query count
      EvaluationException - if a runtime data-level error occurs while a stream is read
    • plan

      public void plan(SemanticModel model, List<RelNode> roots, QueryEventListener listener, PhysicalPlanConsumer consumer)
      Plans each (optimizer-rewritten) root tree — without executing it — and hands the resulting PhysicalNode to consumer, while emitting the planner's physical-decision QueryEvents to listener (each tagged with the owning query's label).

      This is the structured counterpart to explain(SemanticModel, java.util.List, PlanConsumer) (which renders text): use it when the caller needs the plan object — for example to serialize it to JSON for the query bundle. Pass QueryEventListener.NONE when the events are not wanted.

      Parameters:
      model - the semantic model the trees belong to; must not be null
      roots - the trees to plan, parallel to model.rootQueries(); must not be null
      listener - notified on each physical decision; must not be null
      consumer - receives each query's label and planned physical tree; must not be null
      Throws:
      IllegalArgumentException - if roots size differs from the root-query count
    • requireBounded

      public static void requireBounded(List<RelNode> roots, BoundednessSource boundedness)
      Refuses to collect a relation that provably never ends.

      BoundednessChecker runs in the planner and catches a blocking operator over an unbounded input. It cannot catch this one, because here the blocking is done by the consumer: σ (Naturals) contains no blocking node at all, so the plan is legal and it is collect(java.util.function.Consumer<com.darkcollective.relix.processor.internal.RowStreamConsumer>, int) that never returns. Streaming callers are right to be allowed it — emitting an endless relation is what a generator is for — so the guard belongs on the collecting entry points alone, which is why it lives here rather than in the planner.

      The remedy named is the planner's own, so a user who has met either message has met both.

      Package-private rather than private so it can be tested against a real GeneratorBoundednessSource without standing up a whole analysed script; the streaming entry points deliberately never call it.

      Parameters:
      roots - the trees about to be materialised
      boundedness - the per-leaf boundedness for this run
      Throws:
      BoundednessException - if any root is provably unbounded