digraph G {
0 [id="node0" labelType="html" label="<br><b>AdaptiveSparkPlan</b><br><br>" tooltip="AdaptiveSparkPlan isFinalPlan=true"];
subgraph cluster1 {
isCluster="true";
id="cluster1";
label="WholeStageCodegen (2)\n \nduration: 24 ms";
tooltip="WholeStageCodegen (2)";
2 [id="node2" labelType="html" label="<b>HashAggregate</b><br><br>time in aggregation build: 24 ms<br>number of output rows: 1" tooltip="HashAggregate(keys=[], functions=[count(1)])"];
}
3 [id="node3" labelType="html" label="<b>Exchange</b><br><br>shuffle records written: 200<br>remote merged reqs duration: 0 ms<br>remote merged blocks fetched: 0<br>records read: 200<br>local bytes read: 14.1 KiB<br>merged fetch fallback count: 0<br>local blocks read: 200<br>remote merged chunks fetched: 0<br>remote blocks read: 0<br>data size total (min, med, max (stageId: taskId))<br>3.1 KiB (16.0 B, 16.0 B, 16.0 B (stage 102.0: task 11065))<br>local merged bytes read: 0.0 B<br>local merged chunks fetched: 0<br>shuffle write time total (min, med, max (stageId: taskId))<br>32 ms (0 ms, 0 ms, 0 ms (stage 102.0: task 11164))<br>remote merged bytes read: 0.0 B<br>local merged blocks fetched: 0<br>corrupt merged block chunks: 0<br>fetch wait time: 0 ms<br>remote bytes read: 0.0 B<br>number of partitions: 1<br>remote reqs duration: 0 ms<br>remote bytes read to disk: 0.0 B<br>shuffle bytes written total (min, med, max (stageId: taskId))<br>14.1 KiB (72.0 B, 72.0 B, 72.0 B (stage 102.0: task 11065))" tooltip="Exchange SinglePartition, ENSURE_REQUIREMENTS, [plan_id=2150]"];
subgraph cluster4 {
isCluster="true";
id="cluster4";
label="WholeStageCodegen (1)\n \nduration: total (min, med, max (stageId: taskId))\n3 ms (0 ms, 0 ms, 3 ms (stage 102.0: task 11164))";
tooltip="WholeStageCodegen (1)";
5 [id="node5" labelType="html" label="<b>HashAggregate</b><br><br>time in aggregation build total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 102.0: task 11065))<br>number of output rows: 200" tooltip="HashAggregate(keys=[], functions=[partial_count(1)])"];
}
6 [id="node6" labelType="html" label="<b>InMemoryTableScan</b><br><br>number of output rows: 0" tooltip="InMemoryTableScan"];
7 [id="node7" labelType="html" label="<br><b>AdaptiveSparkPlan</b><br><br>" tooltip="AdaptiveSparkPlan isFinalPlan=true"];
subgraph cluster8 {
isCluster="true";
id="cluster8";
label="WholeStageCodegen (2)\n \nduration: 0 ms";
tooltip="WholeStageCodegen (2)";
9 [id="node9" labelType="html" label="<br><b>SerializeFromObject</b><br><br>" tooltip="SerializeFromObject [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#1699, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#1700L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#1701, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#1702L]"];
}
10 [id="node10" labelType="html" label="<br><b>MapGroups</b><br><br>" tooltip="MapGroups org.apache.spark.sql.internal.UDFAdaptors$$$Lambda/0x000000080209d5c0@563ebb8d, invoke(value#1693.toString()), newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), [value#1693], [path#1683, length#1684L, isDir#1685, modificationTime#1686L], obj#1698: org.apache.spark.sql.delta.SerializableFileStatus"];
subgraph cluster11 {
isCluster="true";
id="cluster11";
label="WholeStageCodegen (1)\n \nduration: 0 ms";
tooltip="WholeStageCodegen (1)";
12 [id="node12" labelType="html" label="<b>Sort</b><br><br>sort time: 0 ms<br>peak memory: 0.0 B<br>spill size: 0.0 B" tooltip="Sort [value#1693 ASC NULLS FIRST], false, 0"];
}
13 [id="node13" labelType="html" label="<b>Exchange</b><br><br>shuffle records written: 0<br>remote merged reqs duration: 0 ms<br>remote merged blocks fetched: 0<br>records read: 0<br>local bytes read: 0.0 B<br>merged fetch fallback count: 0<br>local blocks read: 0<br>remote merged chunks fetched: 0<br>remote blocks read: 0<br>data size: 0.0 B<br>local merged bytes read: 0.0 B<br>local merged chunks fetched: 0<br>shuffle write time: 0 ms<br>remote merged bytes read: 0.0 B<br>local merged blocks fetched: 0<br>corrupt merged block chunks: 0<br>fetch wait time: 0 ms<br>remote bytes read: 0.0 B<br>number of partitions: 0<br>remote reqs duration: 0 ms<br>remote bytes read to disk: 0.0 B<br>shuffle bytes written: 0.0 B" tooltip="Exchange hashpartitioning(value#1693, 200), ENSURE_REQUIREMENTS, [plan_id=2018]"];
14 [id="node14" labelType="html" label="<br><b>AppendColumnsWithObject</b><br><br>" tooltip="AppendColumnsWithObject org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x000000080209c000@3c79482b, [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#1683, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#1684L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#1685, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#1686L], [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#1693]"];
15 [id="node15" labelType="html" label="<br><b>MapElements</b><br><br>" tooltip="MapElements org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802093950@5e5aa004, obj#1682: org.apache.spark.sql.delta.SerializableFileStatus"];
16 [id="node16" labelType="html" label="<b>Scan</b><br><br>number of output rows: 0" tooltip="Scan[obj#1667]"];
2->0;
3->2;
5->3;
6->5;
7->6;
9->7;
10->9;
12->10;
13->12;
14->13;
15->14;
16->15;
}
== Physical Plan ==
AdaptiveSparkPlan (25)
+- == Final Plan ==
ResultQueryStage (21), Statistics(sizeInBytes=8.0 EiB)
+- * HashAggregate (20)
+- ShuffleQueryStage (19), Statistics(sizeInBytes=3.1 KiB, rowCount=200)
+- Exchange (18)
+- * HashAggregate (17)
+- TableCacheQueryStage (16), Statistics(sizeInBytes=0.0 B, rowCount=0)
+- InMemoryTableScan (1)
+- InMemoryRelation (2)
+- AdaptiveSparkPlan (15)
+- == Final Plan ==
ResultQueryStage (11)
+- * SerializeFromObject (10)
+- MapGroups (9)
+- * Sort (8)
+- ShuffleQueryStage (7), Statistics(sizeInBytes=0.0 B, rowCount=0)
+- Exchange (6)
+- AppendColumnsWithObject (5)
+- MapElements (4)
+- Scan (3)
+- == Initial Plan ==
SerializeFromObject (14)
+- MapGroups (13)
+- Sort (12)
+- Exchange (6)
+- AppendColumnsWithObject (5)
+- MapElements (4)
+- Scan (3)
+- == Initial Plan ==
HashAggregate (24)
+- Exchange (23)
+- HashAggregate (22)
+- InMemoryTableScan (1)
+- InMemoryRelation (2)
+- AdaptiveSparkPlan (15)
+- == Final Plan ==
ResultQueryStage (11)
+- * SerializeFromObject (10)
+- MapGroups (9)
+- * Sort (8)
+- ShuffleQueryStage (7), Statistics(sizeInBytes=0.0 B, rowCount=0)
+- Exchange (6)
+- AppendColumnsWithObject (5)
+- MapElements (4)
+- Scan (3)
+- == Initial Plan ==
SerializeFromObject (14)
+- MapGroups (13)
+- Sort (12)
+- Exchange (6)
+- AppendColumnsWithObject (5)
+- MapElements (4)
+- Scan (3)
(1) InMemoryTableScan
Output: []
(2) InMemoryRelation
Arguments: [path#1699, length#1700L, isDir#1701, modificationTime#1702L], StorageLevel(disk, memory, deserialized, 1 replicas)
(3) Scan
Output [1]: [obj#1667]
Arguments: obj#1667: org.apache.spark.sql.delta.SerializableFileStatus, MapPartitionsRDD[189] at $anonfun$recordDeltaOperationInternal$1 at DatabricksLogging.scala:128
(4) MapElements
Input [1]: [obj#1667]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802093950@5e5aa004, obj#1682: org.apache.spark.sql.delta.SerializableFileStatus
(5) AppendColumnsWithObject
Input [1]: [obj#1682]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x000000080209c000@3c79482b, [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#1683, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#1684L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#1685, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#1686L], [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#1693]
(6) Exchange
Input [5]: [path#1683, length#1684L, isDir#1685, modificationTime#1686L, value#1693]
Arguments: hashpartitioning(value#1693, 200), ENSURE_REQUIREMENTS, [plan_id=2018]
(7) ShuffleQueryStage
Output [5]: [path#1683, length#1684L, isDir#1685, modificationTime#1686L, value#1693]
Arguments: 0
(8) Sort [codegen id : 1]
Input [5]: [path#1683, length#1684L, isDir#1685, modificationTime#1686L, value#1693]
Arguments: [value#1693 ASC NULLS FIRST], false, 0
(9) MapGroups
Input [5]: [path#1683, length#1684L, isDir#1685, modificationTime#1686L, value#1693]
Arguments: org.apache.spark.sql.internal.UDFAdaptors$$$Lambda/0x000000080209d5c0@563ebb8d, invoke(value#1693.toString()), newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), [value#1693], [path#1683, length#1684L, isDir#1685, modificationTime#1686L], obj#1698: org.apache.spark.sql.delta.SerializableFileStatus
(10) SerializeFromObject [codegen id : 2]
Input [1]: [obj#1698]
Arguments: [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#1699, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#1700L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#1701, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#1702L]
(11) ResultQueryStage
Output [4]: [path#1699, length#1700L, isDir#1701, modificationTime#1702L]
Arguments: 1
(12) Sort
Input [5]: [path#1683, length#1684L, isDir#1685, modificationTime#1686L, value#1693]
Arguments: [value#1693 ASC NULLS FIRST], false, 0
(13) MapGroups
Input [5]: [path#1683, length#1684L, isDir#1685, modificationTime#1686L, value#1693]
Arguments: org.apache.spark.sql.internal.UDFAdaptors$$$Lambda/0x000000080209d5c0@563ebb8d, invoke(value#1693.toString()), newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), [value#1693], [path#1683, length#1684L, isDir#1685, modificationTime#1686L], obj#1698: org.apache.spark.sql.delta.SerializableFileStatus
(14) SerializeFromObject
Input [1]: [obj#1698]
Arguments: [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#1699, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#1700L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#1701, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#1702L]
(15) AdaptiveSparkPlan
Output [4]: [path#1699, length#1700L, isDir#1701, modificationTime#1702L]
Arguments: isFinalPlan=true
(16) TableCacheQueryStage
Output: []
Arguments: 0
(17) HashAggregate [codegen id : 1]
Input: []
Keys: []
Functions [1]: [partial_count(1)]
Aggregate Attributes [1]: [count#1907L]
Results [1]: [count#1908L]
(18) Exchange
Input [1]: [count#1908L]
Arguments: SinglePartition, ENSURE_REQUIREMENTS, [plan_id=2150]
(19) ShuffleQueryStage
Output [1]: [count#1908L]
Arguments: 1
(20) HashAggregate [codegen id : 2]
Input [1]: [count#1908L]
Keys: []
Functions [1]: [count(1)]
Aggregate Attributes [1]: [count(1)#1846L]
Results [1]: [count(1)#1846L AS count#1841L]
(21) ResultQueryStage
Output [1]: [count#1841L]
Arguments: 2
(22) HashAggregate
Input: []
Keys: []
Functions [1]: [partial_count(1)]
Aggregate Attributes [1]: [count#1907L]
Results [1]: [count#1908L]
(23) Exchange
Input [1]: [count#1908L]
Arguments: SinglePartition, ENSURE_REQUIREMENTS, [plan_id=2126]
(24) HashAggregate
Input [1]: [count#1908L]
Keys: []
Functions [1]: [count(1)]
Aggregate Attributes [1]: [count(1)#1846L]
Results [1]: [count(1)#1846L AS count#1841L]
(25) AdaptiveSparkPlan
Output [1]: [count#1841L]
Arguments: isFinalPlan=true