NVIDIA / NVIDIA/cudf

[ENH] benchmark gather then sort vs sort then gather in merge with `sort=True`

Open
#13,630 0 comments 0 reactions 0 assignees View on GitHub
feature request libcudf Performance Python
Dominant language
C++
Stars
9.8k
Forks
1.1k
Avg merge
3d 6m
Merged PRs (30d)
278

Description

**Is your feature request related to a problem? Please describe.**

When we request `sort=True` in a `cudf.merge`, the current implementation does:

1. deduce left and right join columns
2. join, producing left and right gather maps
3. gather left and right columns, and merge results
4. deduce key columns to sort by
5. argsort the key columns
6. gather the result using the argsort return value

Trivially, steps 5 and 6 can be merged into a `sort_by_key` (that's #13557). However, this order probably does more data movement than it needs to. This makes two calls to gather, and one sort-by-key, at the cost of moving the full dataframe through memory twice (once in step 3, once in step 6).

Instead, we could (if sorting) first gather only the key columns we will sort by, argsort those and then use that ordering to sort the left and right gather maps.

1. deduce left and right join columns
2. join, producing left and right gather maps
3. deduce left and right key columns to order by
4. gather left key columns with left map, right key columns with right map
5. sort-by-key the left and right gather maps with the columns from step 4
6. gather left and right columns with new gather maps and merge

This makes four calls to gather and one sort-by-key, but only moves the full dataframe through memory once (in step 6). For dataframes with many non-key columns this might well be an advantage. The latency will be a bit higher, but the total data movement will be less. For example, consider (for simplicity) a left join with one key column and 10 total columns in both left and right dataframes.

The current approach (once the left and right gather maps have been determined) gathers 20 columns in step 3, argsorts one column, then gathers 20 columns again (sort-by-key merges the sort + gather into argsort + gather at the libcudf level).

The proposed alternative would gather 1 column in step 4, sorts-by-key two columns (the two gather maps), then gathers 20 columns. So we move effectively 23 columns through memory rather than 41.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.