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

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

代码说明

  1. 计算列总和:筛选所有以Answer_开头的列,对每列求和,得到包含所有选项总和的单行DataFrame。
  2. 宽表转长表:stack()函数需要指定行数为答题列的总数,每对参数对应一个列名和其求和后的值,将宽表转为便于分组的长表结构。
  3. 解析列名:通过正则表达式从列名中提取选项标识和题目编号,并将题目编号转为整数类型。
  4. 透视聚合:按题目编号分组,用pivot()将选项转为列,计算每个题目下各选项的总被选次数,最后用fillna(0)处理无人选择的选项。

内容的提问来源于stack exchange,提问作者T.Kirtley

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 18:55:23