【问题标题】:Creating a pyspark dataframe from a python dictionary从 python 字典创建 pyspark 数据框
【发布时间】:2021-09-08 00:03:18
【问题描述】:

我想从 python 字典创建一个 pyspark 数据框,但是下面的代码

from pyspark.sql import SparkSession, Row

df_stable = spark.createDataFrame(dict_stable_feature)
df_stable.show()

显示此错误

TypeError: Can not infer schema for type: <class 'str'>

在 stackoverflow 上阅读这篇文章:

Pyspark: Unable to turn RDD into DataFrame due to data type str instead of StringType

我可以推断出问题可能是我错误地使用了 python 标准 str 而不是 StringType 并且 spark 不喜欢它。 我该怎么做才能让它发挥作用??

编辑:

我使用此代码创建了我的字典

Create multiple lists and store them into a dictionary Python

如您所见,创建密钥是在做

cc = str(col)
vv = "_" + str(value)
cv = cc + vv

dict_stable_feature[cv] = t

t 只是10 的二进制列表。

【问题讨论】:

  • 字典中每个键的值的数据类型是什么?它们是标量/单个值还是复杂类型,例如列表/对象?每个键应该是数据框中的新列还是行中的值?您能否与您的字典分享示例代码以及您预期的数据框表应该是什么样的?
  • 我发布了一些其他信息!

标签: string apache-spark dictionary pyspark apache-spark-sql


【解决方案1】:

让我们开始将您的 python 字典转换为正确位置的列表值列表(即用于初始化 spark 数据帧的预期数据结构之一)。

假设字典中的所有列表值的长度相同,您可以尝试以下操作。

column_names = []
dataset = None
for column_name in dict_stable_feature:
    column_names.append(column_name)
    column_values = dict_stable_feature[column_name]
    # initialize dataset ranges
    if dataset is None:
        dataset=[]
        
        for i in range(0,len(column_values)):
            dataset.append([column_values[i]])
    else:
        for ind,val in enumerate(column_values):
            dataset[ind].append(val)

my_df = sparkSession.createDataFrame(dataset,schema=column_names)

如果所有列表值的长度不同,那么您可以尝试以下方法:

max_list_length = max([len(dict_stable_feature[k]) for k in dict_stable_feature])
column_names = []
dataset = [[] for i in range(0,max_list_length)]
default_data_value = None # feel free to change
for column_name in dict_stable_feature:
    column_names.append(column_name)
    column_values = dict_stable_feature[column_name]


    for ind,val in enumerate(column_values):
        dataset[ind].append(val)

    # ensure all columns have the same amount of rows
    no_of_values = len(column_values)
    if  no_of_values < max_list_length:
        for i in range(no_of_values,max_list_length):
            dataset[i].append(default_data_value)

my_df = sparkSession.createDataFrame(dataset,schema=column_names)

让我知道这是否适合你。

【讨论】:

  • 第一个有效!正是我想要的!谢谢!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-01-12
  • 2016-05-09
相关资源
最近更新 更多