如何为Dask DataFrame按行添加基于intensity列的随机数?
问题:Dask DataFrame按行生成自定义范围随机数报错
我想给Dask DataFrame新增一列随机数,每行的随机数范围由原DataFrame的intensity列决定。这段逻辑用Pandas和NumPy能正常运行,但换成Dask和dask.array就报错了。
代码示例
import dask.array as da import dask.dataframe as dd from dask.distributed import Client client = Client() fns = [list-of-filenames] df = dd.read_parquet(fns) # dataframe 包含float类型的intensity列,无缺失值 df['separation_dimension_1'] = da.random.uniform(size=N, low=-noise_level/df.intensity, high=noise_level/df.intensity)
报错信息
ValueError: shape mismatch: objects cannot be broadcast to a single shape. Mismatch is between arg 0 with shape (0,) and arg 1 with shape (33276691,).
完整报错栈
Cell In[21], line 7 5 df['mz_'] = df.mz * 1000000000 6 df['rt_'] = df.scan_time*10 ----> 7 df['separation_dimension_1'] = da.random.uniform(size=N, low=-noise_level/df.intensity, high=noise_level/df.intensity) 8 #df['separation_dimension_2'] = da.random.uniform(size=N, low=-noise_level/df.intensity, high=noise_level/df.intensity) 9 #df['separation_dimension_3'] = da.random.uniform(size=N, low=-noise_level/df.intensity, high=noise_level/df.intensity) 11 df = df[df.intensity > 1e5][['rt_', 'mz_', 'logint']] File ~/miniconda3/envs/dask/lib/python3.9/site-packages/dask/array/random.py:465, in _make_api.<locals>.wrapper(*args, **kwargs) 462 if backend not in _cached_random_states: 463 # Cache the default RandomState object for this backend 464 _cached_random_states[backend] = RandomState() ---> 465 return getattr( 466 _cached_random_states[backend], 467 attr, 468 )(*args, **kwargs) File ~/miniconda3/envs/dask/lib/python3.9/site-packages/dask/array/random.py:423, in RandomState.uniform(self, low, high, size, chunks, **kwargs) 421 @derived_from(np.random.RandomState, skipblocks=1) 422 def uniform(self, low=0.0, high=1.0, size=None, chunks="auto", **kwargs): ---> 423 return self._wrap("uniform", low, high, size=size, chunks=chunks, **kwargs) File ~/miniconda3/envs/dask/lib/python3.9/site-packages/dask/array/random.py:170, in RandomState._wrap(self, funcname, size, chunks, extra_chunks, *args, **kwargs) 165 kwrg[k] = (getitem, lookup[k], slc) 166 vals.append( 167 (_apply_random, self._RandomState, funcname, seed, size, arg, kwrg) 168 ) ---> 170 meta = _apply_random( 171 self._RandomState, 172 funcname, 173 seed, 174 (0,) * len(size), 175 small_args, 176 small_kwargs, 177 ) 179 dsk.update(dict(zip(keys, vals))) 181 graph = HighLevelGraph.from_collections(name, dsk, dependencies=dependencies) File ~/miniconda3/envs/dask/lib/python3.9/site-packages/dask/array/random.py:453, in _apply_random(RandomState, funcname, state_data, size, args, kwargs) 451 state = RandomState(state_data) 452 func = getattr(state, funcname) ---> 453 return func(*args, size=size, **kwargs) File mtrand.pyx:1134, in numpy.random.mtrand.RandomState.uniform() File _common.pyx:600, in numpy.random._common.cont() File _common.pyx:517, in numpy.random._common.cont_broadcast_2() File __init__.pxd:741, in numpy.PyArray_MultiIterNew3() ValueError: shape mismatch: objects cannot be broadcast to a single shape. Mismatch is between arg 0 with shape (0,) and arg 1 with shape (6249365,).
问题原因与解决方案
原因分析
da.random.uniform和np.random.uniform确实存在关键差异:
- NumPy的
uniform允许low和high传入与输出size匹配的数组,实现逐元素的范围控制。 - Dask的
random.uniform则要求low和high是标量值,不能传入Dask Series(即DataFrame的列),因为它无法直接将数组参数与随机生成逻辑分块对齐,导致形状不匹配报错。
解决方案
要实现Dask DataFrame每行自定义范围的随机数,需要用map_partitions对每个分区单独处理,在分区内部用NumPy生成对应范围的随机数:
import numpy as np def add_random_col(df_part, noise_level): # 对每个分区,用intensity列计算每行的随机范围 low = -noise_level / df_part['intensity'] high = noise_level / df_part['intensity'] df_part['separation_dimension_1'] = np.random.uniform(low=low, high=high, size=len(df_part)) return df_part # 应用到整个Dask DataFrame df = df.map_partitions(add_random_col, noise_level=noise_level)
如果需要保证随机数的可复现性,可以在map_partitions中传入分区相关的种子:
def add_random_col_with_seed(df_part, noise_level, partition_seed): np.random.seed(partition_seed) low = -noise_level / df_part['intensity'] high = noise_level / df_part['intensity'] df_part['separation_dimension_1'] = np.random.uniform(low=low, high=high, size=len(df_part)) return df_part # 给每个分区分配唯一种子 df = df.map_partitions(add_random_col_with_seed, noise_level=noise_level, partition_seed=np.arange(df.npartitions))
内容的提问来源于stack exchange,提问作者Soerendip
相关产品推荐
相关产品推荐

