如何用Dask提取MOOCCube数据集的用户视频交互特征?
Dask DataFrame Apply 元数据不匹配报错解决方案
从MOOCCube数据集提取用户视频交互特征时,Pandas在500行小数据集上运行正常,但使用Dask DataFrame的apply(axis=1)时触发报错:
ValueError: The columns in the computed data do not match the columns in the provided metadata Order of columns does not match
尝试过添加index列到元数据、设置元数据为输入列、包装函数返回Dask DataFrame等方法均无效,以下是针对性解决方案。
可复现代码步骤
- 下载数据集片段:
!wget https://lfs.aminer.cn/misc/moocdata/data/mooccube2/relations/user-video.json
- 导入依赖:
import dask.dataframe as dd import pandas as pd import numpy as np
- 特征提取函数:
def segment_feature_extraction(x: dict): segments = x['segment'] num_of_pauses = len(segments) - 1 num_of_skips = sum([1 for i in range(1, len(segments)) if segments[i]['start_point'] > segments[i-1]['end_point']]) total_watched_duration = sum(seg['end_point'] - seg['start_point'] for seg in segments) total_weighted_speed = sum((seg['end_point'] - seg['start_point']) * seg['speed'] for seg in segments) mean_watching_speed = total_weighted_speed / total_watched_duration if total_watched_duration != 0 else 0 total_watching_time = total_watched_duration total_watching_time_with_speed = sum((seg['end_point'] - seg['start_point']) / seg['speed'] for seg in segments) features = { 'video_id': x['video_id'], 'num_of_pauses': num_of_pauses, 'num_of_skips': num_of_skips, 'mean_watching_speed': mean_watching_speed, 'total_watching_time': total_watching_time, 'total_watching_time_with_speed': total_watching_time_with_speed } return pd.Series(features) def full_feature_extraction(x: list): data = pd.DataFrame(columns=['video_id', 'num_of_pauses', 'num_of_skips', 'mean_watching_speed', 'total_watching_time', 'total_watching_time_with_speed']) for vid in x: features = segment_feature_extraction(vid) features = pd.DataFrame([features]) data = pd.concat([data, features], axis=0, ignore_index=True) return data
- 可行的Pandas测试代码:
def pandas_extracting(dataframe: pd.DataFrame): extraction = dataframe.apply(lambda x: full_feature_extraction(x['seq']), axis=1) result = pd.DataFrame(columns=['video_id', 'num_of_pauses', 'num_of_skips', 'mean_watching_speed', 'total_watching_time', 'total_watching_time_with_speed']) for extracts in extraction: result = pd.concat([result, extracts], axis=0, ignore_index=True) return result df = pd.read_json('user-video.json', nrows=500, lines=True) test = pandas_extracting(df) test[0]
- 报错的Dask代码:
def dask_extracting(dataframe): extraction_metadata=[ ('video_id', 'str'), ('num_of_pauses', 'int64'), ('num_of_skip', 'int64'), ('mean_watching_speed', 'float'), ('total_watching_time', 'float'), ('total_watching_time_with_speed', 'float') ] extraction = dataframe.apply(lambda x: full_feature_extraction(x['seq']), meta=extraction_metadata, axis=1) return extraction ddf = dd.read_json('user-video.json', lines=True, blocksize=block_size, meta=[('user_id', 'str'), ('seq', 'object')]) test = dask_extracting(ddf) test.loc[0].compute() # 测试是否可行
解决方案
1. 修复元数据列名错误
报错代码的元数据中列名num_of_skip少了末尾的s,与full_feature_extraction返回的num_of_skips不匹配,先修正该错误:
extraction_metadata=[ ('video_id', 'str'), ('num_of_pauses', 'int64'), ('num_of_skips', 'int64'), # 修正列名 ('mean_watching_speed', 'float'), ('total_watching_time', 'float'), ('total_watching_time_with_speed', 'float') ]
2. 调整函数返回格式适配Dask
Dask的apply(axis=1)默认期望每行返回单行结构(如Series),但原full_feature_extraction返回多行DataFrame,导致元数据匹配失败。推荐改为返回特征列表,再通过explode展开:
def full_feature_extraction(x: list): features_list = [] for vid in x: features = segment_feature_extraction(vid) features_list.append(features.to_dict()) return features_list def dask_extracting(dataframe): # 先获取嵌套特征列表 extraction = dataframe.apply(lambda x: full_feature_extraction(x['seq']), meta=object, axis=1) # 展开嵌套列表为多行 exploded = extraction.explode() # 将字典转为结构化DataFrame并指定元数据 result = exploded.apply(pd.Series, meta={ 'video_id': 'str', 'num_of_pauses': 'int64', 'num_of_skips': 'int64', 'mean_watching_speed': 'float', 'total_watching_time': 'float', 'total_watching_time_with_speed': 'float' }) return result
3. 更高效的分区处理方案
避免在apply中循环拼接DataFrame,改用map_partitions直接处理每个分区的Pandas数据,提升性能:
def process_partition(df: pd.DataFrame): all_features = [] for _, row in df.iterrows(): for vid in row['seq']: features = segment_feature_extraction(vid) all_features.append(features) return pd.DataFrame(all_features) def dask_extracting(dataframe): # 用小样本获取正确元数据(可选) sample_df = pd.read_json('user-video.json', nrows=10, lines=True) sample_meta = process_partition(sample_df).dtypes.to_dict() result = dataframe.map_partitions(process_partition, meta=sample_meta) return result
4. 验证元数据一致性
确保元数据的列名、顺序、数据类型与函数返回结果完全一致。可通过小样本Pandas运行结果直接提取dtypes作为Dask元数据,避免手动编写出错。
内容的提问来源于stack exchange,提问作者Nguyễn Tiến Đạt
相关产品推荐
相关产品推荐

