如何用Dask高效聚类DataFrame中存储NumPy数组的列?
解决Dask中Object类型坐标列转2D分块数组的方案
核心思路
要将存储为object类型的(2,) NumPy数组列,转换为Dask ML KMeans可接受的分块2D Dask Array,关键是并行处理每个DataFrame分区,避免单核心串行操作。
方法1:分区转换+延迟拼接(通用方案)
针对已经读取到Dask DataFrame的object列,通过分区级别的转换生成2D数组,再拼接成分块Dask Array:
- 定义分区转换函数
import numpy as np def convert_partition(series): # 将当前分区的object列转为(行数, 2)的2D NumPy数组 return np.vstack(series.values)
- 并行处理分区并拼接为Dask Array
import dask.array as da from dask.dataframe import map_partitions # 处理每个DataFrame分区,返回每个分区的2D数组(延迟对象) partition_arrays = map_partitions( convert_partition, df["coordinates"], meta=np.zeros((0, 2), dtype=np.float64) # 指定返回结果的元数据 ) # 将延迟对象转为Dask Array分区,再拼接成完整的2D分块数组 X = da.concatenate([ da.from_delayed(p, shape=(None, 2), dtype=np.float64) for p in partition_arrays.to_delayed() ])
方法2:读取时直接展开坐标列(更高效)
如果Parquet文件的坐标列是结构化嵌套格式,可在读取阶段直接拆分为独立列,再转为2D数组:
import dask.dataframe as dd # 读取Parquet文件 df = dd.read_parquet("your_data.parquet", columns=["coordinates"]) # 将每个坐标数组拆分为x、y两列 df["x"] = df["coordinates"].apply(lambda arr: arr[0], meta=('x', float)) df["y"] = df["coordinates"].apply(lambda arr: arr[1], meta=('y', float)) # 直接转为分块2D Dask Array X = df[["x", "y"]].to_dask_array(lengths=True)
后续聚类使用
得到2D分块数组X后,即可直接用于Dask ML的KMeans:
from dask_ml.cluster import KMeans kmeans = KMeans(n_clusters=5, random_state=42) kmeans.fit(X)
内容的提问来源于stack exchange,提问作者drfk
相关产品推荐
相关产品推荐

