Spark SQL全局临时视图转DataFrame选元素报错求助
问题分析与修复
错误原因
- RDD数据格式错误:
parallelize(self.__dict__)会把字典的键(如'model_id'、'conn_string')作为RDD的单个元素,而非字典键值对组成的行数据。Spark尝试用单个字符串匹配StructType定义的多字段结构,因此触发类型不匹配错误。 - 类型不匹配:初始化
Model时传入的model_id是字符串"1",但schema定义的是IntegerType,后续读取会引发类型转换问题。
修复后的代码
from pyspark.sql import SparkSession from pyspark.sql.types import * class Model: def __init__(self, model_id : int = None, model_name : str = None, conn_string : str = None) : self.model_id = model_id self.model_name = model_name self.conn_string = conn_string def save_to_temp(self): spark = SparkSession.builder.getOrCreate() schema: StructType = StructType([ StructField('model_id', IntegerType(), nullable=True), StructField('model_name', StringType(), nullable=True), StructField('conn_string', StringType(), nullable=True) ]) # 将对象的键值对包装为单元素列表,确保每行是完整的字典数据 rdd = spark.sparkContext.parallelize([self.__dict__]) df = spark.createDataFrame(rdd, schema) df.createOrReplaceGlobalTempView("temp_tbl") def read_from_temp(self): spark = SparkSession.builder.getOrCreate() df = spark.sql("select * from global_temp.temp_tbl") self.model_id = df.collect()[0][0] df.show() dataCollect = df.collect() print(dataCollect) if __name__ == "__main__": # 将model_id转为整数,匹配schema定义的类型 model_id = 1 model_name = "model name 1" conn_string = "conn_string 1" model1 = Model(model_id, model_name, conn_string) model1.save_to_temp() model2 = Model() print("Initialized model2") model2.read_from_temp()
关键修复点
- 修正RDD构造逻辑:用
[self.__dict__]将字典包装为列表,让parallelize生成的RDD每个元素是完整的一行数据(字典),与StructType的字段定义匹配。 - 统一数据类型:把
model_id从字符串"1"改为整数1,和schema中的IntegerType保持一致,避免类型转换异常。
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

