perf: Phase 2 — DuckDB SQL fallback for analytical queries
ober
869eac1efb741fc5c820bc259803f79e4b38292e
--- a/benchmarks/mbrainz-bench.ss +++ b/benchmarks/mbrainz-bench.ss @@ -9,7 +9,11 @@ ;;; scheme ... --script benchmarks/mbrainz-bench.ss [--quick] [--scale 0.1] (import (jerboa prelude) - (jerboa-db core)) + (jerboa-db core) + (jerboa-db history) + (jerboa-db query engine) + (jerboa-db query sql-translate) + (jerboa-db analytics)) ;; ---- Config ---- @@ -279,22 +283,32 @@ ;; Q8: Complex aggregation — avg track duration by release status ;; Tests: aggregation + join + grouping -(def (q8 db) - (q '((find ?status (count ?t) (avg ?dur)) - (where (?t track/duration ?dur) - (?t track/artists ?a) - (?r release/artists ?a) - (?r release/status ?status))) - db)) +(def q8-query + '((find ?status (count ?t) (avg ?dur)) + (where (?t track/duration ?dur) + (?t track/artists ?a) + (?r release/artists ?a) + (?r release/status ?status)))) + +(def (q8 db) (q q8-query db)) + +(def (q8-analytics db ae) + (analytical-query (parse-query q8-query) ae (db-value-schema db))) ;; ---- Benchmark runner ---- (def (run-benchmark conn sample-artist-eid sample-artist-name) - (let ([db (db conn)]) - - (displayln "") - (displayln "=== Jerboa-DB MBrainz Benchmark ===") - (displayln "") + (let* ([db (db conn)] + [_ (display "Setting up analytics engine... ")] + [t-ae0 (now-ms)] + [ae (new-analytics-engine (db-value-schema db))] + [_ (analytics-sync! ae db)] + [_ (displayln "done in " + (inexact->exact (round (- (now-ms) t-ae0))) + " ms")] + [_ (displayln "")] + [_ (displayln "=== Jerboa-DB MBrainz Benchmark ===")] + [_ (displayln "")]) (displayln (pad-right "Query" 55) (pad-left "Median" 10) (pad-left "Result" 15)) (displayln (make-string 80 #\-)) @@ -327,12 +341,15 @@ (bench "Q7" "Pull artist attributes" (lambda () (q7 db sample-artist-eid))) (bench "Q8" "Avg track duration by release status" - (lambda () (q8 db))))]) + (lambda () (q8 db))) + (bench "Q8*" " via DuckDB fallback" + (lambda () (q8-analytics db ae))))]) (displayln (make-string 80 #\-)) (let ([total (apply + (map cdr results))]) (displayln (pad-right "Total" 55) (fmt-ms total))) (displayln "") + (analytics-close ae) results)))) ;; ---- Main ---- --- a/lib/jerboa-db/query/engine.ss +++ b/lib/jerboa-db/query/engine.ss @@ -6,7 +6,9 @@ ;;; function clauses, aggregation, and rules. (library (jerboa-db query engine) - (export query-db parse-query explain-query) + (export query-db parse-query explain-query + parsed-query? parsed-query-find-vars parsed-query-in-vars + parsed-query-where-clauses parsed-query-rules) (import (except (chezscheme) make-hash-table hash-table? new file mode 100644 --- /dev/null +++ b/lib/jerboa-db/query/sql-translate.ss @@ -0,0 +1,351 @@ +#!chezscheme +;;; (jerboa-db query sql-translate) — Datalog → SQL translator for OLAP fallback. +;;; +;;; Translates a parsed Datalog query into SQL against the DuckDB analytics +;;; replica's `datoms` table. Used for aggregation-heavy queries (Q4, Q8 in +;;; the MBrainz benchmark) where DuckDB's columnar engine massively +;;; outperforms the row-at-a-time Datalog evaluator. +;;; +;;; Coverage: +;;; - find vars: logic-vars (grouping) + streamable aggregates +;;; - where: data patterns (?e attr ?v); shared vars become joins +;;; - predicates: < <= > >= = != on bound vars + literals +;;; +;;; Falls back (returns #f) on rules, or, not, function clauses, fancy inputs. + +(library (jerboa-db query sql-translate) + (export translate-query-to-sql + analytical-query) + + (import (except (chezscheme) + make-hash-table hash-table? + sort sort! + printf fprintf + path-extension path-absolute? + with-input-from-string with-output-to-string + iota 1+ 1- + partition + make-date make-time + atom? meta) + (jerboa prelude) + (jerboa-db schema) + (jerboa-db analytics) + (jerboa-db query engine) + (jerboa-db query aggregates)) + + (def (logic-var? x) + (and (symbol? x) + (> (string-length (symbol->string x)) 0) + (char=? (string-ref (symbol->string x) 0) #\?))) + + (def (data-pattern? clause) + (and (pair? clause) + (not (pair? (car clause))) + (>= (length clause) 3) + (symbol? (cadr clause)) + (not (logic-var? (cadr clause))))) + + (def (predicate-clause? clause) + (and (pair? clause) (pair? (car clause)) (null? (cdr clause)))) + + (def (value-column-for-type vtype) + (case vtype + [(db.type/long db.type/ref db.type/instant) + (case vtype + [(db.type/long) "v_long"] + [(db.type/ref) "v_ref"] + [(db.type/instant) "v_instant"])] + [(db.type/double) "v_double"] + [(db.type/string db.type/keyword db.type/uuid db.type/symbol) "v_string"] + [(db.type/boolean) "v_bool"] + [else #f])) + + (def (sql-quote-string s) + (let ([buf (open-output-string)]) + (display "'" buf) + (string-for-each + (lambda (c) + (when (char=? c #\') (display "'" buf)) + (display c buf)) + s) + (display "'" buf) + (get-output-string buf))) + + (def (sql-literal v) + (cond + [(string? v) (sql-quote-string v)] + [(symbol? v) (sql-quote-string (symbol->string v))] + [(boolean? v) (if v "TRUE" "FALSE")] + [(number? v) (number->string v)] + [else (error 'sql-literal "Unsupported literal" v)])) + + (def (sql-aggregate-expr name col-expr) + (case name + [(count) (string-append "COUNT(" col-expr ")")] + [(sum) (string-append "SUM(" col-expr ")")] + [(avg) (string-append "AVG(" col-expr ")")] + [(min) (string-append "MIN(" col-expr ")")] + [(max) (string-append "MAX(" col-expr ")")] + [else #f])) + + (def (translate-query-to-sql parsed-q schema) + ;; Returns (sql . slots) or #f + ;; slots: list of (kind . src) per find-var + ;; (agg . agg-name) — column name `agg_N` + ;; (grp . column-alias) — column name from GROUP BY + ;; (lit . value) — literal echoed in output + (let ([find-vars (parsed-query-find-vars parsed-q)] + [in-vars (parsed-query-in-vars parsed-q)] + [where-cl (parsed-query-where-clauses parsed-q)]) + (and (equal? in-vars '($)) + (pair? where-cl) + (for-all (lambda (cl) + (or (data-pattern? cl) + (predicate-clause? cl))) + where-cl) + (let ([data-clauses (filter data-pattern? where-cl)] + [pred-clauses (filter predicate-clause? where-cl)]) + (and (pair? data-clauses) + (build-sql find-vars data-clauses pred-clauses schema)))))) + + (def (build-sql find-vars data-clauses pred-clauses schema) + ;; First pass: assign aliases d1..dN; for each clause record the + ;; column expression for its e-spec and v-spec. Then walk find-vars + ;; to build SELECT, GROUP BY, and the slot mapping for results. + (let loop-assign + ([clauses data-clauses] + [n 1] + [var-cols '()] ;; alist: var -> "dN.col" + [where-parts '()] ;; list of strings + [from-parts '()]) ;; list of strings + (cond + [(null? clauses) + (finalize-sql find-vars (reverse var-cols) + (reverse where-parts) (reverse from-parts) + pred-clauses schema)] + [else + (let* ([cl (car clauses)] + [e-spec (car cl)] [a-spec (cadr cl)] [v-spec (caddr cl)] + [attr (schema-lookup-by-ident schema a-spec)] + [alias (string-append "d" (number->string n))]) + (and attr + (let* ([aid (db-attribute-id attr)] + [vtype (db-attribute-value-type attr)] + [vcol (value-column-for-type vtype)] + [e-expr (string-append alias ".e")] + [v-expr (and vcol (string-append alias "." vcol))]) + (and v-expr + (let* ([from-piece (string-append "datoms " alias)] + [where-pieces + (cons + (string-append alias ".a = " + (number->string aid)) + (cons + (string-append alias ".added") + '()))] + ;; Extra equalities for repeated vars/literals + [extra-where '()] + [vc1 (cond + [(logic-var? e-spec) + (let ([prev (assq e-spec var-cols)]) + (if prev + (begin + (set! extra-where + (cons (string-append e-expr " = " + (cdr prev)) + extra-where)) + var-cols) + (cons (cons e-spec e-expr) + var-cols)))] + [else + (set! extra-where + (cons (string-append e-expr " = " + (sql-literal e-spec)) + extra-where)) + var-cols])] + [vc2 (cond + [(logic-var? v-spec) + (let ([prev (assq v-spec vc1)]) + (if prev + (begin + (set! extra-where + (cons (string-append v-expr " = " + (cdr prev)) + extra-where)) + vc1) + (cons (cons v-spec v-expr) + vc1)))] + [else + (set! extra-where + (cons (string-append v-expr " = " + (sql-literal v-spec)) + extra-where)) + vc1])]) + (loop-assign (cdr clauses) (+ n 1) + vc2 + (append extra-where + where-pieces + where-parts) + (cons from-piece from-parts)))))))]))) + + (def (finalize-sql find-vars var-cols where-parts from-parts pred-clauses schema) + ;; Build SELECT projection + GROUP BY + WHERE for predicates + (let* ([find-info (build-find-info find-vars var-cols)]) + (and find-info + (let* ([pred-where (build-pred-where pred-clauses var-cols)]) + (and pred-where + (build-final-sql find-info from-parts + (append pred-where where-parts))))))) + + ;; find-info: list of slot descriptors (matched to find-vars order): + ;; (grp . sql-expr) + ;; (agg . (agg-name col-alias)) + ;; (lit . value) + (def (build-find-info find-vars var-cols) + (let loop ([fvs find-vars] [acc '()]) + (cond + [(null? fvs) (reverse acc)] + [(logic-var? (car fvs)) + (let ([col (assq (car fvs) var-cols)]) + (and col + (loop (cdr fvs) + (cons (cons 'grp (cdr col)) acc))))] + [(and (pair? (car fvs)) + (memq (caar fvs) '(count sum avg min max)) + (logic-var? (cadar fvs))) + (let ([col (assq (cadar fvs) var-cols)]) + (and col + (let ([alias (string-append "agg_" + (number->string (length acc)))]) + (loop (cdr fvs) + (cons (list 'agg (caar fvs) (cdr col) alias) + acc)))))] + [(not (logic-var? (car fvs))) + (loop (cdr fvs) (cons (cons 'lit (car fvs)) acc))] + [else #f]))) + + (def (build-pred-where pred-clauses var-cols) + (let loop ([cs pred-clauses] [acc '()]) + (cond + [(null? cs) (reverse acc)] + [else + (let* ([form (car (car cs))] + [op (car form)] + [args (cdr form)]) + (and (memq op '(< <= > >= = != not=)) + (= (length args) 2) + (let ([sql-args + (map (lambda (a) + (cond + [(logic-var? a) + (let ([col (assq a var-cols)]) + (and col (cdr col)))] + [else (sql-literal a)])) + args)]) + (and (for-all (lambda (x) x) sql-args) + (let ([sql-op (case op + [(=) "="] + [(!= not=) "<>"] + [else (symbol->string op)])]) + (loop (cdr cs) + (cons (string-append (car sql-args) " " + sql-op " " + (cadr sql-args)) + acc)))))))]))) + + (def (build-final-sql find-info from-parts where-parts) + (let* ([select-pieces + (let loop ([fi find-info] [acc '()]) + (cond + [(null? fi) (reverse acc)] + [else + (let ([slot (car fi)]) + (case (car slot) + [(grp) + (loop (cdr fi) (cons (cdr slot) acc))] + [(agg) + (let* ([name (cadr slot)] + [col (caddr slot)] + [alias (cadddr slot)] + [expr (sql-aggregate-expr + name + (if (eq? name 'count) "*" col))]) + (loop (cdr fi) + (cons (string-append expr " AS " alias) acc)))] + [(lit) + (loop (cdr fi) + (cons (sql-literal (cdr slot)) acc))]))]))] + [has-aggs? + (exists (lambda (s) (eq? (car s) 'agg)) find-info)] + [group-pieces + (if has-aggs? + (let loop ([fi find-info] [acc '()]) + (cond + [(null? fi) (reverse acc)] + [(eq? (caar fi) 'grp) + (loop (cdr fi) (cons (cdar fi) acc))] + [else (loop (cdr fi) acc)])) + '())] + [select-clause (string-append "SELECT " + (sql-join select-pieces ", "))] + [from-clause (string-append "FROM " + (sql-join from-parts ", "))] + [where-clause (if (null? where-parts) "" + (string-append " WHERE " + (sql-join where-parts " AND ")))] + [group-clause (if (null? group-pieces) "" + (string-append " GROUP BY " + (sql-join group-pieces ", ")))] + [sql (string-append select-clause " " + from-clause + where-clause + group-clause)]) + (cons sql find-info))) + + (def (sql-join lst sep) + (cond + [(null? lst) ""] + [(null? (cdr lst)) (car lst)] + [else (string-append (car lst) sep (sql-join (cdr lst) sep))])) + + ;; Run a translated query and shape rows to match Datalog tuple form. + (def (analytical-query parsed-q ae schema) + (let ([translated (translate-query-to-sql parsed-q schema)]) + (and translated + (let* ([sql (car translated)] + [find-info (cdr translated)] + [rows (analytics-query ae sql)]) + (map (lambda (row) + (shape-result-row find-info row)) + rows))))) + + (def (shape-result-row find-info row) + ;; row is alist: ((col-name . val) ...) + ;; For each find-info slot, pull the matching column. + (let loop ([fi find-info] [grp-i 0] [acc '()]) + (cond + [(null? fi) (reverse acc)] + [else + (let ([slot (car fi)]) + (case (car slot) + [(grp) + ;; Pull by index: count of grp-slots before this one → alist position. + ;; Simpler: pull by raw column name (from cdr slot, e.g. "d1.e"). + ;; DuckDB strips alias prefix in column-name; column name is "e", + ;; "v_string", etc. We need a stable shaping, so iterate by row order. + (loop (cdr fi) (+ grp-i 1) + (cons (cdr-of-nth-pair row grp-i) acc))] + [(agg) + (let ([alias (cadddr slot)]) + (loop (cdr fi) grp-i + (cons (cdr (assoc alias row)) acc)))] + [(lit) + (loop (cdr fi) grp-i (cons (cdr slot) acc))]))]))) + + (def (cdr-of-nth-pair lst n) + (cond + [(null? lst) #f] + [(= n 0) (cdar lst)] + [else (cdr-of-nth-pair (cdr lst) (- n 1))])) + +) ;; end library new file mode 100644 --- /dev/null +++ b/tests/test-sql-translate.ss @@ -0,0 +1,61 @@ +;;; Smoke test: end-to-end SQL fallback for an aggregation query. + +(import (chezscheme) + (jerboa-db core) + (jerboa-db schema) + (jerboa-db history) + (jerboa-db query engine) + (jerboa-db query sql-translate) + (jerboa-db analytics)) + +(define conn (connect ":memory:")) + +;; Schema +(transact! conn (list + `((db/ident . track/duration) (db/valueType . db.type/long) (db/cardinality . db.cardinality/one)) + `((db/ident . track/artists) (db/valueType . db.type/ref) (db/cardinality . db.cardinality/many)) + `((db/ident . release/artists)(db/valueType . db.type/ref) (db/cardinality . db.cardinality/many)) + `((db/ident . release/status) (db/valueType . db.type/string)(db/cardinality . db.cardinality/one)))) + +;; Synthetic data: 3 artists, 3 releases (each with one artist), 6 tracks +(define-values (a1 a2 a3) (values 100 101 102)) +(define-values (r1 r2 r3) (values 200 201 202)) +(define t-base 300) + +(transact! conn + `(((db/id . ,a1)) + ((db/id . ,a2)) + ((db/id . ,a3)) + ((db/id . ,r1) (release/artists . ,a1) (release/status . "Official")) + ((db/id . ,r2) (release/artists . ,a2) (release/status . "Official")) + ((db/id . ,r3) (release/artists . ,a3) (release/status . "Bootleg")) + ((db/id . ,(+ t-base 0)) (track/artists . ,a1) (track/duration . 100)) + ((db/id . ,(+ t-base 1)) (track/artists . ,a1) (track/duration . 200)) + ((db/id . ,(+ t-base 2)) (track/artists . ,a2) (track/duration . 300)) + ((db/id . ,(+ t-base 3)) (track/artists . ,a2) (track/duration . 400)) + ((db/id . ,(+ t-base 4)) (track/artists . ,a3) (track/duration . 500)) + ((db/id . ,(+ t-base 5)) (track/artists . ,a3) (track/duration . 600)))) + +(define query + '((find ?status (count ?t) (avg ?dur)) + (where (?t track/duration ?dur) + (?t track/artists ?a) + (?r release/artists ?a) + (?r release/status ?status)))) + +;; Datalog reference +(display "Datalog: ") +(write (q query (db conn))) +(newline) + +;; SQL-translated path +(define ae (new-analytics-engine (db-value-schema (db conn)))) +(analytics-sync! ae (db conn)) +(define p (parse-query query)) +(display "SQL: ") +(write (translate-query-to-sql p (db-value-schema (db conn)))) +(newline) +(display "Result: ") +(write (analytical-query p ae (db-value-schema (db conn)))) +(newline) +(analytics-close ae)