PySpark多选题各选项按题统计求和的问题及解决方案
问题:统计多选题各选项被选人数总和
需求说明
需要统计多选题中每个题目下各选项的被选人数总和。
现有数据格式(示例:3位受访者、100道题)
+---------+---------+ --------------+--------------+ --------------+--------------+ | misc_1 | misc_2 ... Answer_A_1 | Answer_A_2 ... Answer_D_99 | Answer_D_100 | +---------+---------+ --------------+--------------+ --------------+--------------+ | James| 2345 ... 0 1 ... 0 | 1 | | Anna| 5434 ... 1 0 ... 0 | 1 | | Robert| 7890 ... 0 1 ... 1 | 0 | +---------+---------+ --------------+--------------+ --------------+--------------+
期望输出
按题目编号分组,显示各选项被选次数的DataFrame:
+---+---+---+---+----------+ | A | B | C | D | Question +---+---+---+---+----------+ | 1 | 0 | 1 | 1 | 1 | | 2 | 1 | 0 | 1 | 2 | | 0 | 3 | 0 | 0 | 3 | : : : : : : : : : : | 1 | 0 | 0 | 2 | 100 | +---+---+---+---+----------+
尝试代码及错误
尝试编写的PySpark代码:
from pyspark.sql import SparkSession, functions as F def getSums(df): choices = ['A', 'B', 'C', 'D'] arg = {} answers = [column for column in df.columns if column.startswith("Ans")] for a in answers: arg[a] = 'sum' sums = sums.select(*(F.col(i).alias(i.replace("(",'_').replace(')','')) for i in sums.columns)) sums = df.agg(arg).withColumn('idx', F.lit(None)) s = [f",'{l}'"+f",{column}" for column in sums.columns for l in choices if f"_{l}_" in column] unpivotExpr = "stack(4"+''.join(map(str,s))+") as (A,B,C,D)" unpivotDF = sums.select('idx',F.expr(unpivotExpr)) result = unpivotDF return result
执行时出现错误:
AnalysisException: The number of aliases supplied in the AS clause does not match the number of columns output by the UDTF expected 200 aliases but got A,B,C,D
错误原因是误解了stack()函数的工作原理,现寻求不使用pyspark.pandas的替代解决方案。
解决方案
正确的思路是先对所有答题列求和,再拆分列名提取选项和题目编号,最后按题目编号聚合透视:
from pyspark.sql import SparkSession, functions as F def get_multiple_choice_sums(df): # 1. 筛选答题列,计算每列的总和 answer_cols = [col for col in df.columns if col.startswith("Answer_")] sum_exprs = [F.sum(col).alias(col) for col in answer_cols] sums_df = df.select(sum_exprs).limit(1) # 聚合后仅一行数据 # 2. 宽表转长表:将所有答题列转为(列名、总和)的结构 unpivot_expr = f"stack({len(answer_cols)}, " + \ ", ".join([f"'{col}', {col}" for col in answer_cols]) + \ ") as (col_name, total)" long_df = sums_df.select(F.expr(unpivot_expr)) # 3. 从列名提取选项(A/B/C/D)和题目编号 parsed_df = long_df.withColumn( "option", F.regexp_extract("col_name", r"Answer_([A-Z])_\d+", 1) ).withColumn( "Question", F.regexp_extract("col_name", r"Answer_[A-Z]_(\d+)", 1).cast("int") ) # 4. 按题目分组,透视得到各选项的总和,填充未被选择的选项为0 result_df = parsed_df.groupBy("Question").pivot("option").sum("total").fillna(0) return result_df
代码说明
- 计算列总和:筛选所有以
Answer_开头的列,对每列求和,得到包含所有选项总和的单行DataFrame。 - 宽表转长表:
stack()函数需要指定行数为答题列的总数,每对参数对应一个列名和其求和后的值,将宽表转为便于分组的长表结构。 - 解析列名:通过正则表达式从列名中提取选项标识和题目编号,并将题目编号转为整数类型。
- 透视聚合:按题目编号分组,用
pivot()将选项转为列,计算每个题目下各选项的总被选次数,最后用fillna(0)处理无人选择的选项。
内容的提问来源于stack exchange,提问作者T.Kirtley
相关产品推荐
相关产品推荐

