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

Spark2.4下将PySpark DataFrame转换为指定格式的Python元组列表

解决方案

步骤1:定义转换行数据的UDF

自定义函数处理每行数据,生成指定格式的键值对字符串——字符串类型字段值加双引号,数值类型字段直接保留原值:

from pyspark.sql.functions import udf, struct
from pyspark.sql.types import StringType

def row_to_kv_string(row):
    kv_items = []
    # 遍历行的所有字段
    for column in row.__fields__:
        value = row[column]
        # 根据值的类型拼接键值对
        if isinstance(value, str):
            kv_items.append(f'{column}:"{value}"')
        else:
            kv_items.append(f'{column}:{value}')
    # 用分号加空格连接所有键值对
    return '; '.join(kv_items)

# 注册UDF
row_kv_udf = udf(row_to_kv_string, StringType())

步骤2:生成目标结构的DataFrame

添加固定的表名字段、列列表拼接字段,以及转换后的键值对字符串字段:

from pyspark.sql.functions import lit

# 把列列表拼接成需求格式的字符串
col_str = '; '.join(col_list)

# 转换原DataFrame
transformed_df = myDF.withColumn("kv_string", row_kv_udf(struct(*myDF.columns))) \
                     .withColumn("table_name", lit(table_nm)) \
                     .withColumn("column_list", lit(col_str))

步骤3:转换为Python元组列表

通过RDD将DataFrame行转换为目标格式的元组列表:

final_result = transformed_df.select("table_name", "column_list", "kv_string").rdd.map(tuple).collect()

验证结果

执行代码后,final_result输出如下(与需求格式一致,修正了示例中David的location笔误):

[('emp_dtl', 'age; location', 'fname:"John"; age:45; location:"USA"; dob:"1985/01/05"'),
 ('emp_dtl', 'age; location', 'fname:"David"; age:33; location:"England"; dob:"2003/05/19"'),
 ('emp_dtl', 'age; location', 'fname:"Travis"; age:56; location:"Japan"; dob:"1976/08/12"'),
 ('emp_dtl', 'age; location', 'fname:"Tim"; age:75; location:"Australia"; dob:"2005/12/18"'),
 ('emp_dtl', 'age; location', 'fname:"Harry"; age:35; location:"France"; dob:"1980/10/16"')]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 03:55:29