编写Python函数为指定人员生成featureSk二进制列的Pandas DataFrame
PySpark转Pandas宽表实现及错误排查
问题背景
现有PySpark DataFrame productusage结构如下:
| featureSk | PersonNumber |
|---|---|
| A | 1001 |
| B | 1001 |
| C | 1003 |
| C | 1004 |
| A | 1002 |
| B | 1005 |
需要编写Python函数,输入人员编号列表,输出格式如下的Pandas DataFrame:以featureSk的每个取值作为列,若人员在productusage中存在对应featureSk则值为1,否则为0。示例输出:
| PersonNumber | A | B | C |
|---|---|---|---|
| 1001 | 1 | 1 | 0 |
| 1002 | 1 | 0 | 0 |
| 1003 | 0 | 0 | 1 |
注:修正原示例中1002的A值错误,原数据中1002对应featureSk为A,故值应为1。
错误原因分析
报错“featureSk未定义”的常见触发场景:
- 代码中直接使用
featureSk作为变量名,而非通过PySpark的col("featureSk")或productusage.featureSk引用列对象 - 动态生成列逻辑中,未正确处理
featureSk的字符串值,导致Python将其识别为未定义变量
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, count import pandas as pd def get_feature_matrix(person_list): # 1. 获取所有唯一的featureSk值,避免硬编码 features = [row.featureSk for row in productusage.select("featureSk").distinct().collect()] # 2. 创建输入人员的基础DataFrame,确保所有输入人员都被保留 person_df = spark.createDataFrame(pd.DataFrame({"PersonNumber": person_list})) # 3. 分组透视原表,标记每个人员的feature存在情况 pivot_df = productusage.groupBy("PersonNumber")\ .pivot("featureSk", features)\ .agg(count("featureSk"))\ .withColumns({f: when(col(f).isNotNull(), 1).otherwise(0) for f in features}) # 4. 左连接基础人员表,填充缺失特征值为0 result_df = person_df.join(pivot_df, on="PersonNumber", how="left")\ .fillna(0, subset=features) # 5. 转换为Pandas并调整列顺序 pd_result = result_df.toPandas() pd_result = pd_result[["PersonNumber"] + features] return pd_result # 测试示例 if __name__ == "__main__": spark = SparkSession.builder.appName("FeatureMatrix").getOrCreate() # 模拟构建productusage DataFrame data = [("A", 1001), ("B", 1001), ("C", 1003), ("C", 1004), ("A", 1002), ("B", 1005)] productusage = spark.createDataFrame(data, ["featureSk", "PersonNumber"]) # 测试输入人员列表 input_persons = [1001, 1002, 1003] output = get_feature_matrix(input_persons) print(output)
代码说明
- 动态获取特征:通过
distinct()和collect()自动获取所有特征值,适配数据变化 - 保留输入人员:创建基础人员表,确保输入的所有人员(即使不在原表中)都能出现在结果中
- 行转列透视:用
pivot()实现宽表转换,通过count()标记存在性,再转为1/0格式 - 缺失值填充:左连接后用
fillna()将缺失的特征值补0 - 格式对齐:转换为Pandas后调整列顺序,匹配需求输出结构
内容的提问来源于stack exchange,提问作者Brett
相关产品推荐
相关产品推荐

