PySpark Pandas转Spark DF后执行DELETE语句遇非确定性表达式错误求助
问题原因
当你用pyspark.pandas.DataFrame.to_spark()转换得到Spark DataFrame时,pyspark.pandas会自动添加一个名为__natural_order__的列,这个列是通过monotonically_increasing_id()生成的。monotonically_increasing_id()属于非确定性表达式——每次执行生成的ID值不固定,而Spark的DELETE操作不允许在这类DML语句中引用包含非确定性表达式的数据源,因此触发了报错。
从报错的执行计划里也能看到这个线索:
+- Project [__index_level_0__#0L, Name#1, Age#2, monotonically_increasing_id() AS __natural_order__#6L]
解决方法
方法1:转换后移除自动添加的非确定性列
在注册临时表前,手动删除__natural_order__和自动生成的索引列__index_level_0__:
from pyspark import pandas as ps from spark_sdk.sql import SparkSql sql = SparkSql() data = [['tom', '10'], ['nick', '15'], ['juli', '14']] # 创建pyspark-pandas DataFrame df = ps.DataFrame(data, columns=['Name', 'Age']) # 转换后移除自动添加的列 pandas_df = df.to_spark().drop("__natural_order__", "__index_level_0__") s_df = pandas_df # 注册DF为临时表 s_df.registerTempTable("temp_table") temp = sql.execute_query(query="select * from temp_table") print(temp.show()) # 创建表x create_query = f"create table x(Name String, Age String)" sql.execute_query(query=create_query) # 执行DELETE语句 delete_query = f"delete from x as t1 WHERE t1.name in (SELECT distinct Name from temp_table)" sql.execute_query(query=delete_query)
方法2:创建pyspark.pandas DF时禁用自然排序列
通过设置ps.set_option("compute.default_index_type", "distributed"),让pyspark.pandas不生成依赖monotonically_increasing_id()的自然排序列:
from pyspark import pandas as ps from spark_sdk.sql import SparkSql sql = SparkSql() # 禁用默认的自然排序列生成 ps.set_option("compute.default_index_type", "distributed") data = [['tom', '10'], ['nick', '15'], ['juli', '14']] # 创建pyspark-pandas DataFrame df = ps.DataFrame(data, columns=['Name', 'Age']) pandas_df = df.to_spark() s_df = pandas_df # 注册DF为临时表 s_df.registerTempTable("temp_table") temp = sql.execute_query(query="select * from temp_table") print(temp.show()) # 创建表x create_query = f"create table x(Name String, Age String)" sql.execute_query(query=create_query) # 执行DELETE语句 delete_query = f"delete from x as t1 WHERE t1.name in (SELECT distinct Name from temp_table)" sql.execute_query(query=delete_query)
方法3:改用JOIN方式编写DELETE语句
另一种思路是用JOIN替代IN子查询,这种写法也能避开非确定性表达式的问题:
# 替换原来的DELETE语句 delete_query = f""" DELETE FROM x as t1 USING temp_table as t2 WHERE t1.name = t2.Name """ sql.execute_query(query=delete_query)
内容的提问来源于stack exchange,提问作者Puneet Jain
相关产品推荐
相关产品推荐

