从Python字典创建PySpark DataFrame报错无法推断str类型schema如何解决
问题根源
你遇到的报错和Python原生str与Spark StringType的适配无关,本质是你传入spark.createDataFrame的字典结构不符合方法预期:
Spark的createDataFrame默认要求传入的可迭代对象是行维度的集合,但你当前的字典dict_stable_feature是{列名: 该列所有值组成的列表}的列维度结构,Spark解析时误将字符串型的字典键当成了待解析的数据行,才触发了类型推断失败的报错。
解决方法
方案1:借助pandas转换(最简便)
如果你的环境允许使用pandas,直接转成pandas DataFrame再转Spark即可:
import pandas as pd df_stable = spark.createDataFrame(pd.DataFrame(dict_stable_feature)) df_stable.show()
方案2:纯PySpark原生实现(无第三方依赖)
手动将列维度的字典转成行维度的集合,再传入创建DataFrame即可:
from pyspark.sql.types import StructType, StructField, IntegerType # 1. 自定义Schema,明确列名和对应数据类型(你的值都是0/1,用IntegerType即可) schema = StructType([ StructField(col_name, IntegerType(), nullable=False) for col_name in dict_stable_feature.keys() ]) # 2. 将列存储结构转成行存储:将各列的对应位置值打包为一行 row_data = zip(*dict_stable_feature.values()) # 3. 传入行数据和Schema创建DataFrame df_stable = spark.createDataFrame(row_data, schema=schema) df_stable.show()
方案3:基于RDD转换
也可以先将行数据转成RDD再转DataFrame:
row_rdd = spark.sparkContext.parallelize(zip(*dict_stable_feature.values())) df_stable = row_rdd.toDF(schema=list(dict_stable_feature.keys())) df_stable.show()
内容的提问来源于stack exchange,提问作者ianux22
相关产品推荐
相关产品推荐

