如何用Spark DataFrame实现动态调查数据转置?
基于Spark DataFrame的调查数据转换方案
核心思路概述
由于源数据是动态宽表(1000+问题列),核心步骤是先将宽表转成长表结构,再关联元数据表补充问题描述、可选答案等信息,最后按目标格式整理输出。
具体实现步骤
1. 加载源数据与元数据表
先读取用户调查回答表(记为survey_df)和元数据表(记为metadata_df),假设调查表里有固定标识列(比如user_id),其余都是动态问题列(列名对应元数据表的问题字段)。
// Scala示例,Python逻辑一致 val survey_df = spark.read.format("csv").option("header", "true").load("path/to/survey_data.csv") val metadata_df = spark.read.format("csv").option("header", "true").load("path/to/metadata.csv")
2. 动态宽表转长表
因为问题列不固定,需要先提取所有问题列(排除user_id等固定列),然后用stack函数将宽表转为长表,每行对应一个用户的一个问题回答。
// 获取所有问题列(排除固定标识列) val question_cols = survey_df.columns.filter(_ != "user_id") // 构造stack表达式,将每个问题列转为(question_id, answer)行 val stack_expr = question_cols.map(col => s"'${col}', `$col`").mkString(", ") val long_df = survey_df.selectExpr("user_id", s"stack(${question_cols.size}, $stack_expr) as (question_id, user_answer)")
3. 关联元数据表补充信息
将长表与元数据表按question_id关联,获取问题描述、可选答案等字段:
val joined_df = long_df.join(metadata_df, long_df("question_id") === metadata_df("问题"), "left")
4. 处理答案与可选答案格式
根据业务需求整理答案格式,比如如果用户答案是选项编码(如A/B/C),可以关联元数据的可选答案映射为文本;如果是多选答案(用分隔符分隔),可拆分后再处理:
// 示例:假设可选答案是"选项A:描述A|选项B:描述B"格式,拆分并映射用户答案 import org.apache.spark.sql.functions._ val parse_answers = udf((options: String, answer: String) => { val optionMap = options.split("\\|").map(_.split(":")).map(arr => arr(0) -> arr(1)).toMap optionMap.getOrElse(answer, answer) }) val final_df = joined_df.withColumn("user_answer_text", parse_answers(col("可选答案"), col("user_answer")))
5. 整理输出格式
根据目标表格结构,选择需要的字段并调整顺序,最后输出:
final_df.select( col("user_id"), col("描述").alias("问题描述"), col("user_answer_text").alias("用户回答"), col("可选答案").alias("可选答案列表") ).write.format("csv").option("header", "true").save("path/to/output")
关键注意事项
- 动态列适配:通过
columns方法动态获取问题列,避免硬编码,适配问题列的新增或删除。 - 空值处理:关联元数据时用
left join,避免丢失没有元数据的问题;处理用户空回答时需根据业务规则填充默认值。 - 性能优化:1000+列转长表会生成大量数据,可通过分区(如按
user_id分区)提升处理效率;元数据表可缓存减少重复读取。
内容的提问来源于stack exchange,提问作者Hemanth Kumar
相关产品推荐
相关产品推荐

