使用Azure Databricks PySpark删除Azure SQL表行遇报错求助
解决Azure Databricks PySpark删除Azure SQL数据的报错问题
报错原因
你用spark.read执行DELETE语句时,Spark期望从JDBC源获取结果集,但SQL Server要求嵌套在子查询中的INSERT/UPDATE/DELETE/MERGE必须包含OUTPUT子句,否则无法返回结果集给Spark,从而抛出错误。
解决方案
方案1:需要返回被删除的数据
如果需要获取被删除的行数据,在DELETE语句中添加OUTPUT deleted.*来返回被删除的记录,这样Spark就能读取到结果集:
azuresqlOptions={ "driver":jdbcDriver, "url":jdbcUrl, "user":username, "port":jdbcPort, "password":password } # 添加OUTPUT子句返回被删除的数据 query = "(DELETE cone.address OUTPUT deleted.* WHERE Address_ID=756 ) ad1" df1 = spark.read.format("jdbc").option("header","true").options(**azuresqlOptions).option("dbtable",query).load() display(df1)
方案2:仅执行删除操作(无需返回数据)
如果不需要获取被删除的数据,不要用spark.read(它的核心用途是读取数据),直接执行DML语句即可。可以通过Java JDBC连接直接执行:
import java.sql.DriverManager # 加载JDBC驱动 Class.forName(jdbcDriver) # 建立连接 conn = DriverManager.getConnection(jdbcUrl, username, password) # 创建语句对象 stmt = conn.createStatement() # 执行DELETE语句 stmt.executeUpdate("DELETE cone.address WHERE Address_ID=756") # 关闭资源 stmt.close() conn.close()
或者使用Spark的jdbc写入方法结合空DataFrame触发执行(适合批量操作场景):
# 构造空DataFrame(结构不影响,仅用于触发JDBC操作) empty_df = spark.createDataFrame([], schema="Address_ID int") # 执行DELETE语句,用mode="ignore"避免写入空数据时报错 empty_df.write \ .format("jdbc") \ .options(**azuresqlOptions) \ .option("query", "DELETE cone.address WHERE Address_ID=756") \ .mode("ignore") \ .save()
关键说明
spark.read的核心功能是读取数据,用它执行DML语句时,必须保证语句能返回结果集(通过OUTPUT子句)。- 如果只是执行修改操作,优先选择直接执行DML语句的方式,避免
spark.read带来的结果集依赖问题。
内容的提问来源于stack exchange,提问作者Pradeep
相关产品推荐
相关产品推荐

