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

如何用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等方法均无效,以下是针对性解决方案。


可复现代码步骤

  1. 下载数据集片段:
!wget https://lfs.aminer.cn/misc/moocdata/data/mooccube2/relations/user-video.json
  1. 导入依赖:
import dask.dataframe as dd
import pandas as pd
import numpy as np
  1. 特征提取函数:
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
  1. 可行的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]
  1. 报错的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 01:50:21