What is ExternalRDDScan in the DAG?

Viewed 588

What is the meaning of ExternalRDDScan in the DAG?

The whole internet doesn't have an explanation for it.

enter image description here

2 Answers

Based on the source, ExternalRDDScan is a representation of converting existing RDD of arbitrary objects to a dataset of InternalRows, i.e. creating a DataFrame. Let's verify that our understanding is correct:

scala> import spark.implicits._
import spark.implicits._

scala> val rdd = sc.parallelize(Array(1, 2, 3, 4, 5))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:26

scala> rdd.toDF().explain()
== Physical Plan ==
*(1) SerializeFromObject [input[0, int, false] AS value#2]
+- Scan ExternalRDDScan[obj#1]

ExternalRDD is a logical representation of DataFrame/Dataset (not in all cases though) in the query execution plan i.e. in DAG created by the spark.

ExternalRDD(s) are created

  1. when you create a DataFrame from RDD (viz. using createDataFrame(), toDF() )
  2. when you create a DataSet from RDD (viz. using createDataSet(), toDS() )

At runtime, when the ExternalRDD is to be loaded into the memory, a scan operation is done which is represented by ExternalRDDScan (internally the scan strategy is resolved to ExternalRDDScanExec). Look at the example below:

scala> val sampleRDD = sc.parallelize(Seq(1,2,3,4,5))
sampleRDD: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:24

scala> sampleRDD.toDF.queryExecution
res0: org.apache.spark.sql.execution.QueryExecution =
== Parsed Logical Plan ==
SerializeFromObject [input[0, int, false] AS value#2]
+- ExternalRDD [obj#1]

== Analyzed Logical Plan ==
value: int
SerializeFromObject [input[0, int, false] AS value#2]
+- ExternalRDD [obj#1]

== Optimized Logical Plan ==
SerializeFromObject [input[0, int, false] AS value#2]
+- ExternalRDD [obj#1]

== Physical Plan ==
*(1) SerializeFromObject [input[0, int, false] AS value#2]
+- Scan[obj#1]

You can see that in the query execution plan, the DataFrame object is represented by ExternalRDD and the physical plan contains a scan operation which is resolved to ExternalRDDScan (ExternalRDDScanExec) during its execution.

The same holds true for a spark Dataset as well.

scala> val sampleRDD = sc.parallelize(Seq(1,2,3,4,5))
sampleRDD: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:24

scala> sampleRDD.toDS.queryExecution.logical
res9: org.apache.spark.sql.catalyst.plans.logical.LogicalPlan =
SerializeFromObject [input[0, int, false] AS value#23]
+- ExternalRDD [obj#22]

scala> spark.createDataset(sampleRDD).queryExecution.logical
res18: org.apache.spark.sql.catalyst.plans.logical.LogicalPlan =
SerializeFromObject [input[0, int, false] AS value#39]
+- ExternalRDD [obj#38]

The above examples were run in spark version 2.4.2

Reference: https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-LogicalPlan-ExternalRDD.html

Related