ray.data.from_dask#

ray.data.from_dask(df: dask.DataFrame) ray.data.dataset.Dataset[ray.data._internal.arrow_block.ArrowRow][source]#

Create a dataset from a Dask DataFrame.

Parameters

df – A Dask DataFrame.

Returns

Dataset holding Arrow records read from the DataFrame.

PublicAPI: This API is stable across Ray releases.