如何从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
相关产品推荐
相关产品推荐

