基于可变评分量表对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
相关产品推荐
相关产品推荐

