perf(csv): parallel splayed load, batched symfile interning, parallel hash index build - #632
Open
ser-vasilich wants to merge 8 commits into
Open
ser-vasilich wants to merge 8 commits into
ser-vasilich wants to merge 8 commits into
Conversation
Streaming a CSV into a splayed table ran its per-chunk row scan, the per-column writes and the index build on the calling thread; only the field parse used the pool. - Row offsets for a chunk come from the parallel quote-parity scanner over a byte window sized from the running average row length, widened until it holds the chunk (falls back to the serial limited scan). - The chunk's columns are appended by one task per column; SYM cells go through a direct-mapped runtime-id to position cache in the writer, so only a value's first occurrence probes the symfile domain. - ray_splay_build_indexes builds one column per task. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Every distinct string of a splayed load was interned twice, serially: once into the runtime symbol table by the chunk materialisation, then again into the table's symfile domain by the column writer, one cell at a time under the global domain spinlock — the two together were more than half of the load time, and the runtime table grew by the whole vocabulary for nothing. - csv_materialize_rows takes a target domain: SYM chunk columns are built over the symfile domain, and the writer copies their positions. - ray_sym_domain_intern_batch (domain.c): one batch per chunk over all SYM columns. The existing vocabulary is probed in parallel against a snapshot of the reverse index (replaced bucket tables are retired, never freed), misses are deduplicated per hash partition, and the new atoms are built in parallel into per-partition arena regions (ray_arena_alloc_raw / ray_arena_str_at) before the count is published. Index growth re-slots the old table instead of rehashing every atom; the symfile flush packs records into a 1 MB buffer. - SYM fields are hashed by the parse tasks; with a domain target each SYM column's dedupe is split by hash partition across tasks. - The window scan is preceded by a readahead hint and finished chunks are dropped from the mapping, so a long file does not stay resident. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
ray_index_attach_hash grouped the rows in one serial open-addressing pass; on a high-cardinality numeric column of a large table it was the longest single stretch of a splayed load, and running the columns as pool tasks left the build itself single-threaded. Numeric and SYM keys above 64k rows now build in partition-parallel passes: key words and hash-partition buckets over row ranges, a per-partition dedupe that marks each group's first row in a bitmap, then keys, counts, the row scatter and the key table (CAS on the empty slot) per partition. A group's number is the rank of its first row among the marked rows, so the layout is the serial walk's: groups in first-occurrence order, rows ascending inside a group. Every array a pass writes is faulted in beforehand from all workers in slices — a fresh mapping faulted at random from every worker serialised on the page-table locks. STR keys and small columns keep the serial walk. ray_splay_build_indexes defers the columns that qualify for a hash index out of the per-column tasks, builds them one after another with the parallel path and writes all deferred columns together. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
cppcheck cannot parse a declaration of `_Atomic(T)*` type (AST broken on the `=`), which failed static analysis; the two batch inserts now cast at the compare-exchange, the form the rest of domain.c uses. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Rayforce targeted audit passedThe required Rayforce audit gate passed on the latest run. Workflow run: https://github.com/RayforceDB/rayforce/actions/runs/36337382081 |
…n tests The batch intern probes a snapshot of a domain's reverse index outside the lock, relying on replaced tables being retired rather than freed. The external-growth path (a reader catching up with a symfile another process appended to) still freed the table, so a probe running at that moment could read freed memory. It now retires it like the rebuild does. A probe that meets an entry it cannot compare (no atom and no raw bytes, possible only if publishing the raw snapshot failed) no longer counts it as a miss: the whole batch then resolves under the lock, so no string is appended twice. The splayed writer's per-column dispatch goes through ray_pool_par_dispatch_ok like the other sites. The batch-intern header comment no longer claims appends in batch order: new strings get positions grouped by hash partition. Tests: domain/intern_batch (repeats across partitions, pre-interned hits keep their positions, repeat batch appends nothing, "" is 0, flush/reopen); csv_splayed_dedupe_overflow (140k distinct strings in one SYM column overflow a partition's dictionary in debug builds and take the row-by-row fallback); index/hash_large_parallel now creates the pool so the parallel build is the one tested. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…es independent of cores A chunk's byte window chose the scanner's quote mode from its own bytes. In a file with quotes elsewhere, a window without any took the quote-free fast path, where a lone '\r' does not end a row: two rows merged and a value was lost (the serial walk and .csv.read split them). The window now scans in the file's quote mode. The symfile no longer depends on the worker count. The chunk batch is laid out the way the cell-by-cell writer meets the strings — columns in order, each column's strings by first occurrence (partition runs merged by the first row, recorded in the dictionary entry) — and the batch intern appends new strings in batch order. The hash index key table is filled in group order on the calling thread, as the serial build does, so the persisted index is the same bytes on any core count. A replaced reverse-index table is freed at once unless a batch intern is probing a snapshot outside the lock (counted under the lock), so a long-lived reader of a growing symfile keeps no extra tables. Tests: splay/csv_symfile_order (positions follow the writer order over three chunks, parallel batch and partitioned dedupe), and splay/csv_quote_mode_per_file (quotes only in chunk 0, a lone '\r' in chunk 2: same rows and values as .csv.read). Both fail with their fix disabled. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Bug
Streaming a CSV into a splayed table was mostly serial. The per-chunk row scan, the per-column writes and the index build ran on the calling thread; every distinct string was interned twice, serially — once into the runtime symbol table by the chunk materialisation, then again into the table's symfile domain by the column writer, one cell at a time under the global domain spinlock; and the hash index of a high-cardinality numeric column was one serial pass over all rows. On a 100M-row, 105-column file the load took 12–15 minutes on 48 threads (about an hour on 16), and the runtime symbol table grew by the whole vocabulary for nothing.
Fix
Three commits:
Parallel scan, append and index. Row offsets for a chunk come from the parallel quote-parity scanner over a byte window sized from the running average row length; the chunk's columns are appended one task per column;
ray_splay_build_indexesbuilds one column per task.Chunk dictionaries interned straight into the symfile domain.
csv_materialize_rowstakes a target domain: SYM chunk columns are built over the symfile domain and the writer copies their positions. Newray_sym_domain_intern_batchhandles one batch per chunk over all SYM columns: the existing vocabulary is probed in parallel against a snapshot of the reverse index (replaced bucket tables are retired, never freed), misses are deduplicated per hash partition, and the new atoms are built in parallel into per-partition arena regions (ray_arena_alloc_raw/ray_arena_str_at) before the count is published. Index growth re-slots the old table instead of rehashing every atom; the symfile flush is buffered. SYM fields are hashed by the parse tasks and, with a domain target, each column's dedupe is split by hash partition. The window scan gets a readahead hint and finished chunks are dropped from the mapping, so a long file does not stay resident.Partition-parallel hash index build. Numeric and SYM keys above 64k rows build in partition-parallel passes: key words and hash-partition buckets over row ranges, a per-partition dedupe that marks each group's first row, then keys, counts, row scatter and the key table per partition. A group's number is the rank of its first row, so the layout is the serial walk's (groups in first-occurrence order, rows ascending inside a group). Arrays a pass writes are faulted in beforehand from all workers in slices.
ray_splay_build_indexesdefers the columns that qualify for a hash index out of the per-column tasks and writes all deferred columns together.Result
Same file, same machine (48 threads), tables byte-identical to the previous loader (row counts, column sums, distinct counts, string lengths, nulls, date ranges, grouped results):
Peak anonymous memory of the 100M load: 20 GB (the vocabulary of 60M distinct strings in the symfile domain).
Tests:
test/test_index.cindex/hash_large_parallelchecks the parallel build against a serial reference; full suite passes under ASan.