Skip to content

perf(csv): parallel splayed load, batched symfile interning, parallel hash index build - #632

Open
ser-vasilich wants to merge 8 commits into
devfrom
perf/csv-splayed-load
Open

ser-vasilich wants to merge 8 commits into
devfrom
perf/csv-splayed-load

Conversation

@ser-vasilich

@ser-vasilich ser-vasilich commented Sep 27, 2026 •

Copy link
Copy Markdown
Collaborator

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:

  1. 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_indexes builds one column per task.

  2. Chunk dictionaries interned straight into the symfile domain. csv_materialize_rows takes a target domain: SYM chunk columns are built over the symfile domain and the writer copies their positions. New ray_sym_domain_intern_batch handles 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.

  3. 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_indexes defers 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):

rows before after
10M 57.8 s 14 s
100M 700–900 s 130–145 s

Peak anonymous memory of the 100M load: 20 GB (the vocabulary of 60M distinct strings in the symfile domain).

Tests: test/test_index.c index/hash_large_parallel checks the parallel build against a serial reference; full suite passes under ASan.

ser-vasilich and others added 4 commits September 27, 2026 04:25
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>
@github-actions

github-actions Bot commented Sep 27, 2026 •

Copy link
Copy Markdown

Rayforce targeted audit passed

The required Rayforce audit gate passed on the latest run.

Workflow run: https://github.com/RayforceDB/rayforce/actions/runs/36337382081

ser-vasilich and others added 4 commits September 27, 2026 12:25
…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

No deployments
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.

2 participants