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

基于可变评分量表对Spark DataFrame问卷答案列进行分组

解决方案

要实现按问题对应的量表对评分分组的逻辑,核心是先明确每个问题的完整评分范围,再基于范围将最低2个评分标记为poor、最高2个标记为good,剩余标记为average。以下是基于PySpark的实现步骤:

步骤1:准备原始数据与量表配置

首先创建原始问卷DataFrame,并定义每个问题对应的量表最大值(假设量表从1开始到对应最大值):

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, ArrayType

spark = SparkSession.builder.appName("QuestionnaireScoring").getOrCreate()

# 原始问卷数据
data = [
    ("question_1", 1),
    ("question_1", 4),
    ("question_1", 6),
    ("question_2", 2),
    ("question_2", 4)
]

schema = StructType([
    StructField("question", StringType(), True),
    StructField("answer", IntegerType(), True)
])

df = spark.createDataFrame(data, schema)

# 定义每个问题的量表最大值(例如question_1是6分量表,question_2是5分量表)
scale_config = {
    "question_1": 6,
    "question_2": 5
}
scale_df = spark.createDataFrame(scale_config.items(), ["question", "max_score"])

步骤2:生成完整评分映射表

为每个问题生成完整的评分范围,并按规则标记meaning:

# 生成每个问题的所有评分值
def generate_scores(max_score):
    return [(i,) for i in range(1, max_score + 1)]

generate_scores_udf = F.udf(generate_scores, ArrayType(IntegerType()))

# 展开所有评分值
scale_with_scores = scale_df.withColumn("all_scores", generate_scores_udf(F.col("max_score"))) \
    .withColumn("score", F.explode(F.col("all_scores"))) \
    .drop("all_scores", "max_score") \
    .withColumnRenamed("score", "answer")

# 分组确定每个问题的poor/good评分区间
ranked_scores = scale_with_scores.groupBy("question") \
    .agg(F.collect_list("answer").alias("sorted_scores")) \
    .withColumn("poor_scores", F.slice(F.col("sorted_scores"), 1, 2)) \
    .withColumn("good_scores", F.slice(F.col("sorted_scores"), -2, 2)) \
    .drop("sorted_scores")

# 生成poor和good的映射关系
poor_mapping = ranked_scores.select("question", F.explode("poor_scores").alias("answer"), F.lit("poor").alias("meaning"))
good_mapping = ranked_scores.select("question", F.explode("good_scores").alias("answer"), F.lit("good").alias("meaning"))

# 合并映射并填充average
full_mapping = poor_mapping.union(good_mapping)
full_mapping = scale_with_scores.join(full_mapping, on=["question", "answer"], how="left") \
    .withColumn("meaning", F.when(F.col("meaning").isNull(), "average").otherwise(F.col("meaning")))

步骤3:关联原始数据得到结果

将原始DataFrame与评分映射表关联,最终得到带meaning列的结果:

result_df = df.join(full_mapping, on=["question", "answer"], how="left")
result_df.show()

执行后输出结果:

+-----------+------+-------+
|   question|answer|meaning|
+-----------+------+-------+
|question_1|     1|   poor|
|question_1|     4|average|
|question_1|     6|   good|
|question_2|     2|   poor|
|question_2|     4|   good|
+-----------+------+-------+

逻辑说明

  • 针对每个问题的完整量表范围,提取前2个最小值作为poor,最后2个最大值作为good,中间的所有评分自动标记为average
  • 该方案支持4-6分的任意量表,无需修改核心逻辑,只需调整scale_config中的量表最大值即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 12:12:40