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

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

关键说明

  1. 修正语法错误:原字典里verbatim列表的第二个元素引号没闭合,我已经在代码里帮你修正了,不然运行会报错哦。
  2. 行数据生成:通过循环索引0、1、2,把每个字段对应位置的值提取出来,组合成单个行字典,最终得到3个行字典的列表。
  3. 嵌套字典处理:对于reason_critical_process_fg和reason_insight_process_fg这类嵌套字典字段,我们用json.dumps()转换成JSON字符串,这样就能符合你定义的StringType字段要求;如果之后需要解析这些嵌套内容,也可以用from_json()函数再转成StructType。
  4. 数字转字符串:原字典里classification_critical_process_fg和classification_insight_process_fg的值是数字(0),我们用str()转成字符串,匹配你定义的StringType。

验证结果

运行这个函数后,生成的DataFrame会有3行,每行对应一个unique_survey_id,所有字段的值都正确对应,完全符合你的预期输出。

备注:内容来源于stack exchange,提问作者Zikra Noman

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:25:29