kuairand-27k的Parquet 数据导出与上传到 MaxCompute 完整流程(hstu格式)

简介: 本文详解如何将本地kuairand-27k(1257行×14列)Parquet推荐数据集,经探查、类型映射(int64→bigint,list→array<bigint>),通过pyodps上传至阿里云MaxCompute表,含完整环境配置、建表与批量上传代码。

概述

本文介绍如何将本地 Parquet 文件(kuairand-27k 推荐系统数据集)导出并上传到阿里云 MaxCompute 表,包含:数据探查、类型映射、建表、上传的完整代码。


1. 环境准备

依赖安装

pip install pandas pyarrow pyodps
  • pandas + pyarrow:读取 Parquet 文件
  • pyodps:阿里云 MaxCompute Python SDK

凭证配置

通过环境变量传入 AccessKey,避免硬编码:

export ACCESS_ID="your_access_id"
export ACCESS_KEY="your_access_key"

2. 数据探查:读取 Parquet 文件

在上传之前,先了解 Parquet 文件的数据结构:

import pandas as pd

df = pd.read_parquet("kuairand-27k-train-0.parquet")

print(f"数据总行数: {len(df)}")
print(f"数据列数: {len(df.columns)}")
print(f"\n列名列表:")
print(df.columns.tolist())
print(f"\n数据类型:")
print(df.dtypes)

# 逐列查看前5条数据
for col in df.columns:
    print(f"\n--- {col} (dtype: {df[col].dtype}) ---")
    for i in range(min(5, len(df))):
        val = df[col].iloc[i]
        if isinstance(val, (list,)):
            print(f"  [{i}] len={len(val)}, 前10个值: {val[:10]}")
        else:
            print(f"  [{i}] {val}")

探查结果

本数据集共 1257 行、14 列,结构如下:

列名 Python 类型 说明
user_id int64 用户 ID
user_active_degree int64 用户活跃度
follow_user_num_range int64 关注人数区间
fans_user_num_range int64 粉丝人数区间
friend_user_num_range int64 好友人数区间
register_days_range int64 注册天数区间
video_id list(int) 用户历史交互视频序列
action_timestamp list(int) 行为时间戳序列
action_weight list(int) 行为权重序列(bitmask 编码)
watch_time list(int) 观看时长序列
item_video_id list(int) 候选视频 ID 序列
item_action_weight list(int) 候选视频行为标签
item_target_watchtime list(int) 候选视频目标观看时长
item_query_time list(int) 候选请求时间戳

3. 类型映射:Parquet → MaxCompute

Parquet/Python 类型 MaxCompute 类型
int64(标量) bigint
list(int)(数组) array<bigint>

4. 完整上传脚本

#!/usr/bin/env python3
"""
将 kuairand-27k-train-0.parquet 数据上传到 MaxCompute 表 pairec_kuairand_train
"""

import os
import pandas as pd
import numpy as np
from odps import ODPS
from odps.models import TableSchema as Schema, Column

# ========== 1. 配置连接参数 ==========
project_name = "pairec_mc"
access_id = os.environ["ACCESS_ID"]
access_key = os.environ["ACCESS_KEY"]
endpoint = "http://service.cn.maxcompute.aliyun.com/api"

# ========== 2. 连接 MaxCompute ==========
odps = ODPS(access_id, access_key, project_name, endpoint=endpoint)
print("MaxCompute 连接成功")

# ========== 3. 建表 ==========
TABLE_NAME = "pairec_kuairand_train"

# 先删除已有表(如需覆盖写入)
if odps.exist_table(TABLE_NAME):
    print(f"表 {TABLE_NAME} 已存在,正在删除...")
    odps.delete_table(TABLE_NAME)
    print(f"表 {TABLE_NAME} 已删除")

# 定义表 schema
columns = [
    # 用户侧标量字段 (bigint)
    Column(name="user_id", type="bigint"),
    Column(name="user_active_degree", type="bigint"),
    Column(name="follow_user_num_range", type="bigint"),
    Column(name="fans_user_num_range", type="bigint"),
    Column(name="friend_user_num_range", type="bigint"),
    Column(name="register_days_range", type="bigint"),
    # 历史序列字段 (array<bigint>)
    Column(name="video_id", type="array<bigint>"),
    Column(name="action_timestamp", type="array<bigint>"),
    Column(name="action_weight", type="array<bigint>"),
    Column(name="watch_time", type="array<bigint>"),
    # 候选物料字段 (array<bigint>)
    Column(name="item_video_id", type="array<bigint>"),
    Column(name="item_action_weight", type="array<bigint>"),
    Column(name="item_target_watchtime", type="array<bigint>"),
    Column(name="item_query_time", type="array<bigint>"),
]

schema = Schema(columns=columns)
odps.create_table(TABLE_NAME, schema)
print(f"表 {TABLE_NAME} 创建成功")

# ========== 4. 读取 Parquet 并上传数据 ==========
PARQUET_PATH = "kuairand-27k-train-0.parquet"
df = pd.read_parquet(PARQUET_PATH)
print(f"Parquet 读取完成,共 {len(df)} 行")

# 将 numpy 数组转为 Python list(PyODPS 要求原生 Python 类型)
array_cols = [
    "video_id", "action_timestamp", "action_weight", "watch_time",
    "item_video_id", "item_action_weight", "item_target_watchtime", "item_query_time"
]
for col in array_cols:
    df[col] = df[col].apply(
        lambda x: list(x) if isinstance(x, np.ndarray) else (x if isinstance(x, list) else [])
    )

# 标量列确保为 Python int
scalar_cols = [
    "user_id", "user_active_degree", "follow_user_num_range",
    "fans_user_num_range", "friend_user_num_range", "register_days_range"
]
for col in scalar_cols:
    df[col] = df[col].astype(int)

# 使用 Tunnel 上传
table = odps.get_table(TABLE_NAME)
print(f"开始上传数据(共 {len(df)} 行)...")

with table.open_writer() as writer:
    records = []
    for idx, row in df.iterrows():
        record = [
            int(row["user_id"]),
            int(row["user_active_degree"]),
            int(row["follow_user_num_range"]),
            int(row["fans_user_num_range"]),
            int(row["friend_user_num_range"]),
            int(row["register_days_range"]),
            list(row["video_id"]),
            list(row["action_timestamp"]),
            list(row["action_weight"]),
            list(row["watch_time"]),
            list(row["item_video_id"]),
            list(row["item_action_weight"]),
            list(row["item_target_watchtime"]),
            list(row["item_query_time"]),
        ]
        records.append(table.new_record(record))

        if (idx + 1) % 100 == 0:
            print(f"  已处理 {idx + 1}/{len(df)} 行")

    writer.write(records)

print(f"数据上传完成!共上传 {len(records)} 行到 {TABLE_NAME}")

5. 执行

source env.sh  # 设置 ACCESS_ID、ACCESS_KEY 环境变量
python upload_to_odps.py

预期输出:

project_name: pairec
endpoint: http://service.cn.maxcompute.aliyun.com/api
MaxCompute 连接成功
表 pairec_kuairand_train 创建成功
Parquet 读取完成,共 1257 行
开始上传数据(共 1257 行)...
  已处理 100/1257 行
  ...
  已处理 1200/1257 行
数据上传完成!共上传 1257 行到 pairec_kuairand_train

6. 常见问题与注意事项

Endpoint 与 Project 不匹配

MaxCompute 的 Project 绑定到特定 Region 的 Endpoint。如果报 Project not found,需确认 Project 所在的 Region 并切换对应 Endpoint:

Endpoint 适用场景
service.cn.maxcompute.aliyun.com 公网(华东2)

AccessKey 与 Endpoint 网络域不匹配

公网 AccessKey 只能用于公网 Endpoint,内网 AccessKey 只能用于内网 Endpoint,混用会报 AccessKeyIdNotFound

PyODPS 要求原生 Python 类型

open_writer() 写入数据时,不支持 numpy 的 int64ndarray 类型,需显式转换为 Python 原生的 intlist,否则会抛出类型错误。

批量上传性能

对于大数据量(百万行以上),建议分批写入(如每批 5000 条),避免单次 writer.write() 内存过大:

BATCH_SIZE = 5000
with table.open_writer() as writer:
    batch = []
    for idx, row in df.iterrows():
        batch.append(table.new_record([...]))
        if len(batch) >= BATCH_SIZE:
            writer.write(batch)
            batch = []
    if batch:
        writer.write(batch)
相关实践学习
机器学习概览及常见算法
机器学习(Machine Learning, ML)是人工智能的核心,专门研究计算机怎样模拟或实现人类的学习行为,以获取新的知识或技能,重新组织已有的知识结构使之不断改善自身的性能,它是使计算机具有智能的根本途径,其应用遍及人工智能的各个领域。 本课程将带你入门机器学习,掌握机器学习的概念和常用的算法。
相关文章
|
5月前
|
并行计算 算法框架/工具 iOS开发
TorchRec在macos ARM芯片(Apple Silicon)上无法安装
JaggedTensor等在macOS ARM芯片上无法运行,主因是ARM64与x86_64架构不兼容,且TorchRec深度依赖CUDA——而Apple Silicon仅支持Metal。fbgemm-gpu缺失、Rosetta 2不支持CUDA指令,导致关键操作失败。建议改用MLX框架或标准PyTorch张量替代。
472 4
|
网络安全 开发工具 对象存储
OSS 的C++ SDK编译安装指南
OSS 的C++ SDK编译安装指南
OSS 的C++ SDK编译安装指南
|
2月前
|
人工智能 运维 安全
阿里云通义千问大模型介绍:核心功能、性能优势、行业落地场景与官方定价解析
作为阿里云自主研发的**国产旗舰级通用AI大模型**,通义千问(Qwen)历经多次版本迭代,目前已形成Qwen3全系列模型矩阵,涵盖轻量极速、均衡通用、旗舰推理、多模态创作全梯度能力。依托阿里云万亿级算力底座与海量数据训练,千问大模型在逻辑推理、代码生成、长文本处理、多模态理解、复杂任务拆解等核心维度达到国内顶尖、国际一流水平,同时通过国家大模型标准符合性评测,具备合规性强、稳定性高、国产化适配完善、成本可控等核心优势,是目前国内个人创作、企业数字化、行业智能化改造落地最广泛的通用大模型。本文将全面详解2026版阿里云通义千问大模型的核心功能、差异化性能优势、全行业落地场景、官方最新价格明细
2299 1
|
2月前
|
人工智能 自然语言处理 数据可视化
2026最新AI工作流平台实测指南:场景落地与选型参考
本文深度解析2026年AI工作流平台如何破解办公割裂困局:聚焦任务自动化、会议落地、数据报告、知识沉淀、跨部门协同五大核心场景,对比Dify、阿里云百炼等主流平台特性,并提供分角色选型建议与实用落地经验。(239字)
|
3月前
|
分布式计算 Shell MaxCompute
阿里云 PAI-DLC PyTorchJob 任务提交参数的介绍
本文详解PAI-DLC中`dlc submit pytorchjob`命令的两类核心参数:DLC平台控制参数(如`--name`、`--data_sources`、`--priority`等,用于定义任务属性与资源)和Command执行指令(含环境安装、`torchrun`分布式训练、模型导出等Shell逻辑),并强调关键注意事项。
391 1
|
3月前
|
SQL 人工智能 安全
别再让 AI 温柔地夸你的烂代码了:Code Review 提示词该这样写
AI代码审查不能只求“温柔”,而要像资深工程师一样犀利。本文揭示:模糊提示=无效审查,必须用高精度角色锚点(如Google Staff Engineer)、硬性约束(P0-P3风险分级、可运行重构代码)和结构化输出,让AI真正成为生产级审查助手。提示词,已是工程规范新一环。
688 0
别再让 AI 温柔地夸你的烂代码了:Code Review 提示词该这样写
|
3月前
|
机器学习/深度学习 分布式计算 搜索推荐
推荐系统中的主要陷阱
本文剖析推荐系统六大核心陷阱:线上线下特征/数据不一致、评估指标失真、探索与利用两难、算法精准度与体验矛盾、工程实现漏洞(代码/特征穿越/收敛问题),以及目标模糊的系统性挑战。附阿里PAI-Rec等实战工具方案。(239字)
327 0
|
4月前
|
搜索推荐
PAI-Rec 多路召回截断实践:用 PriorityAdjustCountFilter 和 SnakeFilter 控制精排入口数量
PAI-Rec推荐开发平台提供PriorityAdjustCountFilter(按优先级截取)与SnakeFilter(按权重蛇形交错)两种多路召回截断策略,无需粗排即可将数百候选精准压缩至200个以内进入精排,兼顾保量性、多样性与业务可控性。
344 0
|
8月前
|
人工智能 自然语言处理 安全
企业级 AI API 接入架构实践:从多密钥混乱到统一治理的工程解法
本文探讨企业级AI API接入的治理难题,指出密钥管理混乱、多模型/多团队协同难等痛点,提出“统一入口+分组策略+策略层解耦”的可演进架构,助力AI能力从实验走向稳定、安全、可持续的基础设施。
|
机器学习/深度学习 人工智能 自然语言处理
构建企业级数据分析助手:Data Agent 开发实践
本篇将介绍DMS的一款数据分析智能体(Data Agent for Analytics )产品的技术思考和实践。Data Agent for Analytics 定位为一款企业级数据分析智能体, 基于Agentic AI 技术,帮助用户查数据、做分析、生成报告、深入洞察。由于不同产品的演进路径,背景都不一样,所以只介绍最核心的部分,来深入剖析如何构建企业级数据分析助手:能力边界定义,技术内核,企业级能力。希望既能作为Data Agent for Analytics产品的技术核心介绍,也能作为读者的开发实践的参考。
3069 3
构建企业级数据分析助手:Data Agent 开发实践