NVIDIA / NVIDIA/cudf

[FEA] dask_cudf cross-partition type coercions

Open
#7,742 3 comments 0 reactions 0 assignees View on GitHub
dask feature request 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.**

It's been frustrating adapting cudf -> dask_cudf kernels in two basic areas around cross-partition type mismatches:

* ingest: loading json, csv, etc. that vary in column types across partitions: existence, nans, int vs float, etc. When the code writer isn't the user -- so a library, piece of software, a UI, this is common and you can't just workaround by specifying dtypes ahead of time

* compute: when doing data cleaning (ex: date inference) or some algs, it's unclear what `meta` should be ahead of time, only after you actually do the calc. dask will sample the first df... which is often wrong

**Describe the solution you'd like**

dask_cudf ingest operators: an auto-coercion flag ("when columns are in conflict across partitions, coerce to the closest common type, like float or str")

dask_cudf map, concat, etc: same thing

**Describe alternatives you've considered**

It may also be possible to make each operator smarter via sampling or other tricks. dask core and some cudf io seems to be experimenting here.

I like explicit flags b/c of their predictability/reliability, and uniformity... but ultimately, whatever work :)

**Additional context**

By default, I'm guessing this issue will be ignored & deprioritized ;-)

Before doing that, it may be worth polling dask_cudf users -- not devs -- how they feel about this ;-) my bet is people spend a surprising % of their time on a few issues around here, well before actual perf

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.