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 (5)";
tooltip="WholeStageCodegen (5)";
2 [id="node2" labelType="html" label="<br><b>SerializeFromObject</b><br><br>" tooltip="SerializeFromObject [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#4352]"];
3 [id="node3" labelType="html" label="<br><b>MapElements</b><br><br>" tooltip="MapElements org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802109d40@7bea9f07, obj#4351: java.lang.String"];
4 [id="node4" labelType="html" label="<br><b>DeserializeToObject</b><br><br>" tooltip="DeserializeToObject invoke(path#4215.toString()), obj#4350: java.lang.String"];
5 [id="node5" labelType="html" label="<br><b>EmptyRelation</b><br><br>" tooltip="EmptyRelation [plan_id=6282]"];
6 [id="node6" labelType="html" label="<br><b>Project</b><br><br>" tooltip="Project [path#4215]"];
7 [id="node7" labelType="html" label="<br><b>EmptyRelation</b><br><br>" tooltip="EmptyRelation Filter (count#4217L = 1)"];
8 [id="node8" labelType="html" label="<br><b>Filter</b><br><br>" tooltip="Filter (count#4217L = 1)"];
9 [id="node9" labelType="html" label="<br><b>EmptyRelation</b><br><br>" tooltip="EmptyRelation LogicalQueryStage Aggregate [path#4215], [path#4215, count(1) AS count#4217L], HashAggregate(keys=[path#4215], functions=[count(1)])"];
10 [id="node10" labelType="html" label="<br><b>LogicalQueryStage</b><br><br>" tooltip="LogicalQueryStage Aggregate [path#4215], [path#4215, count(1) AS count#4217L], HashAggregate(keys=[path#4215], functions=[count(1)])"];
11 [id="node11" labelType="html" label="<br><b>HashAggregate</b><br><br>" tooltip="HashAggregate(keys=[path#4215], functions=[count(1)])"];
}
12 [id="node12" labelType="html" label="<b>Exchange</b><br><br>shuffle records written: 0<br>data size: 0.0 B<br>shuffle write time: 0 ms<br>number of partitions: 200<br>shuffle bytes written: 0.0 B" tooltip="Exchange hashpartitioning(path#4215, 200), ENSURE_REQUIREMENTS, [plan_id=6234]"];
subgraph cluster13 {
isCluster="true";
id="cluster13";
label="WholeStageCodegen (4)\n \nduration: total (min, med, max (stageId: taskId))\n1 ms (0 ms, 0 ms, 1 ms (stage 251.0: task 23824))";
tooltip="WholeStageCodegen (4)";
14 [id="node14" labelType="html" label="<b>HashAggregate</b><br><br>spill size: 0.0 B<br>time in aggregation build total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 251.0: task 23787))<br>peak memory total (min, med, max (stageId: taskId))<br>50.0 MiB (256.0 KiB, 256.0 KiB, 256.0 KiB (stage 251.0: task 23787))<br>number of output rows: 0<br>number of sort fallback tasks: 0<br>avg hash probes per key: 0" tooltip="HashAggregate(keys=[path#4215], functions=[partial_count(1)])"];
15 [id="node15" 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.commands.VacuumCommand$FileNameAndSize, true])).path()))) AS path#4215]"];
}
16 [id="node16" labelType="html" label="<br><b>MapPartitions</b><br><br>" tooltip="MapPartitions org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x00000008020dc000@6b8c52a1, obj#4212: org.apache.spark.sql.delta.commands.VacuumCommand$FileNameAndSize"];
17 [id="node17" labelType="html" label="<br><b>DeserializeToObject</b><br><br>" tooltip="DeserializeToObject newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), obj#4209: org.apache.spark.sql.delta.SerializableFileStatus"];
subgraph cluster18 {
isCluster="true";
id="cluster18";
label="WholeStageCodegen (3)\n \nduration: total (min, med, max (stageId: taskId))\n181 ms (0 ms, 1 ms, 1 ms (stage 251.0: task 23787))";
tooltip="WholeStageCodegen (3)";
19 [id="node19" labelType="html" label="<b>Filter</b><br><br>number of output rows: 0" tooltip="Filter ((modificationTime#3957L < 1788027607910) OR isDir#3956)"];
}
20 [id="node20" labelType="html" label="<b>InMemoryTableScan</b><br><br>number of output rows: 0" tooltip="InMemoryTableScan [path#3954, length#3955L, isDir#3956, modificationTime#3957L], [((modificationTime#3957L < 1788027607910) OR isDir#3956)]"];
21 [id="node21" labelType="html" label="<br><b>AdaptiveSparkPlan</b><br><br>" tooltip="AdaptiveSparkPlan isFinalPlan=true"];
subgraph cluster22 {
isCluster="true";
id="cluster22";
label="WholeStageCodegen (2)\n \nduration: 0 ms";
tooltip="WholeStageCodegen (2)";
23 [id="node23" 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#3954, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#3955L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#3956, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#3957L]"];
}
24 [id="node24" labelType="html" label="<br><b>MapGroups</b><br><br>" tooltip="MapGroups org.apache.spark.sql.internal.UDFAdaptors$$$Lambda/0x000000080209d5c0@754784c3, invoke(value#3948.toString()), newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), [value#3948], [path#3938, length#3939L, isDir#3940, modificationTime#3941L], obj#3953: org.apache.spark.sql.delta.SerializableFileStatus"];
subgraph cluster25 {
isCluster="true";
id="cluster25";
label="WholeStageCodegen (1)\n \nduration: 0 ms";
tooltip="WholeStageCodegen (1)";
26 [id="node26" 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#3948 ASC NULLS FIRST], false, 0"];
}
27 [id="node27" 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#3948, 200), ENSURE_REQUIREMENTS, [plan_id=5143]"];
28 [id="node28" 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#3938, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#3939L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#3940, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#3941L], [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#3948]"];
29 [id="node29" labelType="html" label="<br><b>MapElements</b><br><br>" tooltip="MapElements org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802093950@5e5aa004, obj#3937: org.apache.spark.sql.delta.SerializableFileStatus"];
30 [id="node30" labelType="html" label="<b>Scan</b><br><br>number of output rows: 0" tooltip="Scan[obj#3922]"];
2->0;
3->2;
4->3;
5->4;
6->5;
7->6;
8->7;
9->8;
10->9;
11->10;
12->11;
14->12;
15->14;
16->15;
17->16;
19->17;
20->19;
21->20;
23->21;
24->23;
26->24;
27->26;
28->27;
29->28;
30->29;
}
== Physical Plan ==
AdaptiveSparkPlan (43)
+- == Final Plan ==
ResultQueryStage (5), Statistics(sizeInBytes=8.0 EiB)
+- * SerializeFromObject (4)
+- * MapElements (3)
+- * DeserializeToObject (2)
+- * EmptyRelation (1)
+- Project (unknown)
+- EmptyRelation (unknown)
+- == Initial Plan ==
SerializeFromObject (42)
+- MapElements (41)
+- DeserializeToObject (40)
+- SortMergeJoin LeftAnti (39)
:- Sort (30)
: +- Project (29)
: +- Filter (28)
: +- HashAggregate (27)
: +- Exchange (26)
: +- HashAggregate (25)
: +- SerializeFromObject (24)
: +- MapPartitions (23)
: +- DeserializeToObject (22)
: +- Filter (21)
: +- InMemoryTableScan (6)
: +- InMemoryRelation (7)
: +- AdaptiveSparkPlan (20)
+- == Final Plan ==
ResultQueryStage (16)
+- * SerializeFromObject (15)
+- MapGroups (14)
+- * Sort (13)
+- ShuffleQueryStage (12), Statistics(sizeInBytes=0.0 B, rowCount=0)
+- Exchange (11)
+- AppendColumnsWithObject (10)
+- MapElements (9)
+- Scan (8)
+- == Initial Plan ==
SerializeFromObject (19)
+- MapGroups (18)
+- Sort (17)
+- Exchange (11)
+- AppendColumnsWithObject (10)
+- MapElements (9)
+- Scan (8)
+- Sort (38)
+- Exchange (37)
+- Project (36)
+- Filter (35)
+- SerializeFromObject (34)
+- MapPartitions (33)
+- DeserializeToObject (32)
+- Scan ExistingRDD Delta Table State #0 - hdlfs://065c2b6f-41bb-4a20-b7a2-eb965f9cb071.files.hdl.prod-eu20.hanacloud.ondemand.com:443/crp-workload-determination-service/in/workload-dupl-lock-v3/_delta_log (31)
(1) EmptyRelation [codegen id : 5]
Output [1]: [path#4215]
Arguments: [plan_id=6282]
(2) DeserializeToObject [codegen id : 5]
Input [1]: [path#4215]
Arguments: invoke(path#4215.toString()), obj#4350: java.lang.String
(3) MapElements [codegen id : 5]
Input [1]: [obj#4350]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802109d40@7bea9f07, obj#4351: java.lang.String
(4) SerializeFromObject [codegen id : 5]
Input [1]: [obj#4351]
Arguments: [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#4352]
(5) ResultQueryStage
Output [1]: [value#4352]
Arguments: 3
(6) InMemoryTableScan
Output [4]: [path#3954, length#3955L, isDir#3956, modificationTime#3957L]
Arguments: [path#3954, length#3955L, isDir#3956, modificationTime#3957L], [((modificationTime#3957L < 1788027607910) OR isDir#3956)]
(7) InMemoryRelation
Arguments: [path#3954, length#3955L, isDir#3956, modificationTime#3957L], StorageLevel(disk, memory, deserialized, 1 replicas)
(8) Scan
Output [1]: [obj#3922]
Arguments: obj#3922: org.apache.spark.sql.delta.SerializableFileStatus, MapPartitionsRDD[462] at $anonfun$recordDeltaOperationInternal$1 at DatabricksLogging.scala:128
(9) MapElements
Input [1]: [obj#3922]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802093950@5e5aa004, obj#3937: org.apache.spark.sql.delta.SerializableFileStatus
(10) AppendColumnsWithObject
Input [1]: [obj#3937]
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#3938, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#3939L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#3940, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#3941L], [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#3948]
(11) Exchange
Input [5]: [path#3938, length#3939L, isDir#3940, modificationTime#3941L, value#3948]
Arguments: hashpartitioning(value#3948, 200), ENSURE_REQUIREMENTS, [plan_id=5143]
(12) ShuffleQueryStage
Output [5]: [path#3938, length#3939L, isDir#3940, modificationTime#3941L, value#3948]
Arguments: 0
(13) Sort [codegen id : 1]
Input [5]: [path#3938, length#3939L, isDir#3940, modificationTime#3941L, value#3948]
Arguments: [value#3948 ASC NULLS FIRST], false, 0
(14) MapGroups
Input [5]: [path#3938, length#3939L, isDir#3940, modificationTime#3941L, value#3948]
Arguments: org.apache.spark.sql.internal.UDFAdaptors$$$Lambda/0x000000080209d5c0@754784c3, invoke(value#3948.toString()), newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), [value#3948], [path#3938, length#3939L, isDir#3940, modificationTime#3941L], obj#3953: org.apache.spark.sql.delta.SerializableFileStatus
(15) SerializeFromObject [codegen id : 2]
Input [1]: [obj#3953]
Arguments: [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#3954, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#3955L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#3956, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#3957L]
(16) ResultQueryStage
Output [4]: [path#3954, length#3955L, isDir#3956, modificationTime#3957L]
Arguments: 1
(17) Sort
Input [5]: [path#3938, length#3939L, isDir#3940, modificationTime#3941L, value#3948]
Arguments: [value#3948 ASC NULLS FIRST], false, 0
(18) MapGroups
Input [5]: [path#3938, length#3939L, isDir#3940, modificationTime#3941L, value#3948]
Arguments: org.apache.spark.sql.internal.UDFAdaptors$$$Lambda/0x000000080209d5c0@754784c3, invoke(value#3948.toString()), newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), [value#3948], [path#3938, length#3939L, isDir#3940, modificationTime#3941L], obj#3953: org.apache.spark.sql.delta.SerializableFileStatus
(19) SerializeFromObject
Input [1]: [obj#3953]
Arguments: [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).path()))) AS path#3954, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).length()) AS length#3955L, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).isDir()) AS isDir#3956, invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.SerializableFileStatus, true])).modificationTime()) AS modificationTime#3957L]
(20) AdaptiveSparkPlan
Output [4]: [path#3954, length#3955L, isDir#3956, modificationTime#3957L]
Arguments: isFinalPlan=true
(21) Filter
Input [4]: [path#3954, length#3955L, isDir#3956, modificationTime#3957L]
Condition : ((modificationTime#3957L < 1788027607910) OR isDir#3956)
(22) DeserializeToObject
Input [4]: [path#3954, length#3955L, isDir#3956, modificationTime#3957L]
Arguments: newInstance(class org.apache.spark.sql.delta.SerializableFileStatus), obj#4209: org.apache.spark.sql.delta.SerializableFileStatus
(23) MapPartitions
Input [1]: [obj#4209]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x00000008020dc000@6b8c52a1, obj#4212: org.apache.spark.sql.delta.commands.VacuumCommand$FileNameAndSize
(24) SerializeFromObject
Input [1]: [obj#4212]
Arguments: [static_invoke(UTF8String.fromString(invoke(knownnotnull(assertnotnull(input[0, org.apache.spark.sql.delta.commands.VacuumCommand$FileNameAndSize, true])).path()))) AS path#4215]
(25) HashAggregate
Input [1]: [path#4215]
Keys [1]: [path#4215]
Functions [1]: [partial_count(1)]
Aggregate Attributes [1]: [count#4290L]
Results [2]: [path#4215, count#4292L]
(26) Exchange
Input [2]: [path#4215, count#4292L]
Arguments: hashpartitioning(path#4215, 200), ENSURE_REQUIREMENTS, [plan_id=6085]
(27) HashAggregate
Input [2]: [path#4215, count#4292L]
Keys [1]: [path#4215]
Functions [1]: [count(1)]
Aggregate Attributes [1]: [count(1)#4221L]
Results [2]: [path#4215, count(1)#4221L AS count#4217L]
(28) Filter
Input [2]: [path#4215, count#4217L]
Condition : (count#4217L = 1)
(29) Project
Output [1]: [path#4215]
Input [2]: [path#4215, count#4217L]
(30) Sort
Input [1]: [path#4215]
Arguments: [path#4215 ASC NULLS FIRST], false, 0
(31) Scan ExistingRDD Delta Table State #0 - hdlfs://065c2b6f-41bb-4a20-b7a2-eb965f9cb071.files.hdl.prod-eu20.hanacloud.ondemand.com:443/crp-workload-determination-service/in/workload-dupl-lock-v3/_delta_log [codegen id : 1]
Output [10]: [txn#2402, add#2403, remove#2404, metaData#2405, protocol#2406, cdc#2407, checkpointMetadata#2408, sidecar#2409, domainMetadata#2410, commitInfo#2411]
Arguments: [txn#2402, add#2403, remove#2404, metaData#2405, protocol#2406, cdc#2407, checkpointMetadata#2408, sidecar#2409, domainMetadata#2410, commitInfo#2411], Delta Table State #0 - hdlfs://065c2b6f-41bb-4a20-b7a2-eb965f9cb071.files.hdl.prod-eu20.hanacloud.ondemand.com:443/crp-workload-determination-service/in/workload-dupl-lock-v3/_delta_log MapPartitionsRDD[283] at $anonfun$recordDeltaOperationInternal$1 at DatabricksLogging.scala:128, ExistingRDD, UnknownPartitioning(0)
(32) DeserializeToObject
Input [10]: [txn#2402, add#2403, remove#2404, metaData#2405, protocol#2406, cdc#2407, checkpointMetadata#2408, sidecar#2409, domainMetadata#2410, commitInfo#2411]
Arguments: newInstance(class org.apache.spark.sql.delta.actions.SingleAction), obj#3914: org.apache.spark.sql.delta.actions.SingleAction
(33) MapPartitions
Input [1]: [obj#3914]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x000000080208f3e8@4d3dda6d, obj#3915: java.lang.String
(34) SerializeFromObject
Input [1]: [obj#3915]
Arguments: [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#3916]
(35) Filter
Input [1]: [value#3916]
Condition : isnotnull(value#3916)
(36) Project
Output [1]: [value#3916 AS path#3917]
Input [1]: [value#3916]
(37) Exchange
Input [1]: [path#3917]
Arguments: hashpartitioning(path#3917, 200), ENSURE_REQUIREMENTS, [plan_id=6091]
(38) Sort
Input [1]: [path#3917]
Arguments: [path#3917 ASC NULLS FIRST], false, 0
(39) SortMergeJoin
Left keys [1]: [path#4215]
Right keys [1]: [path#3917]
Join type: LeftAnti
Join condition: None
(40) DeserializeToObject
Input [1]: [path#4215]
Arguments: invoke(path#4215.toString()), obj#4350: java.lang.String
(41) MapElements
Input [1]: [obj#4350]
Arguments: org.apache.spark.sql.delta.commands.VacuumCommand$$$Lambda/0x0000000802109d40@7bea9f07, obj#4351: java.lang.String
(42) SerializeFromObject
Input [1]: [obj#4351]
Arguments: [static_invoke(UTF8String.fromString(input[0, java.lang.String, true])) AS value#4352]
(43) AdaptiveSparkPlan
Output [1]: [value#4352]
Arguments: isFinalPlan=true