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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 10:05:31