Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion .github/workflows/clojure-integration-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -83,12 +83,15 @@ jobs:
clojure -M:test -m cognitect.test-runner \
-n pgloader.cast-test \
-n pgloader.batch-test \
-n pgloader.prefetch-test \
-n pgloader.ddl-test \
-n pgloader.ddl.citus-test \
-n pgloader.load-file.parser-test \
-n pgloader.transforms-test \
-n pgloader.pg-service-test \
-n pgloader.cli-test
-n pgloader.cli-test \
-n pgloader.core-test \
-n pgloader.source.mssql-test

# ── Build test-runner Docker image (contains v3 pgloader binary) ───────────
# Built once and shared with all integration jobs via artifact.
Expand Down
5 changes: 4 additions & 1 deletion clojure/Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -48,12 +48,15 @@ test-unit:
clojure -M:test -m cognitect.test-runner \
-n pgloader.cast-test \
-n pgloader.batch-test \
-n pgloader.prefetch-test \
-n pgloader.ddl-test \
-n pgloader.ddl.citus-test \
-n pgloader.load-file.parser-test \
-n pgloader.transforms-test \
-n pgloader.pg-service-test \
-n pgloader.cli-test
-n pgloader.cli-test \
-n pgloader.core-test \
-n pgloader.source.mssql-test

# ─── E2E integration tests ────────────────────────────────────────────────────
# All suite management lives in tests/Makefile.
Expand Down
32 changes: 26 additions & 6 deletions clojure/src/pgloader/core.clj
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@
(:import [org.postgresql PGConnection]
[java.sql Connection DriverManager]
[java.io File]
[java.util.concurrent Executors ExecutorService Future TimeUnit])
[java.util.concurrent Executors ExecutorService Future TimeUnit
ExecutionException])
(:require [clojure.tools.logging :as log]))

(set! *warn-on-reflection* true)
Expand Down Expand Up @@ -408,6 +409,25 @@
:reset-sequences true}
with-options))

(defn- await-table-futures!
"Wait for every table worker, finish index work, then propagate a strict failure."
[table-futs ^ExecutorService idx-executor]
(let [first-failure (volatile! nil)]
(doseq [^Future f table-futs]
(try
(.get f)
(catch ExecutionException e
(when-not @first-failure
(vreset! first-failure (or (.getCause e) e))))
(catch Exception e
(when-not @first-failure
(vreset! first-failure e)))))
(when idx-executor
(.shutdown idx-executor)
(.awaitTermination idx-executor Long/MAX_VALUE TimeUnit/NANOSECONDS))
(when (and copy/*on-error-stop* @first-failure)
(throw @first-failure))))

(defn run-command
[cmd opts]
(if (= :archive (:load-type cmd))
Expand Down Expand Up @@ -786,7 +806,7 @@
(fn [[_i t]]
(.submit ^ExecutorService workers-pool
^java.util.concurrent.Callable
(fn []
(bound-fn []
(let [worker-src (source-from-uri source-uri table-spec
with-options source-overrides (:decoding-as cmd))
worker-pg (postgres-connection target-uri)
Expand All @@ -806,7 +826,8 @@
table (:table-name t)
cols (:columns t)
table-label (table-stats-label schema table)]
(when-not (@failed-tables table)
(when-not (or (@failed-tables table)
(and copy/*on-error-stop* @load-failed))
(let [;; Generated columns exist on the target (DDL emits
;; GENERATED ALWAYS AS) but must be excluded from COPY
;; since PostgreSQL cannot accept values for them.
Expand Down Expand Up @@ -835,7 +856,7 @@
(mapv (fn [part-src]
(.submit ^ExecutorService part-exec
^java.util.concurrent.Callable
(fn []
(bound-fn []
(let [part-pg (postgres-connection target-uri)]
(when pg-params
(doseq [param pg-params]
Expand Down Expand Up @@ -949,8 +970,7 @@
(map-indexed vector cat))]
(.shutdown ^ExecutorService workers-pool)
(.awaitTermination ^ExecutorService workers-pool Long/MAX_VALUE TimeUnit/NANOSECONDS)
(doseq [^Future f table-futs]
(try (.get f) (catch Exception _)))
(await-table-futures! table-futs idx-executor)
(stats/update-entry! :post "COPY Wall-Clock Time"
:rows workers
:bytes (:bytes (stats/get-totals :data))
Expand Down
39 changes: 25 additions & 14 deletions clojure/src/pgloader/prefetch.clj
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
[org.postgresql.util PSQLException]
[java.sql Connection]
[java.util.concurrent LinkedBlockingQueue
BlockingQueue]
BlockingQueue TimeUnit]
[java.util.concurrent.atomic AtomicBoolean AtomicLong]
[java.nio.charset StandardCharsets])
(:require [pgloader.batch :as batch]
Expand Down Expand Up @@ -75,7 +75,7 @@
(defn- send-batch-or-retry!
"Send a single batch, or handle errors and return updated counters.
Returns {:status :ok :rows-ok ... :errors ... :ws-nanos ... :bytes ... :reject-paths ...}
for success/retry, or throws for non-retryable errors."
for success/retry. In strict mode, rolls back and propagates COPY errors."
[^PGConnection pg-conn table-spec ^String copy-sql-str
b rows-ok errors ws-nanos bytes reject-paths]
(let [batch-start (System/nanoTime)]
Expand All @@ -95,7 +95,13 @@
(if cause (.getMessage ^Throwable cause) "unknown"))))
(throw e))
(catch PSQLException e
(.rollback ^Connection pg-conn)
(try
(.rollback ^Connection pg-conn)
(catch Exception rollback-error
(.addSuppressed e rollback-error)
(throw e)))
(when copy/*on-error-stop*
(throw e))
(log/info "Entering error recovery.")
(let [retry-result (batch/retry-batch! b table-spec e pg-conn)]
{:status :retry
Expand All @@ -112,8 +118,8 @@
(defn writer-task
"Virtual thread task that drains batches from the pipeline queue
and sends them to PostgreSQL via CopyManager.
Each batch gets its own transaction. On data errors, retry-batch!
handles per-row recovery with independent sub-batch commits.
Each batch gets its own transaction. In resume mode, retry-batch! handles
per-row recovery with independent sub-batch commits.
Returns {:rows-ok n :rows-bad n :ws-nanos n :bytes n :reject-paths {...}}."
[^PGConnection pg-conn table-spec ^CopyPipeline pipeline]
(let [copy-sql-str (copy/copy-sql table-spec)
Expand All @@ -123,20 +129,25 @@
ws-nanos (long 0)
bytes (long 0)
reject-paths nil]
(let [item (.take ^BlockingQueue (.queue pipeline))]
(if (= :end-of-data item)
(let [item (.poll ^BlockingQueue (.queue pipeline)
100 TimeUnit/MILLISECONDS)]
(if (or (= :end-of-data item)
(and (nil? item)
(.get ^AtomicBoolean (.done pipeline))))
{:rows-ok rows-ok
:rows-bad errors
:ws-nanos (- (System/nanoTime) start)
:bytes bytes
:reject-paths reject-paths}
(let [^batch/Batch b item
result (send-batch-or-retry!
pg-conn table-spec copy-sql-str
b rows-ok errors ws-nanos bytes reject-paths)]
(recur (long (:rows-ok result)) (long (:errors result))
(long (:ws-nanos result)) (long (:bytes result))
(:reject-paths result))))))))
(if (nil? item)
(recur rows-ok errors ws-nanos bytes reject-paths)
(let [^batch/Batch b item
result (send-batch-or-retry!
pg-conn table-spec copy-sql-str
b rows-ok errors ws-nanos bytes reject-paths)]
(recur (long (:rows-ok result)) (long (:errors result))
(long (:ws-nanos result)) (long (:bytes result))
(:reject-paths result)))))))))

(defn copy-table!
"Orchestrate the full copy of a single table.
Expand Down
15 changes: 3 additions & 12 deletions clojure/src/pgloader/source/mssql.clj
Original file line number Diff line number Diff line change
Expand Up @@ -215,24 +215,15 @@
meta (.getMetaData rs)
n (.getColumnCount meta)]
((fn thisfn []
(when (try (.next rs)
(catch Exception e
(log/warn (str "MSSQL row advance error in " table-name ": " (.getMessage e)))
false))
(when (.next rs)
(lazy-seq
(cons (loop [i 1 result (transient [])]
(if (<= i n)
(recur (inc i)
(conj! result
;; getString preserves full decimal/numeric precision (#1615, #1619).
;; Per-column error recovery: substitute nil on any error.
(try
(.getString rs i)
(catch Exception e
(log/warn (str "MSSQL column " i " read error in "
table-name " (substituting NULL): "
(.getMessage e)))
nil))))
;; Driver errors must abort rather than change data to NULL.
(.getString rs i)))
(persistent! result)))
(thisfn)))))))
(catch Exception e
Expand Down
19 changes: 18 additions & 1 deletion clojure/test/pgloader/cli_test.clj
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
(ns pgloader.cli-test
(:require [clojure.test :refer [deftest is testing]]
[pgloader.cli :as cli]))
[pgloader.cli :as cli]
[pgloader.core :as core]
[pgloader.load-file.parser :as parser]))

(deftest test-parse-args-basic
(testing "positional .load file"
Expand Down Expand Up @@ -81,3 +83,18 @@
(is (= ["quote identifiers" "include drop"] (:with-opts opts)))
(is (= "sqlite:///tmp/test.db" (:source-uri opts)))
(is (= "pgsql:///target" (:target-uri opts))))))

(deftest strict-load-file-failure-stops-later-files
(let [calls (atom [])
failure (ex-info "copy failed" {:file "first.load"})]
(with-redefs [parser/parse-file (fn [file] {:ok {:file file}})
core/run-command (fn [cmd _opts]
(swap! calls conj (:file cmd))
(when (= "first.load" (:file cmd))
(throw failure)))]
(let [thrown (try
(cli/run ["first.load" "second.load"])
nil
(catch Exception e e))]
(is (identical? failure thrown))
(is (= ["first.load"] @calls))))))
24 changes: 24 additions & 0 deletions clojure/test/pgloader/core_test.clj
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
(ns pgloader.core-test
(:require [clojure.test :refer [deftest is testing]]
[pgloader.copy :as copy]
[pgloader.core :as core])
(:import [java.util.concurrent CompletableFuture Executors ExecutorService]))

(deftest await-table-futures-propagates-the-first-strict-failure
(let [failure (ex-info "copy failed" {:table "items"})
failed (CompletableFuture/failedFuture failure)
complete (CompletableFuture/completedFuture :ok)]
(testing "strict mode unwraps and propagates a worker failure"
(let [^ExecutorService index-executor (Executors/newSingleThreadExecutor)]
(.submit index-executor ^Runnable (fn [] nil))
(binding [copy/*on-error-stop* true]
(let [thrown (try
(#'core/await-table-futures! [failed complete] index-executor)
nil
(catch Exception e e))]
(is (identical? failure thrown))
(is (.isShutdown index-executor))
(is (.isTerminated index-executor))))))
(testing "resume mode still waits for failures without aborting the load"
(binding [copy/*on-error-stop* false]
(is (nil? (#'core/await-table-futures! [failed complete] nil)))))))
124 changes: 124 additions & 0 deletions clojure/test/pgloader/prefetch_test.clj
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
(ns pgloader.prefetch-test
(:require [clojure.test :refer [deftest is testing]]
[pgloader.batch :as batch]
[pgloader.copy :as copy]
[pgloader.prefetch :as prefetch])
(:import [java.lang.reflect InvocationHandler Proxy]
[java.nio.charset StandardCharsets]
[java.sql Connection]
[java.util.concurrent.atomic AtomicBoolean]
[org.postgresql PGConnection]
[org.postgresql.util PSQLException PSQLState]))

(defn- pg-connection
([rolled-back?] (pg-connection rolled-back? nil))
([rolled-back? rollback-error]
(Proxy/newProxyInstance
(.getClassLoader PGConnection)
(into-array Class [PGConnection Connection])
(reify InvocationHandler
(invoke [_ _ method _]
(when (= "rollback" (.getName method))
(reset! rolled-back? true)
(when rollback-error
(throw rollback-error))))))))

(deftest writer-finishes-when-reader-fails-with-a-full-queue
(binding [copy/*prefetch-queue-capacity* 1]
(let [pipeline (prefetch/make-pipeline 10 1024)
test-batch (batch/batch-add-row!
(batch/make-batch 10 1024)
(.getBytes "1\n" StandardCharsets/UTF_8))
writer (atom nil)]
(.put (.queue pipeline) test-batch)
(is (zero? (.remainingCapacity (.queue pipeline))))
(.set ^AtomicBoolean (.done pipeline) true)
(with-redefs [copy/copy-sql (constantly "COPY public.items FROM STDIN")
batch/send-batch! (fn [& _] {:rows 1})]
(reset! writer
(future
(prefetch/writer-task
(pg-connection (atom false))
{:target-schema "public" :target-table "items"}
pipeline)))
(try
(is (= 1 (:rows-ok (deref @writer 500 ::timed-out))))
(finally
(future-cancel @writer)))))))

(deftest on-error-stop-does-not-retry-a-failed-copy-batch
(let [rolled-back? (atom false)
retried? (atom false)
test-batch (batch/batch-add-row!
(batch/make-batch 10 1024)
(.getBytes "1\n" StandardCharsets/UTF_8))
error (PSQLException. "bad value" PSQLState/DATA_ERROR)]
(with-redefs [batch/send-batch! (fn [& _] (throw error))
batch/retry-batch! (fn [& _]
(reset! retried? true)
{:rows-ok 0 :errors 1})]
(binding [copy/*on-error-stop* true]
(is (thrown-with-msg?
PSQLException #"bad value"
(#'prefetch/send-batch-or-retry!
(pg-connection rolled-back?)
{:target-schema "public" :target-table "items"}
"COPY public.items FROM STDIN"
test-batch 0 0 0 0 nil)))))
(testing "the active transaction is rolled back before propagation"
(is @rolled-back?))
(testing "strict mode never enters row rejection"
(is (false? @retried?)))))

(deftest resume-mode-still-retries-a-failed-copy-batch
(let [rolled-back? (atom false)
retried? (atom false)
test-batch (batch/batch-add-row!
(batch/make-batch 10 1024)
(.getBytes "1\n" StandardCharsets/UTF_8))
error (PSQLException. "bad value" PSQLState/DATA_ERROR)]
(with-redefs [batch/send-batch! (fn [& _] (throw error))
batch/retry-batch! (fn [& _]
(reset! retried? true)
{:rows-ok 1 :errors 1})]
(binding [copy/*on-error-stop* false]
(is (= {:status :retry
:rows-ok 1
:errors 1
:bytes 2
:reject-paths nil}
(select-keys
(#'prefetch/send-batch-or-retry!
(pg-connection rolled-back?)
{:target-schema "public" :target-table "items"}
"COPY public.items FROM STDIN"
test-batch 0 0 0 0 nil)
[:status :rows-ok :errors :bytes :reject-paths])))))
(is @rolled-back?)
(is @retried?)))

(deftest rollback-failure-preserves-the-copy-error-and-aborts
(let [rolled-back? (atom false)
retried? (atom false)
test-batch (batch/batch-add-row!
(batch/make-batch 10 1024)
(.getBytes "1\n" StandardCharsets/UTF_8))
copy-error (PSQLException. "bad value" PSQLState/DATA_ERROR)
rollback-error (java.sql.SQLException. "rollback failed")]
(with-redefs [batch/send-batch! (fn [& _] (throw copy-error))
batch/retry-batch! (fn [& _]
(reset! retried? true)
{:rows-ok 0 :errors 1})]
(binding [copy/*on-error-stop* false]
(let [thrown (try
(#'prefetch/send-batch-or-retry!
(pg-connection rolled-back? rollback-error)
{:target-schema "public" :target-table "items"}
"COPY public.items FROM STDIN"
test-batch 0 0 0 0 nil)
nil
(catch PSQLException e e))]
(is (identical? copy-error thrown))
(is (= [rollback-error] (vec (.getSuppressed thrown)))))))
(is @rolled-back?)
(is (false? @retried?))))
Loading