如何将PySpark Row列表转换为指定结构的DataFrame?
问题:将PySpark Row嵌套列表转换为结构化DataFrame
原始数据格式
[[Row(database_description_item='Catalog Name', database_description_value='spark_catalog'), Row(database_description_item='Namespace Name', database_description_value='dummydb'), Row(database_description_item='Comment', database_description_value=''), Row(database_description_item='Location', database_description_value='physical/location/of/database/'), Row(database_description_item='Owner', database_description_value='username')], [Row(database_description_item='Catalog Name', database_description_value='spark_catalog'), Row(database_description_item='Namespace Name', database_description_value='dummydb2'), Row(database_description_item='Comment', database_description_value=''), Row(database_description_item='Location', database_description_value='physical/location/of/database/'), Row(database_description_item='Owner', database_description_value='username')], [Row(database_description_item='Catalog Name', database_description_value='spark_catalog'), Row(database_description_item='Namespace Name', database_description_value='dummydb3'), Row(database_description_item='Comment', database_description_value='some description here'), Row(database_description_item='Location', database_description_value='physical/location/of/database/'), Row(database_description_item='Owner', database_description_value='username')]
需求
将上述数据转换为包含Catalog_Name、Namespace_Name、Comment、Location、Owner列的结构化DataFrame。
原尝试代码及问题
原代码将整个子列表作为Row的单个元素,导致转换后每个单元格都是列表,无法得到预期的结构化结果:
from pyspark.sql import Row rows=[] for r in l: row = Row(r) rows.append(row) rdd = spark.sparkContext.parallelize(rows) df = rdd.toDF().toPandas() df3 = pd.DataFrame(df['_1'].to_list(), columns=['Catalog_Name','Namespace_Name','Comment','Location','Owner'])
解决方案
方法一:纯PySpark实现(推荐,适合大数据场景)
直接提取每个子列表中对应字段的值,构建目标结构的Row后创建DataFrame:
from pyspark.sql import Row processed_rows = [] for sub_list in l: # 匹配每个字段对应的value,构建字典 row_data = { 'Catalog_Name': next(r.database_description_value for r in sub_list if r.database_description_item == 'Catalog Name'), 'Namespace_Name': next(r.database_description_value for r in sub_list if r.database_description_item == 'Namespace Name'), 'Comment': next(r.database_description_value for r in sub_list if r.database_description_item == 'Comment'), 'Location': next(r.database_description_value for r in sub_list if r.database_description_item == 'Location'), 'Owner': next(r.database_description_value for r in sub_list if r.database_description_item == 'Owner') } processed_rows.append(Row(**row_data)) # 创建PySpark DataFrame df = spark.createDataFrame(processed_rows) df.show()
方法二:修复Pandas转换方法
先提取每个子列表中Row的database_description_value值,再传入Pandas创建DataFrame:
import pandas as pd # 提取所有子列表的value值 structured_data = [] for sub_list in l: row_values = [item.database_description_value for item in sub_list] structured_data.append(row_values) # 创建Pandas DataFrame df = pd.DataFrame(structured_data, columns=['Catalog_Name','Namespace_Name','Comment','Location','Owner']) print(df)
内容的提问来源于stack exchange,提问作者Pysparker
相关产品推荐
相关产品推荐

