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
相关产品推荐
相关产品推荐

