如何使用Dask创建带有PyArrow数据类型的数组?
在Dask中创建带PyArrow数据类型的数组
问题描述
我在Pandas里可以这样创建带PyArrow数据类型的Series:
>>> import pandas as pd >>> s = pd.Series([1,2,3]).astype("int64[pyarrow]") >>> s.dtype int64[pyarrow]
但在Dask里尝试了几种方法都没成功:
- 配置
array.pyarrow_dtype后创建数组,返回的还是numpy int64类型:
>>> import dask.config >>> import dask.array as da >>> dask.config.set({"array.pyarrow_dtype": True}) >>> s = da.array([1,2,3]) >>> s
- 直接指定
"int64[pyarrow]"作为dtype,触发类型错误:
>>> import dask.array as da >>> s = da.array([1,2,3], dtype="int64[pyarrow]") TypeError: data type 'int64[pyarrow]' not understood
- 传入PyArrow的
pa.int64()类型,同样触发错误:
>>> import dask.array as da >>> import pyarrow as pa >>> s = da.array([1,2,3], pa.int64()) TypeError: Cannot interpret 'DataType(int64)' as a data type
请问是否可以用Dask创建带有PyArrow数据类型的数组?
解决方案
可以通过以下两种方式实现:
方法1:借助Dask DataFrame转换
先创建带PyArrow dtype的Pandas Series,再转为Dask Series或Array:
import pandas as pd import dask.dataframe as dd # 创建PyArrow dtype的Pandas Series pd_series = pd.Series([1,2,3]).astype("int64[pyarrow]") # 转为Dask Series dd_series = dd.from_pandas(pd_series, npartitions=1) # 提取Dask Array da_array = dd_series.to_dask_array() print(da_array.dtype) # 输出 int64[pyarrow]
方法2:基于PyArrow数组直接创建
用PyArrow数组作为输入生成Dask Array:
import pyarrow as pa import dask.array as da # 创建PyArrow数组 pa_array = pa.array([1,2,3], type=pa.int64()) # 转为Dask Array,指定分块大小 da_array = da.from_array(pa_array, chunks=2) print(da_array.dtype) # 输出 int64[pyarrow]
注意:当前Dask Array对PyArrow dtype的原生支持仍在迭代完善,部分操作可能需要依赖Dask DataFrame的PyArrow集成,建议使用2024.01.0及以上版本的Dask以获得更好的兼容性。
内容的提问来源于stack exchange,提问作者Adrien Pacifico
相关产品推荐
相关产品推荐

