Dask DataFrame操作异常:set_index参数报错与merge类型错误求助
合并Dask DataFrame时遇到的两类问题
问题1:调用set_index()传入divisions参数触发TypeError
提示“set_index() got an unexpected keyword argument 'divisions'”,报错详情:
--------------------------------------------------------------------------- TypeError Traceback (most recent call last) Cell In [12], line 5 2 import dask.dataframe as dd 4 df_201_1_2_3_sorted = df_201_1_2_3.set_index("docdb_family_id", divisions=unique_divisions2) ----> 5 df_225_228_sorted = df_225_228.set_index("docdb_family_id", divisions=unique_divisions2) File ~/.local/lib/python3.8/site-packages/pandas/util/_decorators.py:311, in deprecate_nonkeyword_arguments.<locals>.decorate.<locals>.wrapper(*args, **kwargs) 305 if len(args) > num_allow_args: 306 warnings.warn( 307 msg.format(arguments=arguments), 308 FutureWarning, 309 stacklevel=stacklevel, 310 ) --> 311 return func(*args, **kwargs) TypeError: set_index() got an unexpected keyword argument 'divisions'
解决方法
Dask 2022.9.1版本的set_index()不支持直接传入divisions参数,可根据数据状态选择以下方式:
- 若数据已按
docdb_family_id排序:先执行df.set_index("docdb_family_id"),再调用repartition(divisions=unique_divisions2)设置分区边界。 - 若数据未排序:使用
set_index("docdb_family_id", shuffle='disk')通过磁盘shuffle完成排序并设置索引。
问题2:merge操作触发TypeError,提示无法识别Dask DataFrame
即使不使用divisions参数,执行merge时仍报错“Can only merge Series or DataFrame objects, a <class 'dask.dataframe.core.DataFrame'> was passed”,报错详情:
--------------------------------------------------------------------------- TypeError Traceback (most recent call last) Cell In [11], line 3 1 ####PROBLEMA: QUESTO FA RIPARTIRE IL KERNEL!!!! SI STOPPA. COME FARE????? ----> 3 large_join2 = df_225_228_sorted.merge( 4 df_201_1_2_3_sorted, 5 how="left", 6 on=["docdb_family_id"] 7 #left_index=True, 8 #right_index=True 9 ).persist() File ~/.local/lib/python3.8/site-packages/pandas/core/frame.py:9351, in DataFrame.merge(self, right, how, on, left_on, right_on, left_index, right_index, sort, suffixes, copy, indicator, validate) 9332 @Substitution("") 9333 @Appender(_merge_doc, indents=2) 9334 def merge( (...) 9347 validate: str | None = None, 9348 ) -> DataFrame: 9349 from pandas.core.reshape.merge import merge -> 9351 return merge( 9352 self, 9353 right, 9354 how=how, 9355 on=on, 9356 left_on=left_on, 9357 right_on=right_on, 9358 left_index=left_index, 9359 right_index=right_index, 9360 sort=sort, 9361 suffixes=suffixes, 9362 copy=copy, 9363 indicator=indicator, 9364 validate=validate, 9365 ) File ~/.local/lib/python3.8/site-packages/pandas/core/reshape/merge.py:107, in merge(left, right, how, on, left_on, right_on, left_index, right_index, sort, suffixes, copy, indicator, validate) 90 @Substitution(" left : DataFrame or named Series") 91 @Appender(_merge_doc, indents=0) 92 def merge( (...) 105 validate: str | None = None, 106 ) -> DataFrame: --> 107 op = _MergeOperation( 108 left, 109 right, 110 how=how, 111 on=on, 112 left_on=left_on, 113 right_on=right_on, 114 left_index=left_index, 115 right_index=right_index, 116 sort=sort, 117 suffixes=suffixes, 118 copy=copy, 119 indicator=indicator, 120 validate=validate, 121 ) 122 return op.get_result() File ~/.local/lib/python3.8/site-packages/pandas/core/reshape/merge.py:629, in _MergeOperation.__init__(self, left, right, how, on, left_on, right_on, axis, left_index, right_index, sort, suffixes, copy, indicator, validate) 611 def __init__( 612 self, 613 left: DataFrame | Series, (...) 626 validate: str | None = None, 627 ): 628 _left = _validate_operand(left) --> 629 _right = _validate_operand(right) 630 self.left = self.orig_left = _left 631 self.right = self.orig_right = _right File ~/.local/lib/python3.8/site-packages/pandas/core/reshape/merge.py:2285, in _validate_operand(obj) 2283 return obj.to_frame() 2284 else: --> 2285 raise TypeError( 2286 f"Can only merge Series or DataFrame objects, a {type(obj)} was passed" 2287 ) TypeError: Can only merge Series or DataFrame objects, a <class 'dask.dataframe.core.DataFrame'> was passed
解决方法
报错显示调用了pandas的merge方法,但该方法无法处理Dask DataFrame对象,可按以下步骤修复:
- 确认两个待合并对象均为Dask DataFrame:用
type(df)检查,若存在pandas DataFrame,需用dd.from_pandas(pandas_df, npartitions=4)(分区数按需调整)转换为Dask对象。 - 使用Dask原生merge方法替代DataFrame实例的merge方法:
large_join2 = dd.merge( df_225_228_sorted, df_201_1_2_3_sorted, how="left", on=["docdb_family_id"] ).persist()
当前使用的Dask版本
dask 2022.9.1 pyhd8ed1ab_0 conda-forge dask-core 2022.9.1 pyhd8ed1ab_0 conda-forge
内容的提问来源于stack exchange,提问作者Lusian
相关产品推荐
相关产品推荐

