Spark中Delete结合Exists报错:不允许非确定性表达式
替代方案:删除TABLE1中与TABLE2主键匹配的记录
针对Spark 3.2.1 Databricks环境中EXISTS在DELETE语句里触发非确定性表达式报错的问题,以下是几种可行的替代实现方式:
方案1:使用IN子句替换EXISTS
直接通过IN子句匹配TABLE2的主键组合,语法简洁且能避开EXISTS的判定问题:
val op = spark.sql(s""" DELETE FROM TABLE1 WHERE (DAY, DATA_STREAM) IN ( SELECT DAY, DATA_STREAM FROM TABLE2 ) """)
方案2:使用MERGE INTO语句
利用MERGE INTO的匹配删除逻辑,这是Spark处理这类数据同步场景的标准方式之一:
val op = spark.sql(s""" MERGE INTO TABLE1 AS t USING TABLE2 AS s ON t.DAY = s.DAY AND t.DATA_STREAM = s.DATA_STREAM WHEN MATCHED THEN DELETE """)
方案3:先筛选待删除记录再删除
如果TABLE2数据量较小,可以先提取主键集合,再作为过滤条件删除:
// 先获取TABLE2的主键数据 val deleteKeys = spark.sql("SELECT DAY, DATA_STREAM FROM TABLE2").collect() // 构造删除条件 val deleteCondition = deleteKeys.map { row => s"(DAY = '${row.getAs[String]("DAY")}' AND DATA_STREAM = '${row.getAs[String]("DATA_STREAM")}')" }.mkString(" OR ") val op = spark.sql(s""" DELETE FROM TABLE1 WHERE $deleteCondition """)
注意:此方案仅适用于TABLE2数据量不大的场景,避免collect()导致Driver内存溢出。
内容的提问来源于stack exchange,提问作者ultraInstinct
相关产品推荐
相关产品推荐

