You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 10:45:44