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

如何从PySpark DataFrame构建保留索引的用户稀疏特征向量

PySpark构建用户-日期维度的稀疏特征向量

问题背景

现有如下结构的PySpark DataFrame(注意:原数据存在重复列名3mon_views,需先重命名处理):

+--------------------+-------+----------+----------+----------+----------+--------+
|             user_id|game_id|3mon_views|3mon_carts|3mon_trans|3mon_views|      dt|
+--------------------+-------+----------+----------+----------+----------+--------+
|0006e38c8968431f8...|0418034|       1.0|       0.0|       0.0|       0.0|20230813|
|0006e38c8968431f8...|0501080|       0.0|       1.0|       0.0|       0.0|20230813|
|0006e38c8968431f8...|0601010|       3.0|       0.0|       0.0|       0.0|20230813|
|0006e38c8968431f8...|0602002|       0.0|       2.0|       0.0|       0.0|20230813|
|0006e38c8968431f8...|0603006|       0.0|       0.0|       5.0|       0.0|20230813|
|0006e38c8968431f8...|0605004|       0.0|       0.0|       0.0|       1.0|20230813|
|0006e38c8968431f8...|0608002|       0.0|       0.0|       0.0|       2.0|20230813|
|0006e38c8968431f8...|0608006|       0.0|       0.0|       2.0|       0.0|20230813|
|0006e38c8968431f8...|0608007|       0.0|       0.0|       0.0|       4.0|20230813|
|0006e38c8968431f8...|0611004|       0.0|       1.0|       0.0|       0.0|20230813|
|0006e38c8968431f8...|0614001|       0.0|       0.0|       0.0|       1.0|20230813|
|0006e38c8968431f8...|0614008|       0.0|       0.0|       0.0|       2.0|20230813|
|0006e38c8968431f8...|0615007|       9.0|       0.0|       0.0|       0.0|20230813|
|0006e38c8968431f8...|1101004|      10.0|       0.0|      15.0|       0.0|20230813|
|0006e38c8968431f8...|1101007|       0.0|       0.0|       5.0|       3.0|20230813|
+--------------------+-------+----------+----------+----------+----------+--------+

需求:按user_id和dt分组,构建形状为(4*1000,)的稀疏向量,每个索引对应特定game_id与特征的组合(存在则取对应值,否则为0),最终得到包含user_id、sparsefeat_vec、dt的结果DataFrame。

实现步骤

1. 预处理:重命名重复列

原DataFrame有两个3mon_views列,先重命名避免冲突:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, array, explode, struct, collect_list, udf, lit
from pyspark.ml.linalg import SparseVector

# 初始化SparkSession
spark = SparkSession.builder.appName("SparseFeatureVector").getOrCreate()

# 重命名重复列
df = df.withColumnRenamed("3mon_views", "3mon_views_1") \
       .withColumnRenamed("3mon_views", "3mon_views_2")

2. 定义特征规则并生成game_id索引映射

指定4个特征的顺序,给每个唯一game_id分配基础索引:

# 定义特征顺序,对应后续的索引偏移
feature_list = ["3mon_views_1", "3mon_carts", "3mon_trans", "3mon_views_2"]
feature_count = len(feature_list)
total_dim = 1000 * feature_count  # 总维度:4*1000

# 获取所有唯一game_id并分配基础索引
game_id_index = df.select("game_id").distinct().rdd \
                  .zipWithIndex().toDF(["game_id_struct", "base_idx"]) \
                  .select(col("game_id_struct.game_id").alias("game_id"), col("base_idx"))

3. 宽表转长表并关联索引

把每个特征列拆分为单独行,同时关联对应的game_id基础索引:

# 将宽表转长表:每条记录拆分为4条(每个特征一行)
long_df = df.select("user_id", "game_id", "dt", 
                    explode(array(*[struct(lit(f).alias("feature"), col(f).alias("value")) for f in feature_list])).alias("feat_struct")) \
            .select("user_id", "game_id", "dt", "feat_struct.feature", "feat_struct.value")

# 关联game_id的基础索引
long_df = long_df.join(game_id_index, on="game_id", how="left")

4. 计算全局稀疏向量索引

根据特征在列表中的位置,计算每个(game_id+feature)组合的全局索引:

# 给每个特征分配偏移量
feature_index_map = {f: i for i, f in enumerate(feature_list)}
feature_index_udf = udf(lambda f: feature_index_map[f])

# 计算全局索引:base_idx * 特征数 + 特征偏移量
long_df = long_df.withColumn("feature_idx", feature_index_udf(col("feature"))) \
                 .withColumn("global_idx", col("base_idx") * feature_count + col("feature_idx"))

5. 分组构建稀疏向量

收集每个用户-日期下的有效(索引,值)对,过滤0值后构建稀疏向量:

# 定义构建稀疏向量的UDF
def build_sparse_vector(indices_values, total_dim):
    indices = []
    values = []
    for idx, val in indices_values:
        if val != 0.0:
            indices.append(idx)
            values.append(val)
    return SparseVector(total_dim, indices, values)

sparse_vec_udf = udf(lambda x: build_sparse_vector(x, total_dim))

# 分组生成结果
result_df = long_df.groupBy("user_id", "dt") \
                   .agg(collect_list(struct("global_idx", "value")).alias("indices_values")) \
                   .withColumn("sparsefeat_vec", sparse_vec_udf(col("indices_values"))) \
                   .select("user_id", "sparsefeat_vec", "dt")

额外说明

  • 若提前已知全部1000个game_id,可直接硬编码索引映射,省去通过RDD生成索引的步骤
  • 过滤值为0的记录能有效减少稀疏向量的存储冗余
  • 最终的sparsefeat_vec列类型为SparseVector,可直接用于后续机器学习任务

内容的提问来源于stack exchange,提问作者Chris_007

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 20:25:04