PySpark字典转DataFrame:解决单行存储问题,生成对应3行的目标DataFrame
PySpark字典转DataFrame:解决单行存储问题,生成对应3行的目标DataFrame
嗨,我来帮你搞定这个问题!你当前的代码把整个字典当成单个行数据传入createDataFrame,所以才会出现所有值挤在一行的情况。咱们需要把字典转换成「3个独立行数据」的列表,让每个unique_survey_id对应一行,这样Spark就能生成你想要的结果啦。
问题原因分析
你之前的写法spark.createDataFrame([inferenced_df],schema),是把整个字典包装成了一个列表元素,相当于告诉Spark“这是一行数据”。但实际上你的字典里每个键对应的值都是长度为3的序列(列表/嵌套字典列表),我们需要把每个索引位置的各字段值组合起来,形成3个独立的行字典。
修正后的完整代码
from pyspark.sql.types import StructType, StructField, StringType import json def alerts_inference(inferenced_df): # 先修正原字典里的语法错误:verbatim列表第二个元素的引号没闭合 inferenced_df['verbatim'][1] = "I am 23 yrs old," # 定义Schema(和你原来的一致,注意嵌套字典字段用StringType,我们转成JSON字符串存储) schema = StructType([ StructField("unique_survey_id", StringType(), True), StructField("verbatim", StringType(), True), StructField("classification_critical_crc_escalation_fg", StringType(), True), StructField("reason_critical_crc_escalation_fg", StringType(), True), StructField("classification_critical_technical_fg", StringType(), True), StructField("reason_critical_technical_fg", StringType(), True), StructField("classification_critical_process_fg", StringType(), True), StructField("reason_critical_process_fg", StringType(), True), StructField("classification_insight_experience_fg", StringType(), True), StructField("reason_insight_experience_fg", StringType(), True), StructField("classification_insight_process_fg", StringType(), True), StructField("reason_insight_process_fg", StringType(), True) ]) # 生成3行数据的列表:遍历每个索引,组合对应位置的字段值 rows = [] for idx in range(3): row = { "unique_survey_id": inferenced_df['unique_survey_id'][idx], "verbatim": inferenced_df['verbatim'][idx], "classification_critical_crc_escalation_fg": inferenced_df['classification_critical_crc_escalation_fg'][idx], "reason_critical_crc_escalation_fg": inferenced_df['reason_critical_crc_escalation_fg'][idx], "classification_critical_technical_fg": inferenced_df['classification_critical_technical_fg'][idx], "reason_critical_technical_fg": inferenced_df['reason_critical_technical_fg'][idx], "classification_critical_process_fg": str(inferenced_df['classification_critical_process_fg'][idx]), # 嵌套字典转成JSON字符串,符合StringType要求 "reason_critical_process_fg": json.dumps(inferenced_df['reason_critical_process_fg'][idx]), "classification_insight_experience_fg": inferenced_df['classification_insight_experience_fg'][idx], "reason_insight_experience_fg": inferenced_df['reason_insight_experience_fg'][idx], "classification_insight_process_fg": str(inferenced_df['classification_insight_process_fg'][idx]), # 同样处理嵌套字典 "reason_insight_process_fg": json.dumps(inferenced_df['reason_insight_process_fg'][idx]) } rows.append(row) # 用行列表创建DataFrame inferenced_df = spark.createDataFrame(rows, schema) return inferenced_df
关键说明
- 修正语法错误:原字典里
verbatim列表的第二个元素引号没闭合,我已经在代码里帮你修正了,不然运行会报错哦。 - 行数据生成:通过循环索引0、1、2,把每个字段对应位置的值提取出来,组合成单个行字典,最终得到3个行字典的列表。
- 嵌套字典处理:对于
reason_critical_process_fg和reason_insight_process_fg这类嵌套字典字段,我们用json.dumps()转换成JSON字符串,这样就能符合你定义的StringType字段要求;如果之后需要解析这些嵌套内容,也可以用from_json()函数再转成StructType。 - 数字转字符串:原字典里
classification_critical_process_fg和classification_insight_process_fg的值是数字(0),我们用str()转成字符串,匹配你定义的StringType。
验证结果
运行这个函数后,生成的DataFrame会有3行,每行对应一个unique_survey_id,所有字段的值都正确对应,完全符合你的预期输出。
备注:内容来源于stack exchange,提问作者Zikra Noman
相关产品推荐
相关产品推荐

