在Databricks Notebook中按条件删SQL Server行出现全量删除问题
解决Spark连接SQL Server时按条件删除行失效的问题
你的两种方法为什么会全量删除?
方法一的问题
spark.read.jdbc(...).write.jdbc(..., mode='overwrite')这行代码会把读取到的表A数据全量覆盖写回表A,mode='overwrite'的逻辑是先删除原表再重建写入,这一步已经清空了原表内容,后续的DELETE操作根本无法生效。而且write.jdbc返回的是None,你调用connection.cursor()会直接抛出异常,实际导致全量删除的就是这个overwrite操作。- 你混淆了Spark的DataFrame读写逻辑和原生JDBC连接操作,Spark的JDBC读写API不是用来获取数据库连接对象的。
方法二的问题
spark.sql(delete_query)执行的是Spark SQL引擎的语句,不会直接发送到SQL Server执行,Spark无法通过该API直接执行外部数据库的DML语句。- 后续的
write.jdbc(..., mode='overwrite')同样会触发全量覆盖逻辑,导致所有数据被删除。
正确实现按条件删除的两种方式
方式一:使用原生JDBC连接执行DELETE语句
直接通过JDBC驱动创建数据库连接,执行DELETE语句,这是最直接高效的方式:
from java.sql import DriverManager def delete_old_records(jdbcUrl, connectionProperties): delete_query = """ DELETE FROM A WHERE EventId < 1000 """ # 创建原生JDBC连接 conn = DriverManager.getConnection( jdbcUrl, connectionProperties["user"], connectionProperties["password"] ) try: stmt = conn.createStatement() affected_rows = stmt.executeUpdate(delete_query) conn.commit() print(f"成功删除 {affected_rows} 条符合条件的记录") except Exception as e: conn.rollback() print(f"删除失败: {str(e)}") finally: stmt.close() conn.close() # 传入你的JDBC地址和连接属性调用函数 delete_old_records(jdbcUrl, connectionProperties)
注:Databricks环境通常已预装SQL Server JDBC驱动,若未预装需手动添加依赖。
方式二:使用Spark过滤后重写表(仅适合小表)
如果表数据量较小,可以先读取需要保留的数据,再覆盖写回原表,间接实现删除逻辑:
# 读取表中需要保留的记录(EventId >= 1000) df = spark.read.jdbc(url=jdbcUrl, table="A", properties=connectionProperties) filtered_df = df.filter(df.EventId >= 1000) # 覆盖写回原表,注意该操作会先删除原表再写入数据 filtered_df.write.jdbc(url=jdbcUrl, table="A", mode="overwrite", properties=connectionProperties)
⚠️ 注意:该方式不适合大表,全量读写性能差,且并发场景下易出现数据丢失,优先选择方式一。
内容的提问来源于stack exchange,提问作者LunaLoveDove
相关产品推荐
相关产品推荐

