digraph G {
0 [labelType="html" label="<br><b>Project</b><br><br>"];
subgraph cluster1 {
isCluster="true";
label="WholeStageCodegen (1)";
2 [labelType="html" label="<br><b>SerializeFromObject</b><br><br>"];
3 [labelType="html" label="<br><b>MapElements</b><br><br>"];
4 [labelType="html" label="<br><b>DeserializeToObject</b><br><br>"];
5 [labelType="html" label="<br><b>Project</b><br><br>"];
6 [labelType="html" label="<br><b>Filter</b><br><br>"];
7 [labelType="html" label="<br><b>Scan ExistingRDD Delta Table State #0 - file:/home/jovyan/spark-warehouse/test.db/fhvhv_trip/_delta_log</b><br><br>"];
}
2->0;
3->2;
4->3;
5->4;
6->5;
7->6;
}
8
Project [path#1210, partitionValues#1211, size#1212L, modificationTime#1213L, dataChange#1214, from_json(StructField(numRecords,LongType,true), StructField(minValues,StructType(StructField(hvfhs_license_num,StringType,true),StructField(dispatching_base_num,StringType,true),StructField(originating_base_num,StringType,true),StructField(request_datetime,TimestampNTZType,true),StructField(on_scene_datetime,TimestampNTZType,true),StructField(pickup_datetime,TimestampNTZType,true),StructField(dropoff_datetime,TimestampNTZType,true),StructField(PULocationID,IntegerType,true),StructField(DOLocationID,IntegerType,true),StructField(trip_miles,DoubleType,true),StructField(trip_time,LongType,true),StructField(base_passenger_fare,DoubleType,true),StructField(tolls,DoubleType,true),StructField(bcf,DoubleType,true),StructField(sales_tax,DoubleType,true),StructField(congestion_surcharge,DoubleType,true),StructField(airport_fee,DoubleType,true),StructField(tips,DoubleType,true),StructField(driver_pay,DoubleType,true),StructField(shared_request_flag,StringType,true),StructField(shared_match_flag,StringType,true),StructField(access_a_ride_flag,StringType,true),StructField(wav_request_flag,StringType,true),StructField(wav_match_flag,StringType,true),StructField(cbd_congestion_fee,DoubleType,true)),true), StructField(maxValues,StructType(StructField(hvfhs_license_num,StringType,true),StructField(dispatching_base_num,StringType,true),StructField(originating_base_num,StringType,true),StructField(request_datetime,TimestampNTZType,true),StructField(on_scene_datetime,TimestampNTZType,true),StructField(pickup_datetime,TimestampNTZType,true),StructField(dropoff_datetime,TimestampNTZType,true),StructField(PULocationID,IntegerType,true),StructField(DOLocationID,IntegerType,true),StructField(trip_miles,DoubleType,true),StructField(trip_time,LongType,true),StructField(base_passenger_fare,DoubleType,true),StructField(tolls,DoubleType,true),StructField(bcf,DoubleType,true),StructField(sales_tax,DoubleType,true),StructField(congestion_surcharge,DoubleType,true),StructField(airport_fee,DoubleType,true),StructField(tips,DoubleType,true),StructField(driver_pay,DoubleType,true),StructField(shared_request_flag,StringType,true),StructField(shared_match_flag,StringType,true),StructField(access_a_ride_flag,StringType,true),StructField(wav_request_flag,StringType,true),StructField(wav_match_flag,StringType,true),StructField(cbd_congestion_fee,DoubleType,true)),true), StructField(nullCount,StructType(StructField(hvfhs_license_num,LongType,true),StructField(dispatching_base_num,LongType,true),StructField(originating_base_num,LongType,true),StructField(request_datetime,LongType,true),StructField(on_scene_datetime,LongType,true),StructField(pickup_datetime,LongType,true),StructField(dropoff_datetime,LongType,true),StructField(PULocationID,LongType,true),StructField(DOLocationID,LongType,true),StructField(trip_miles,LongType,true),StructField(trip_time,LongType,true),StructField(base_passenger_fare,LongType,true),StructField(tolls,LongType,true),StructField(bcf,LongType,true),StructField(sales_tax,LongType,true),StructField(congestion_surcharge,LongType,true),StructField(airport_fee,LongType,true),StructField(tips,LongType,true),StructField(driver_pay,LongType,true),StructField(shared_request_flag,LongType,true),StructField(shared_match_flag,LongType,true),StructField(access_a_ride_flag,LongType,true),StructField(wav_request_flag,LongType,true),StructField(wav_match_flag,LongType,true),StructField(cbd_congestion_fee,LongType,true)),true), stats#1215, Some(Etc/UTC)) AS stats#1233, tags#1216, deletionVector#1217, baseRowId#1218L, defaultRowCommitVersion#1219L, clusteringProvider#1220]
SerializeFromObject [staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).path, true, false, true) AS path#1210, externalmaptocatalyst(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -1), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -1), StringType, ObjectType(class java.lang.String)), true, false, true), lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -2), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -2), StringType, ObjectType(class java.lang.String)), true, false, true), knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).partitionValues) AS partitionValues#1211, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).size AS size#1212L, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).modificationTime AS modificationTime#1213L, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).dataChange AS dataChange#1214, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).stats, true, false, true) AS stats#1215, externalmaptocatalyst(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -3), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -3), StringType, ObjectType(class java.lang.String)), true, false, true), lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -4), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -4), StringType, ObjectType(class java.lang.String)), true, false, true), knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).tags) AS tags#1216, if (isnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector)) null else named_struct(storageType, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).storageType, true, false, true), pathOrInlineDv, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).pathOrInlineDv, true, false, true), offset, unwrapoption(IntegerType, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).offset), sizeInBytes, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).sizeInBytes, cardinality, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).cardinality, maxRowIndex, unwrapoption(LongType, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).maxRowIndex)) AS deletionVector#1217, unwrapoption(LongType, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).baseRowId) AS baseRowId#1218L, unwrapoption(LongType, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).defaultRowCommitVersion) AS defaultRowCommitVersion#1219L, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, unwrapoption(ObjectType(class java.lang.String), knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).clusteringProvider), true, false, true) AS clusteringProvider#1220]
MapElements org.apache.spark.sql.Dataset$$Lambda$4933/0x0000000841c25840@426e6b96, obj#1209: org.apache.spark.sql.delta.actions.AddFile
DeserializeToObject newInstance(class scala.Tuple1), obj#1208: scala.Tuple1
Project [add#1123]
Filter isnotnull(add#1123)
Scan ExistingRDD Delta Table State #0 - file:/home/jovyan/spark-warehouse/test.db/fhvhv_trip/_delta_log[txn#1122,add#1123,remove#1124,metaData#1125,protocol#1126,cdc#1127,checkpointMetadata#1128,sidecar#1129,domainMetadata#1130,commitInfo#1131]
WholeStageCodegen (1)
== Physical Plan ==
Project (7)
+- * SerializeFromObject (6)
+- * MapElements (5)
+- * DeserializeToObject (4)
+- * Project (3)
+- * Filter (2)
+- * Scan ExistingRDD Delta Table State #0 - file:/home/jovyan/spark-warehouse/test.db/fhvhv_trip/_delta_log (1)
(1) Scan ExistingRDD Delta Table State #0 - file:/home/jovyan/spark-warehouse/test.db/fhvhv_trip/_delta_log [codegen id : 1]
Output [10]: [txn#1122, add#1123, remove#1124, metaData#1125, protocol#1126, cdc#1127, checkpointMetadata#1128, sidecar#1129, domainMetadata#1130, commitInfo#1131]
Arguments: [txn#1122, add#1123, remove#1124, metaData#1125, protocol#1126, cdc#1127, checkpointMetadata#1128, sidecar#1129, domainMetadata#1130, commitInfo#1131], Delta Table State #0 - file:/home/jovyan/spark-warehouse/test.db/fhvhv_trip/_delta_log MapPartitionsRDD[38] at Spark Connect - session_id: "58d45445-3af3-4104-9739-56b4041ec0f9"
plan {
command {
sql_command {
sql: "select count(*) from test..., ExistingRDD, UnknownPartitioning(0)
(2) Filter [codegen id : 1]
Input [10]: [txn#1122, add#1123, remove#1124, metaData#1125, protocol#1126, cdc#1127, checkpointMetadata#1128, sidecar#1129, domainMetadata#1130, commitInfo#1131]
Condition : isnotnull(add#1123)
(3) Project [codegen id : 1]
Output [1]: [add#1123]
Input [10]: [txn#1122, add#1123, remove#1124, metaData#1125, protocol#1126, cdc#1127, checkpointMetadata#1128, sidecar#1129, domainMetadata#1130, commitInfo#1131]
(4) DeserializeToObject [codegen id : 1]
Input [1]: [add#1123]
Arguments: newInstance(class scala.Tuple1), obj#1208: scala.Tuple1
(5) MapElements [codegen id : 1]
Input [1]: [obj#1208]
Arguments: org.apache.spark.sql.Dataset$$Lambda$4933/0x0000000841c25840@426e6b96, obj#1209: org.apache.spark.sql.delta.actions.AddFile
(6) SerializeFromObject [codegen id : 1]
Input [1]: [obj#1209]
Arguments: [staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).path, true, false, true) AS path#1210, externalmaptocatalyst(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -1), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -1), StringType, ObjectType(class java.lang.String)), true, false, true), lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -2), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -2), StringType, ObjectType(class java.lang.String)), true, false, true), knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).partitionValues) AS partitionValues#1211, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).size AS size#1212L, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).modificationTime AS modificationTime#1213L, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).dataChange AS dataChange#1214, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).stats, true, false, true) AS stats#1215, externalmaptocatalyst(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -3), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_key, ObjectType(class java.lang.Object), true, -3), StringType, ObjectType(class java.lang.String)), true, false, true), lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -4), staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, validateexternaltype(lambdavariable(ExternalMapToCatalyst_value, ObjectType(class java.lang.Object), true, -4), StringType, ObjectType(class java.lang.String)), true, false, true), knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).tags) AS tags#1216, if (isnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector)) null else named_struct(storageType, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).storageType, true, false, true), pathOrInlineDv, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).pathOrInlineDv, true, false, true), offset, unwrapoption(IntegerType, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).offset), sizeInBytes, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).sizeInBytes, cardinality, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).cardinality, maxRowIndex, unwrapoption(LongType, knownnotnull(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).deletionVector).maxRowIndex)) AS deletionVector#1217, unwrapoption(LongType, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).baseRowId) AS baseRowId#1218L, unwrapoption(LongType, knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).defaultRowCommitVersion) AS defaultRowCommitVersion#1219L, staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, unwrapoption(ObjectType(class java.lang.String), knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.actions.AddFile, true])).clusteringProvider), true, false, true) AS clusteringProvider#1220]
(7) Project
Output [11]: [path#1210, partitionValues#1211, size#1212L, modificationTime#1213L, dataChange#1214, from_json(StructField(numRecords,LongType,true), StructField(minValues,StructType(StructField(hvfhs_license_num,StringType,true),StructField(dispatching_base_num,StringType,true),StructField(originating_base_num,StringType,true),StructField(request_datetime,TimestampNTZType,true),StructField(on_scene_datetime,TimestampNTZType,true),StructField(pickup_datetime,TimestampNTZType,true),StructField(dropoff_datetime,TimestampNTZType,true),StructField(PULocationID,IntegerType,true),StructField(DOLocationID,IntegerType,true),StructField(trip_miles,DoubleType,true),StructField(trip_time,LongType,true),StructField(base_passenger_fare,DoubleType,true),StructField(tolls,DoubleType,true),StructField(bcf,DoubleType,true),StructField(sales_tax,DoubleType,true),StructField(congestion_surcharge,DoubleType,true),StructField(airport_fee,DoubleType,true),StructField(tips,DoubleType,true),StructField(driver_pay,DoubleType,true),StructField(shared_request_flag,StringType,true),StructField(shared_match_flag,StringType,true),StructField(access_a_ride_flag,StringType,true),StructField(wav_request_flag,StringType,true),StructField(wav_match_flag,StringType,true),StructField(cbd_congestion_fee,DoubleType,true)),true), StructField(maxValues,StructType(StructField(hvfhs_license_num,StringType,true),StructField(dispatching_base_num,StringType,true),StructField(originating_base_num,StringType,true),StructField(request_datetime,TimestampNTZType,true),StructField(on_scene_datetime,TimestampNTZType,true),StructField(pickup_datetime,TimestampNTZType,true),StructField(dropoff_datetime,TimestampNTZType,true),StructField(PULocationID,IntegerType,true),StructField(DOLocationID,IntegerType,true),StructField(trip_miles,DoubleType,true),StructField(trip_time,LongType,true),StructField(base_passenger_fare,DoubleType,true),StructField(tolls,DoubleType,true),StructField(bcf,DoubleType,true),StructField(sales_tax,DoubleType,true),StructField(congestion_surcharge,DoubleType,true),StructField(airport_fee,DoubleType,true),StructField(tips,DoubleType,true),StructField(driver_pay,DoubleType,true),StructField(shared_request_flag,StringType,true),StructField(shared_match_flag,StringType,true),StructField(access_a_ride_flag,StringType,true),StructField(wav_request_flag,StringType,true),StructField(wav_match_flag,StringType,true),StructField(cbd_congestion_fee,DoubleType,true)),true), StructField(nullCount,StructType(StructField(hvfhs_license_num,LongType,true),StructField(dispatching_base_num,LongType,true),StructField(originating_base_num,LongType,true),StructField(request_datetime,LongType,true),StructField(on_scene_datetime,LongType,true),StructField(pickup_datetime,LongType,true),StructField(dropoff_datetime,LongType,true),StructField(PULocationID,LongType,true),StructField(DOLocationID,LongType,true),StructField(trip_miles,LongType,true),StructField(trip_time,LongType,true),StructField(base_passenger_fare,LongType,true),StructField(tolls,LongType,true),StructField(bcf,LongType,true),StructField(sales_tax,LongType,true),StructField(congestion_surcharge,LongType,true),StructField(airport_fee,LongType,true),StructField(tips,LongType,true),StructField(driver_pay,LongType,true),StructField(shared_request_flag,LongType,true),StructField(shared_match_flag,LongType,true),StructField(access_a_ride_flag,LongType,true),StructField(wav_request_flag,LongType,true),StructField(wav_match_flag,LongType,true),StructField(cbd_congestion_fee,LongType,true)),true), stats#1215, Some(Etc/UTC)) AS stats#1233, tags#1216, deletionVector#1217, baseRowId#1218L, defaultRowCommitVersion#1219L, clusteringProvider#1220]
Input [11]: [path#1210, partitionValues#1211, size#1212L, modificationTime#1213L, dataChange#1214, stats#1215, tags#1216, deletionVector#1217, baseRowId#1218L, defaultRowCommitVersion#1219L, clusteringProvider#1220]