这是我正在运行的示例代码.
使用mod列作为分区创建测试镶木地板数据集.
scala> val test = spark.range(0 , 100000000).withColumn("mod", $"id".mod(40))
test: org.apache.spark.sql.DataFrame = [id: bigint, mod: bigint]
scala> test.write.partitionBy("mod").mode("overwrite").parquet("test_pushdown_filter")
Run Code Online (Sandbox Code Playgroud)
之后,我将这些数据作为数据框架读取并在分区列上应用过滤器即mod.
scala> val df = spark.read.parquet("test_pushdown_filter").filter("mod = 5")
df: org.apache.spark.sql.Dataset[org.apache.spark.sql.Row] = [id: bigint, mod: int]
scala> df.queryExecution.executedPlan
res1: org.apache.spark.sql.execution.SparkPlan =
*FileScan parquet [id#16L,mod#17] Batched: true, Format: Parquet, Location: InMemoryFileIndex[file:/C:/Users/kprajapa/WorkSpace/places/test_pushdown_filter], PartitionCount: 1, PartitionFilters: [
isnotnull(mod#17), (mod#17 = 5)], PushedFilters: [], ReadSchema: struct<id:bigint>
Run Code Online (Sandbox Code Playgroud)
你可以在执行计划中看到它只读取1个分区.
但是,如果您将相同的过滤器应用于数据集.它读取所有分区,然后应用过滤器.
scala> case class Test(id: Long, mod: Long)
defined class Test
scala> val ds = spark.read.parquet("test_pushdown_filter").as[Test].filter(_.mod==5)
ds: …Run Code Online (Sandbox Code Playgroud) apache-spark apache-spark-sql apache-spark-dataset catalyst-optimizer
我正在阅读High Performance Spark,作者做出以下声明:
虽然 Catalyst 优化器非常强大,但它目前遇到挑战的情况之一是非常大的查询计划。这些查询计划往往是迭代算法的结果,例如图算法或机器学习算法。一个简单的解决方法是将数据转换为 RDD,并在每次迭代结束时转换回 DataFrame/Dataset,如例 3-58 所示。
示例 3-58 被标记为“通过 RDD 进行往返以削减查询计划”,并复制如下:
val rdd = df.rdd
rdd.cache()
sqlCtx.createDataFrame(rdd. df.schema)
Run Code Online (Sandbox Code Playgroud)
有谁知道需要此解决方法的根本原因是什么?
作为参考,已针对此问题提交了错误报告,可通过以下链接获取: https ://issues.apache.org/jira/browse/SPARK-13346
似乎没有解决办法,但维护者已经解决了这个问题,并且似乎不认为他们需要解决它。