常规
好问题。
通过测试场景和推断找不到合适的文档来推断答案。第二次尝试,由于网络上的各种语句无法备份。
我认为这个问题与 AQE Spark 3.x 方面无关,但它是
关于说,作为 Spark 应用程序阶段 N 的一部分的数据帧具有
通过了从静态源获取数据的阶段,即
受到应用多个谓词的过滤。
那么中心点是谓词的方式是否重要
Spark (Catalyst) 是否对谓词重新排序以最小化
要完成的工作?
- 这里的前提是,首先过滤出最大数量的数据比评估过滤非常多的谓词更有意义
很少出来。
- 这是一个众所周知的 RDBMS 点,它指的是 sargable 谓词(受定义随着时间的推移而演变)。
- 很多讨论都集中在索引上,Spark、Hive 没有这个,但 DF 是柱状的。
第 1 点
你可以试试%sql
EXPLAIN EXTENDED select k, sum(v) from values (1, 2), (1, 3) t(k, v) group by k;
从这里你可以看到如果有重新安排会发生什么
谓词,但我在非 AQE 的物理计划中没有看到这些方面
Databricks 上的模式。参考
https://docs.databricks.com/sql/language-manual/sql-ref-syntax-qry-explain.html.
Catalyst 可以重新安排我在这里和那里读到的过滤。到什么
程度,是大量的研究;我无法确认这一点。
也是一个有趣的阅读:
https://www.waitingforcode.com/apache-spark-sql/catalyst-optimizer-in-spark-sql/read
第 2 点
我用相同的方式运行了以下可悲的人为示例
功能查询,但谓词颠倒,使用具有
高基数并测试了实际上不存在的值
然后在调用时比较 UDF 中使用的累加器的计数。
场景 1
import org.apache.spark.sql.functions._
def randomInt1to1000000000 = scala.util.Random.nextInt(1000000000)+1
def randomInt1to10 = scala.util.Random.nextInt(10)+1
def randomInt1to1000000 = scala.util.Random.nextInt(1000000)+1
val df = sc.parallelize(Seq.fill(1000000){(randomInt1to1000000,randomInt1to1000000000,randomInt1to10)}).toDF("nuid","hc", "lc").withColumn("text", lpad($"nuid", 3, "0")).withColumn("literal",lit(1))
val accumulator = sc.longAccumulator("udf_call_count")
spark.udf.register("myUdf", (x: String) => {accumulator.add(1)
x.length}
)
accumulator.reset()
df.where("myUdf(text) = 3 and hc = -4").select(max($"text")).show(false)
println(s"Number of UDF calls ${accumulator.value}")
返回:
+---------+
|max(text)|
+---------+
|null |
+---------+
Number of UDF calls 1000000
场景 2
import org.apache.spark.sql.functions._
def randomInt1to1000000000 = scala.util.Random.nextInt(1000000000)+1
def randomInt1to10 = scala.util.Random.nextInt(10)+1
def randomInt1to1000000 = scala.util.Random.nextInt(1000000)+1
val dfA = sc.parallelize(Seq.fill(1000000){(randomInt1to1000000,randomInt1to1000000000,randomInt1to10)}).toDF("nuid","hc", "lc").withColumn("text", lpad($"nuid", 3, "0")).withColumn("literal",lit(1))
val accumulator = sc.longAccumulator("udf_call_count")
spark.udf.register("myUdf", (x: String) => {accumulator.add(1)
x.length}
)
accumulator.reset()
dfA.where("hc = -4 and myUdf(text) = 3").select(max($"text")).show(false)
println(s"Number of UDF calls ${accumulator.value}")
返回:
+---------+
|max(text)|
+---------+
|null |
+---------+
Number of UDF calls 0
我的结论是:
第 3 点
我从手册中注意到
https://docs.databricks.com/spark/latest/spark-sql/udf-scala.html
以下:
求值顺序和空值检查
Spark SQL(包括 SQL 以及 DataFrame 和 Dataset API)不
保证子表达式的求值顺序。特别是,
不一定要评估运算符或函数的输入
从左到右或以任何其他固定顺序。例如,逻辑与
和 OR 表达式没有从左到右的“短路”
语义。
因此,依赖副作用或顺序是危险的
布尔表达式的评估,以及 WHERE 和 HAVING 的顺序
子句,因为这样的表达式和子句可以在
查询优化和规划。具体来说,如果 UDF 依赖于
SQL 中用于空检查的短路语义,没有
保证在调用 UDF 之前会进行空值检查。为了
例如,
spark.udf.register("strlen", (s: String) => s.length)
spark.sql("select s from test1 where s is not null and strlen(s) > 1") // no guarantee
此 WHERE 子句不保证调用 strlen UDF
过滤掉空值后。
要执行正确的 null 检查,我们建议您执行以下任一操作
以下:
使 UDF 本身具有 null 感知能力并在 UDF 内部进行 null 检查
本身使用 IF 或 CASE WHEN 表达式进行空检查并调用
条件分支中的 UDF。
spark.udf.register("strlen_nullsafe", (s: String) => if (s != null) s.length else -1)
spark.sql("select s from test1 where s is not null and strlen_nullsafe(s) > 1") // ok
spark.sql("select s from test1 where if(s is not null, strlen(s), null) > 1") // ok
有点矛盾。