From 46937e143a87f4647134e4ac30a2e4301b98e23d Mon Sep 17 00:00:00 2001 From: Alejandro Alonso Date: Thu, 16 Jul 2026 08:56:03 +0200 Subject: [PATCH] WIP arrow api --- backend/deps.edn | 7 ++- backend/dev/graph_arrow_spike.clj | 89 +++++++++++++++++++++++++++++++ backend/scripts/_env | 3 +- backend/scripts/run.template.sh | 2 +- 4 files changed, 97 insertions(+), 4 deletions(-) create mode 100644 backend/dev/graph_arrow_spike.clj diff --git a/backend/deps.edn b/backend/deps.edn index 9bdf318dd7..d43c81c954 100644 --- a/backend/deps.edn +++ b/backend/deps.edn @@ -66,13 +66,16 @@ software.amazon.awssdk/s3 {:mvn/version "2.46.18"} software.amazon.awssdk/sts {:mvn/version "2.46.18"} - com.ladybugdb/lbug {:mvn/version "0.18.0"}} + com.ladybugdb/lbug {:mvn/version "0.18.0"} + ;; Required by Arrow RootAllocator (lbug only pulls arrow-memory-core). + org.apache.arrow/arrow-memory-netty {:mvn/version "18.2.0"}} :paths ["src" "resources" "target/classes"] :aliases {:dev {:jvm-opts ["--sun-misc-unsafe-memory-access=allow" - "--enable-native-access=ALL-UNNAMED"] + "--enable-native-access=ALL-UNNAMED" + "--add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED"] :extra-deps {com.bhauman/rebel-readline {:mvn/version "0.1.11"} clojure-humanize/clojure-humanize {:mvn/version "0.2.2"} diff --git a/backend/dev/graph_arrow_spike.clj b/backend/dev/graph_arrow_spike.clj new file mode 100644 index 0000000000..e374c0b04b --- /dev/null +++ b/backend/dev/graph_arrow_spike.clj @@ -0,0 +1,89 @@ +;; Spike: Ladybug createArrowTable → native COPY (no CSV). +;; Run inside devenv from backend/: +;; clojure -M:dev -m graph-arrow-spike +;; +;; Requires --add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED +;; (configured in deps.edn :dev :jvm-opts). + +(ns graph-arrow-spike + (:gen-class) + (:import + (com.ladybugdb Connection Database QueryResult) + (java.nio.charset StandardCharsets) + (java.util ArrayList List) + (org.apache.arrow.memory RootAllocator) + (org.apache.arrow.vector VarCharVector VectorSchemaRoot) + (org.apache.arrow.vector.types.pojo ArrowType$Utf8 Field FieldType Schema))) + +(defn- check! + [^QueryResult result label] + (when-not (.isSuccess result) + (throw (ex-info (str label ": " (.getErrorMessage result)) + {:label label + :err (.getErrorMessage result)}))) + result) + +(defn- page-root + [^RootAllocator alloc] + (let [varchar-type (org.apache.arrow.vector.types.pojo.ArrowType$Utf8.) + field-id (Field. "id" (FieldType/nullable varchar-type) nil) + field-name (Field. "name" (FieldType/nullable varchar-type) nil) + schema (Schema. [field-id field-name]) + root (VectorSchemaRoot/create schema alloc) + ^VarCharVector idv (.getVector root "id") + ^VarCharVector nv (.getVector root "name")] + (.allocateNew idv 2) + (.allocateNew nv 2) + (.setSafe idv 0 (.getBytes "p1" StandardCharsets/UTF_8)) + (.setSafe idv 1 (.getBytes "p2" StandardCharsets/UTF_8)) + (.setSafe nv 0 (.getBytes "Home" StandardCharsets/UTF_8)) + (.setSafe nv 1 (.getBytes "About" StandardCharsets/UTF_8)) + (.setValueCount idv 2) + (.setValueCount nv 2) + (.setRowCount root 2) + root)) + +(defn- try-query! + [^Connection conn cypher label] + (with-open [^QueryResult r (.query conn cypher)] + (println label + "success?" (.isSuccess r) + "tuples" (when (.isSuccess r) (.getNumTuples r)) + "err" (when-not (.isSuccess r) (.getErrorMessage r))) + (.isSuccess r))) + +(defn -main + [& _] + (with-open [^RootAllocator alloc (RootAllocator.) + ^Database db (Database.) + ^Connection conn (Connection. db)] + (.setQueryTimeout conn 0) + (println "=== 1) native DDL ===") + (with-open [r (.query conn "CREATE NODE TABLE Page(id STRING, name STRING, PRIMARY KEY(id));")] + (check! r "ddl")) + + (println "=== 2) createArrowTable staging ===") + (let [root2 (page-root alloc) + batches (doto (ArrayList.) (.add root2))] + (with-open [r (.createArrowTable conn "_stg_Page" ^List batches alloc)] + (check! r "createArrowTable")) + + (println "=== 3) query staging ===") + (try-query! conn "MATCH (n:_stg_Page) RETURN n.id, n.name;" "stg") + + (println "=== 4) COPY Page FROM _stg_Page ===") + (when-not (try-query! conn "COPY Page FROM _stg_Page;" "copy-ident") + (println "=== 4b) COPY via MATCH subquery ===") + (try-query! conn + "COPY Page FROM (MATCH (n:_stg_Page) RETURN n.id AS id, n.name AS name);" + "copy-subq")) + + (println "=== 5) query native Page ===") + (try-query! conn "MATCH (n:Page) RETURN n.id, n.name;" "page") + + (println "=== 6) dropArrowTable ===") + (with-open [r (.dropArrowTable conn "_stg_Page")] + (println "drop success?" (.isSuccess r) "err" (.getErrorMessage r))) + + (.close root2) + (println "DONE")))) diff --git a/backend/scripts/_env b/backend/scripts/_env index 58ea02ede0..33ab6f0f15 100644 --- a/backend/scripts/_env +++ b/backend/scripts/_env @@ -84,7 +84,8 @@ export JAVA_OPTS="\ -XX:-OmitStackTraceInFastThrow \ --sun-misc-unsafe-memory-access=allow \ --enable-preview \ - --enable-native-access=ALL-UNNAMED"; + --enable-native-access=ALL-UNNAMED \ + --add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED"; function setup_minio() { if [ "${PENPOT_OBJECTS_STORAGE_BACKEND}" != "s3" ]; then diff --git a/backend/scripts/run.template.sh b/backend/scripts/run.template.sh index cff4afc870..143f14146d 100644 --- a/backend/scripts/run.template.sh +++ b/backend/scripts/run.template.sh @@ -18,7 +18,7 @@ if [ -f ./environ ]; then source ./environ fi -export JAVA_OPTS="-Djava.util.logging.manager=org.apache.logging.log4j.jul.LogManager -Dlog4j2.configurationFile=log4j2.xml -XX:-OmitStackTraceInFastThrow --sun-misc-unsafe-memory-access=allow --enable-native-access=ALL-UNNAMED --enable-preview $JVM_OPTS $JAVA_OPTS" +export JAVA_OPTS="-Djava.util.logging.manager=org.apache.logging.log4j.jul.LogManager -Dlog4j2.configurationFile=log4j2.xml -XX:-OmitStackTraceInFastThrow --sun-misc-unsafe-memory-access=allow --enable-native-access=ALL-UNNAMED --add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED --enable-preview $JVM_OPTS $JAVA_OPTS" ENTRYPOINT=${1:-app.main};