使用Dask读取多数据集并处理类别列编码及数据分块的问题
问题解决:Dask处理大尺寸睡眠分期数据集的编码与分块方案
1. 优化Dask数据加载
首先替换低效的逐文件追加方式,直接用Dask读取目录下所有Parquet文件,自动维护分区结构:
import dask.dataframe as dd import numpy as np data_dir = r'/content/drive/MyDrive/uyku_parquet' # 读取目录下所有Parquet文件,自动生成合理分区 data = dd.read_parquet(f"{data_dir}/*.parquet")
2. 解决SleepStaging列编码错误
原错误源于尝试将全局计算的Pandas Series赋值给Dask DataFrame,导致分区对齐失败。改为先获取全局类别映射,再逐分区编码:
# 获取全量数据集的唯一睡眠分期类别 unique_categories = data['SleepStaging'].unique().compute() # 创建类别到整数的映射字典 category_map = {cat: idx for idx, cat in enumerate(unique_categories)} # 对每个分区的SleepStaging列应用映射编码,无需全局对齐 data['SleepStaging'] = data['SleepStaging'].map(category_map, meta=('SleepStaging', int))
3. 实现按类别分块逻辑
定义分区处理函数,对每个分区内的类别分组后切割为6000行的特征块,再合并所有分区结果:
def chunk_partition(df, chunk_size=6000): feature_chunks = [] label_chunks = [] # 提取特征列(排除SleepStaging) feature_cols = [col for col in df.columns if col != 'SleepStaging'] for label, group in df.groupby('SleepStaging'): features = group[feature_cols].to_numpy() # 计算可完整切割的块数 num_chunks = features.shape[0] // chunk_size for i in range(num_chunks): start = i * chunk_size end = start + chunk_size feature_chunks.append(features[start:end, :]) label_chunks.append(label) return np.array(feature_chunks), np.array(label_chunks) # 对每个分区应用分块函数 partition_results = data.map_partitions(chunk_partition, meta=('object', 'int')) # 计算所有分区结果并合并为最终训练数据 X_parts, y_parts = zip(*partition_results.compute()) X = np.concatenate(X_parts, axis=0) y = np.concatenate(y_parts, axis=0)
关键说明
- 数据加载:
dd.read_parquet自动处理分区,避免手动追加导致的分区混乱 - 编码逻辑:基于全局类别映射逐分区处理,既保证编码一致性,又避免内存溢出
- 分块逻辑:每个分区独立处理后合并,大幅降低单步内存占用,适配百万行级数据集
内容的提问来源于stack exchange,提问作者AbuMuhandisAlTurki
相关产品推荐
相关产品推荐

