Repository navigation
VE-5598915 + VE-5602339: parallel fetch with thread-safe write-back - #24
Conversation
e05b2c3 to
17c0cd6
Compare
mkottakota1
left a comment
There was a problem hiding this comment.
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
- Possible out-of-bounds read on strings: a
SQL_NO_TOTAL(-4) indicator slips throughemitCell()'s string handling and reachescopy()as a negative length. (ODBCLoader.cpp:702) - CI isn't actually exercising concurrency —
Threading = 0is 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.
| 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; |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
(len == SQL_NTS || len < 0) — SQL_NTS is -3, so it's already covered by len < 0. The first half can go.
There was a problem hiding this comment.
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.
|
|
||
| // 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) { |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.
|
|
||
| SQLRETURN fetchRet = SQL_SUCCESS; | ||
| while (!done && !workerStatus[workerIdx].failed) { | ||
| if (isCanceled()) break; // US 5598915: prompt exit on cancel |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.
| queue.notEmpty.notify_one(); | ||
| } | ||
|
|
||
| if (!done && !SQL_SUCCEEDED(fetchRet) && fetchRet != SQL_NO_DATA && |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| // 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); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| \timing off | ||
|
|
||
| -- Set output to a fixed timezone regardless of where this is being tested | ||
| set time zone to 'EST'; |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
|
||
| **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``: |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
|
||
| 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``. |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Done. Reworded to "a small fixed number of converted batches (currently 8)" so it's clearly tied to the constant rather than restating it.
| 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..." |
There was a problem hiding this comment.
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
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
left a comment
There was a problem hiding this comment.
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.
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.
std::atomic<bool>that workers read.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.