Skip to main content

Inside GPU execution

After decoding, the input is a collection of buffers in GPU VRAM. The engine runs one compiled schedule for each native fragment. Its tasks call GPU operations, retain the intermediate results still in use and release completed resources.

Columns and memory access​

customer_id7972amount12589status offsets0481220status characterspaidopenpaidrefunded
Each offset pair selects one string's bytes; refunded needs more space than paid. Row position connects values across columns. Optional validity bitmaps track nulls separately.

Fixed-width values sit next to one another within a column. Nullable columns carry a validity bitmap; strings use offsets into a character buffer. A scan over amount can read adjacent values without loading unrelated fields from each row. Adjacent GPU threads can combine adjacent accesses into fewer memory transactions, much as locality reduces cache-line traffic in CPU code.

Different operators create different access patterns and intermediate data:

OperationMemory access and temporary state
FilterEvaluate a column predicate and gather selected rows. Reading input can greatly exceed the size of the result.
Hash joinBuild a hash table from one input's keys, probe with the other input, then gather matching row pairs. Probes can touch scattered buckets; duplicate keys can expand the result.
Grouped aggregateCombine rows by key into group state, then merge partial state when input is chunked. More distinct keys require more state.
Sort / top-kCompare and reorder values; chunk candidates require a merge before they form a global result.

These access patterns explain why a compressed file's size cannot stand in for the query's VRAM demand. Decoded columns, operator state and output can coexist. Column reads and intermediate writes also consume device-memory bandwidth, even after input transfers over PCIe have completed.

CUDA streams and dependencies​

CudaStream identifies an ordered queue of GPU work. A CPU worker submits operations and can return before the GPU completes them. Operations submitted to one stream follow its ordering; work on separate streams needs explicit completion dependencies before sharing results.

Stream AStream BRoot streamwait for both producersDecode + filterChunk 1Partial sumsGroup state for chunk 1Decode + filterChunk 2Partial sumsGroup state for chunk 2Merge partial resultsPreserve the grouped-sum result
Both streams operate on the same GPU's VRAM. The merge waits for both partial results; independent chunks may overlap when resources permit.

An execution lane combines a worker's owned task state with a leased stream. Independent source chunks or tasks can use multiple lanes, bounded by service resources and the work available. The coordinator merges their outputs under the required SQL semantics. A DataFusion CPU partition does not determine lane ownership or width.

Cross-stream dependencies use completion fences, such as CUDA events. The consumer waits for the producer's work before reading its buffers, and buffer ownership lasts through the final use. More streams permit overlap; device resources and dependencies determine whether work overlaps in practice. The streams share the same VRAM and memory interface, so concurrent kernels can compete for memory bandwidth. All lanes allocate inside one query memory cap.

Relational data into GPU algorithms​

A SQL table function can run cuGraph, cuVS or cuML inside a larger query. cuDF columns supply edges or features; the algorithm returns columns for subsequent filters, joins and aggregation. Format conversion and result ownership determine which boundaries borrow buffers and which perform device copies.

Continue with GPU table functions for a PageRank example, the prepared-export lifetime and the differences between the three libraries. Then follow admission and resource release across the attempt.