PySpark执行spark.sql关联DataFrame报表不存在报错如何解决
报错触发原因
spark.sql() 方法仅能识别已注册到Spark SQL元数据的表或临时视图,代码中定义的df、df1本质是Python运行环境中的DataFrame实例,未完成临时视图注册流程,SQL引擎解析语句时无法找到对应数据源,因此抛出Table or view not found异常。
可行解决方案
方案1:保留SQL写法,提前注册临时视图
执行SQL前调用createOrReplaceTempView方法将两个DataFrame注册为临时视图,视图名与SQL语句中引用的表名保持一致即可:
df = spark.createDataFrame(r, schema=column) # 注册临时视图 df.createOrReplaceTempView("df") df1.createOrReplaceTempView("df1") # 执行内连接查询 df_final = spark.sql(""" SELECT * FROM df INNER JOIN df1 ON df.a = df1.b """)
临时视图生命周期与当前SparkSession绑定,SparkSession停止后视图会自动清除,不会持久化到存储。
方案2:使用DataFrame原生API实现内连接
无需注册临时视图,直接调用DataFrame的join方法完成连接,更适配PySpark的Python开发流程:
df = spark.createDataFrame(r, schema=column) # 传入三个参数:待连接的DataFrame、连接条件、连接类型 df_final = df.join(df1, on=df.a == df1.b, how="inner")
注意:如果两个DataFrame存在同名字段,select *的写法会返回重复列,后续调用列时可能触发列歧义报错,可根据业务需求手动指定需要返回的字段列表。
内容的提问来源于stack exchange,提问作者roger montez
相关产品推荐
相关产品推荐

