Skip to content

VE-5598915 + VE-5602339: parallel fetch with thread-safe write-back - #24

Merged
SandeepDave2 merged 3 commits into
mainfrom
VE-5598915_5602339_threaded_fetch_writeback
Sep 8, 2026
Merged

SandeepDave2 merged 3 commits into
mainfrom
VE-5598915_5602339_threaded_fetch_writeback

Conversation

@SandeepDave2

@SandeepDave2 SandeepDave2 commented Aug 21, 2026 •

Copy link
Copy Markdown
Collaborator

Builds on #23 (VE-5603316), which split the source query into N slices run sequentially. This PR runs those slices concurrently.

Design - producer/consumer, single shared queue

N worker threads each open their own ODBC connection and fetch one slice in parallel, pushing row batches into a shared bounded queue. The main thread is the sole consumer and the only thread that touches the writer.

  • Workers do only ODBC fetch + byte conversion into their own buffers - never writer, allocator, getStringRef, vt_report_error, or log.
  • Column metadata, including the ODBC C type per column, is computed once in setup() and shared read-only. Workers make no Vertica SDK calls at all - cancellation is published by the main thread into a std::atomic<bool> that workers read.
  • Worker errors are marshalled as (message, sqlstate) and re-raised on the main thread; a failed worker aborts the COPY and Vertica's rollback prevents partial commit.
  • Bounded queue gives backpressure; a worker that fails sets the queue's shutdown flag so its siblings and the consumer stop immediately rather than draining the rest. Workers are joined on every path (completion, error, cancel), with a safety-net join in destroy(). If thread creation itself fails, the unstarted slices are marked failed and discounted from the producer count so the consumer can't wait on a drain that will never arrive.
  • thread_count is capped at min(slices, thread_count, MAX_THREAD=64).
  • split_method='modulo' is now also enabled for MariaDB.

Backward compatibility

With thread_count unset or 1 (or after any fallback), no threads spawn and the path is byte-identical to the pre-change loader.

Deployment note

For genuine parallelism, the driver's odbcinst.ini section must set Threading = 0; otherwise unixODBC serializes calls and slices run sequentially (results stay correct, speedup is lost). This is documented in the README under Usage -> Parallel fetch, and the shipped test config (tests/config/odbcinst.ini) now sets it, so CI exercises real concurrency.

Testing

10M-row source, N=4/8/16 - row count and sum(id) identical across single/range/modulo; string byte-fidelity verified via exact numeric(38,0) content hashes over fixed and variable-length VARCHARs, repeated across runs; NULL/negative-key regressions hold; fallback (bad column, bad DSN) and cancellation paths return cleanly with no orphan threads.

CI suite (tests/multithread_test.sql) follows the existing harness - covers all thread modes, backpressure, fallback branches, bounds validation, negative split keys, a mid-fetch worker failure (asserting the COPY aborts and the target table is left empty), and an external-table path. This is correctness/regression coverage; concurrency under contention is proven by the 10M-row benchmark.

Performance

Benchmarked with Vertica and the Postgres source on separate 12-core / 32 GB RHEL 8.10 boxes, loading 10M rows over the network.

Loading took about 80 seconds on a single connection via direct COPY, and about 77 seconds via the external-table route. With 8 threads direct COPY dropped to roughly 66 seconds (~1.21x), and the external-table route to about 61 seconds, improving slightly further to ~60 seconds at 16 threads (~1.28x).

The gain is modest by design. Only about a third of the load time is spent fetching data; the rest is Vertica writing it to disk, which threads can't speed up. So parallel fetch shrinks the fetch portion and the load time settles near the disk-write floor instead of dropping by the number of threads. Around 8 threads is where the benefit levels off - past that you're adding connections without adding throughput, and on a 12-core box you also run out of cores.

@SandeepDave2
SandeepDave2 marked this pull request as ready for review August 21, 2026 12:53
@SandeepDave2
SandeepDave2 force-pushed the VE-5598915_5602339_threaded_fetch_writeback branch from e05b2c3 to 17c0cd6 Compare August 26, 2026 19:30

@mkottakota1 mkottakota1 left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Overall the design is solid — the bounded queue with backpressure, clean worker teardown (join before freeing handles), errors marshalled back to the main thread instead of thrown across threads, and the order-independent aggregate tests are all well thought out. Nice work.

I have a few comments to address before I approve:

Must-fix

  1. Possible out-of-bounds read on strings: a SQL_NO_TOTAL (-4) indicator slips through emitCell()'s string handling and reaches copy() as a negative length. (ODBCLoader.cpp:702)
  2. CI isn't actually exercising concurrency — Threading = 0 is missing from [MYODBC], so unixODBC serializes the workers. (odbcinst.ini:10)

Should-fix
3. The type-conversion switch is duplicated between emitCell() and the single-slice loop — will drift. (ODBCLoader.cpp:473)
4. Confirm isCanceled() is safe to call from worker threads, or use an atomic flag. (ODBCLoader.cpp:682)
5. Add a test for a worker failing mid-fetch (the new error/rollback path is untested). (multithread_test.sql:88)

Plus a few minor cleanups (fail-fast on worker error, per-cell allocation overhead, redundant SQL_NTS check, README/config Threading mismatch).

Happy to approve once these are addressed.

Comment thread ODBCLoader.cpp Outdated
row[i].isNull = false;
// A negative indicator (SQL_NTS/SQL_NO_TOTAL) gives no length, so
// copy the whole field; emitCell() re-measures it on the main thread.
size_t n = (len == SQL_NTS || len < 0) ? (size_t)stype[i] : (size_t)len;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Possible crash on strings when the driver doesn't report a length

ODBCLoader.cpp:702 → consumed at ODBCLoader.cpp:507

In the worker, when the driver returns a negative length indicator we copy the full field but then save the raw indicator into lenIndicator. That raw value can be SQL_NO_TOTAL (-4), not just SQL_NTS (-3). Later in emitCell() the string branch only fixes up SQL_NTS, so a -4 gets passed straight to writer->getStringRef().copy(buf, len) as the length. A negative length becomes a huge unsigned number and we read past the buffer → likely crash. Two easy fixes: either store the number of bytes we actually copied (n) instead of the raw indicator, or in emitCell() treat any len < 0 the same as SQL_NTS (re-measure with strnlen). I'd lean toward storing n.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

(len == SQL_NTS || len < 0) — SQL_NTS is -3, so it's already covered by len < 0. The first half can go.

@SandeepDave2 SandeepDave2 Aug 31, 2026 •

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. The worker now stores SQL_NTS for any negative indicator instead of the raw value, and emitCell()'s string branch triggers on data.len < 0 rather than == SQL_NTS, so a negative length can never reach copy() from either path. I kept the re-measure approach rather than storing n, because when the driver reports no length we copy the whole fixed-width field - storing n would include the trailing padding and corrupt the value.
Also dropped the redundant len == SQL_NTS || -as len < 0 already covers it.

Comment thread ODBCLoader.cpp

// Converts one already-fetched cell (raw driver bytes) into the writer; main
// thread only. NOTE: the single-slice loop in process() still has its own copy.
void emitCell(ServerInterface &srvInterface, SQLUSMALLINT i, SQLPOINTER buf, SQLLEN len) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The big type-conversion switch now exists in two places

ODBCLoader.cpp:473 and the original at ODBCLoader.cpp:893

emitCell() is basically a copy-paste of the existing conversion switch in the single-slice loop — the comment even says so. The tricky bits (dates, intervals, numerics, the Oracle int8 special case) are now duplicated, and the next person who fixes a bug in one will almost certainly forget the other. Could we make the single-slice path build a small Buf and call emitCell() too, so there's just one conversion routine to maintain?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. Removed the inline switch from the single-slice loop - it now builds the same pointer/length pair and calls emitCell(), so there's one conversion routine with two callers. The NULL check and the setNull() pre-fill stayed in the caller, so the single-connection path is unchanged.

Comment thread ODBCLoader.cpp Outdated

SQLRETURN fetchRet = SQL_SUCCESS;
while (!done && !workerStatus[workerIdx].failed) {
if (isCanceled()) break; // US 5598915: prompt exit on cancel

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are we allowed to call isCanceled() from the worker threads?

ODBCLoader.cpp:682 and ODBCLoader.cpp:722

We were careful to keep the workers away from writer/ServerInterface, but the workers still call isCanceled(), which is a Vertica SDK call. I couldn't find anything saying it's safe to call off the main thread — if it reads shared session/plan state internally, that's a data race. Safer pattern: only the main thread calls isCanceled(), and it flips an std::atomic that the workers check. Can we confirm the SDK guarantee, or switch to the atomic flag?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done - switched to the atomic flag rather than rely on an unconfirmed guarantee. Added a std::atomic<bool>, the main thread publishes isCanceled() in the consumer loop and signals the queue's shutdown path so a blocked worker wakes, and workers only read the atomic. No SDK calls off the main thread now.

Comment thread ODBCLoader.cpp
queue.notEmpty.notify_one();
}

if (!done && !SQL_SUCCEEDED(fetchRet) && fetchRet != SQL_NO_DATA &&

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When one worker fails, the rest keep working for nothing

If a worker hits an error it just leaves its own loop, but the other workers keep fetching and the main thread keeps converting/emitting rows until the whole queue drains — and only then do we raise the error and roll back. On a big load that's a lot of wasted CPU and memory after we already know the load is going to fail. Setting queue.shutdown = true as soon as the first worker fails would let everyone stop early.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. A worker that records an error now sets queue.shutdown = true and notifies both condition variables on its way out, so the siblings and the consumer stop immediately instead of draining the whole queue.

Comment thread ODBCLoader.cpp
// copy the whole field; emitCell() re-measures it on the main thread.
size_t n = (len == SQL_NTS || len < 0) ? (size_t)stype[i] : (size_t)len;
if (n > (size_t)stype[i]) n = (size_t)stype[i];
row[i].bytes.assign((char*)wresp[i] + (size_t)stype[i]*j, n);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One heap allocation per cell may eat the speedup

Every non-null value does bytes.assign(...), which is a separate std::string allocation for every cell of every row. With wide tables and a large rowset, that allocation/copy overhead on the worker side could cancel out a good chunk of what we gained from parallelism. Worth considering copying each fetched batch into one contiguous buffer and tracking per-cell offset/length, instead of a std::string per cell.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed it's a real cost on wide tables, but I'd prefer it as a follow-up - it's an optimisation rather than a defect, and the contiguous-buffer version would rework the exact byte-marshalling path this PR just fixed and verified byte-for-byte (10M rows, fixed- and variable-width VARCHARs, content hashes identical across single / 8-thread / 16-thread runs).
It also doesn't look like the limiting factor today: measured end-to-end at 1.21x on direct COPY (79.8 s -> 65.7 s at 8 threads) and 1.28x via the external-table route (77.1 s -> 60.1 s at 16 threads) - and that was measured on this same allocation strategy, since the per-cell std::string is unchanged by these fixes. The real ceiling is storage: fetch is ~36% of load time, storage ~64% and disk-bound, which caps the gain near 1.5x regardless of thread count.

Comment thread tests/multithread_test.sql Outdated
\timing off

-- Set output to a fixed timezone regardless of where this is being tested
set time zone to 'EST';

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Leftover set time zone that this test doesn't need

this test only selects i and v, no date/time columns, so set time zone to 'EST' doesn't affect any output here. Looks copied from the other test files. Can drop it to avoid implying there's a temporal dependency.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, removed the line and its comment.


-- Multi-threaded fetch: thread_count / split_column / split_method.
-- Slice order is not deterministic, so every check below is an order-independent
-- aggregate over testdb.test_source (9 data rows plus 1 all-NULL row). A dropped

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The expected totals (9 data rows + 1 NULL, sum_i=45, len_v=54) hard-code the current contents of testdb.test_source. That's the same assumption copy_test.sql already makes, so I'm fine with it — just flagging that if anyone ever edits that shared fixture, this test breaks too. A one-line comment pointing at where the fixture is defined would help the next person.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. Added a header comment noting the totals track the shared testdb.test_source fixture created by .github/workflows/ci.yml, so anyone editing it knows this test needs updating too.

Comment thread README.md

**Memory.** Each worker binds its own ``rowset``-sized fetch buffers and up to 8 converted batches may be queued, so the buffered memory is roughly ``thread_count x rowset x row_width``. ``thread_count`` multiplies the cost of ``rowset``: raising both at once is the easy way to exceed ``FencedUDxMemoryLimitMB``.

**unixODBC.** Parallel fetch only delivers a speed-up if the driver manager is allowed to run driver calls concurrently. unixODBC serialises them unless the driver's section in ``odbcinst.ini`` sets ``Threading = 0``:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This section correctly says you need Threading = 0 for real concurrency, but our own odbcinst.ini doesn't set it (see comment A). Anyone copying our CI config as a starting point will silently get no parallelism. Let's make the two consistent — ideally set it in the config and keep this note.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. Added Threading = 0 to both driver sections in tests/config/odbcinst.ini, so the shipped config matches the README and CI now actually exercises concurrency.

Comment thread README.md Outdated

The loader falls back to a single connection - without failing the load - when ``split_column`` is missing, is not a bare identifier, is not integer-valued, or has been pruned out of the remote query. The reason is written to the UDx log. Rows are **not** returned in a deterministic order when more than one slice is used.

**Memory.** Each worker binds its own ``rowset``-sized fetch buffers and up to 8 converted batches may be queued, so the buffered memory is roughly ``thread_count x rowset x row_width``. ``thread_count`` multiplies the cost of ``rowset``: raising both at once is the easy way to exceed ``FencedUDxMemoryLimitMB``.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

up to 8 converted batches may be queued" mirrors MAX_QUEUE_BATCHES in the code. If someone tunes that constant later, this doc silently goes stale. Not a blocker — maybe just phrase it as "a small fixed number of batches (currently 8)" so it's clearly tied to the constant.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. Reworded to "a small fixed number of converted batches (currently 8)" so it's clearly tied to the constant rather than restating it.

Comment thread Makefile
s/(Error parsing .* )\(.*\)$$/$$1(...)/; \
s/mariadb/MySQL/ig; ' $(TMPDIR)/federated_queries.out) && echo "[✓] Test validation passed" || (echo "[✗] Test validation failed" && exit 1)
@echo ""
@echo "[*] Running Multi-threaded Fetch Test..."

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new multithread target follows the same tee/normalize/diff pattern as the federated test, and the extra s/^vsql: ERROR \d+:.*message: /vsql: ERROR: / rule is what lets Test 14's two error lines match deterministically. No concerns here.

…conversion routine, atomic cancel flag, fail-fast on worker error, CI threading config, and mid-fetch failure test
@SandeepDave2

Copy link
Copy Markdown
Collaborator Author

Overall the design is solid — the bounded queue with backpressure, clean worker teardown (join before freeing handles), errors marshalled back to the main thread instead of thrown across threads, and the order-independent aggregate tests are all well thought out. Nice work.

I have a few comments to address before I approve:

Must-fix

  1. Possible out-of-bounds read on strings: a SQL_NO_TOTAL (-4) indicator slips through emitCell()'s string handling and reaches copy() as a negative length. (ODBCLoader.cpp:702)
  2. CI isn't actually exercising concurrency — Threading = 0 is missing from [MYODBC], so unixODBC serializes the workers. (odbcinst.ini:10)

Should-fix 3. The type-conversion switch is duplicated between emitCell() and the single-slice loop — will drift. (ODBCLoader.cpp:473) 4. Confirm isCanceled() is safe to call from worker threads, or use an atomic flag. (ODBCLoader.cpp:682) 5. Add a test for a worker failing mid-fetch (the new error/rollback path is untested). (multithread_test.sql:88)

Plus a few minor cleanups (fail-fast on worker error, per-cell allocation overhead, redundant SQL_NTS check, README/config Threading mismatch).

Happy to approve once these are addressed.

All five points and the minor cleanups are addressed - details on each thread. The only one deferred is the per-cell allocation (reasoning on that thread).

Verified two ways: CI is green, and I re-ran the full 10M-row suite locally against this commit — split / NULL-key / negative-key / parallel at N=4, 8, 16 / modulo all return 10,000,000 rows with identical sums, the string content hashes come back identical to the pre-fix build, and the fallback, bad-DSN and cancellation paths all return cleanly with no orphaned threads.

@mkottakota1 mkottakota1 left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the thorough turnaround — all points are addressed. String-length fix and the atomic cancel flag look correct, the shared emitCell() removes the duplication nicely, and Test 16 covers the mid-fetch failure/rollback path. Once CI is green I'm happy to approve.

@SandeepDave2
SandeepDave2 merged commit 07805c7 into main Sep 8, 2026
2 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants