digraph G {
0 [labelType="html" label="<br><b>AdaptiveSparkPlan</b><br><br>"];
subgraph cluster1 {
isCluster="true";
label="WholeStageCodegen (3)\n \nduration: 1 ms";
2 [labelType="html" label="<b>HashAggregate</b><br><br>time in aggregation build: 0 ms<br>number of output rows: 1"];
}
3 [labelType="html" label="<b>Exchange</b><br><br>shuffle records written: 9<br>local merged chunks fetched: 0<br>shuffle write time total (min, med, max (stageId: taskId))<br>4 ms (0 ms, 0 ms, 0 ms (stage 36.0: task 117))<br>remote merged bytes read: 0.0 B<br>local merged blocks fetched: 0<br>corrupt merged block chunks: 0<br>remote merged reqs duration: 0 ms<br>remote merged blocks fetched: 0<br>records read: 9<br>local bytes read: 531.0 B<br>fetch wait time: 0 ms<br>remote bytes read: 0.0 B<br>merged fetch fallback count: 0<br>local blocks read: 9<br>remote merged chunks fetched: 0<br>remote blocks read: 0<br>data size total (min, med, max (stageId: taskId))<br>144.0 B (0.0 B, 16.0 B, 16.0 B (stage 36.0: task 109))<br>local merged 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>531.0 B (0.0 B, 59.0 B, 59.0 B (stage 36.0: task 109))"];
subgraph cluster4 {
isCluster="true";
label="WholeStageCodegen (2)\n \nduration: total (min, med, max (stageId: taskId))\n10.7 s (396 ms, 1.3 s, 1.4 s (stage 36.0: task 110))";
5 [labelType="html" label="<b>HashAggregate</b><br><br>spill size total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>time in aggregation build total (min, med, max (stageId: taskId))<br>10.7 s (395 ms, 1.3 s, 1.4 s (stage 36.0: task 110))<br>peak memory total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>number of output rows: 9<br>number of sort fallback tasks: 0<br>avg hash probes per key: 0"];
6 [labelType="html" label="<b>HashAggregate</b><br><br>spill size total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>time in aggregation build total (min, med, max (stageId: taskId))<br>10.4 s (380 ms, 1.3 s, 1.3 s (stage 36.0: task 110))<br>peak memory total (min, med, max (stageId: taskId))<br>1320.0 MiB (40.0 MiB, 160.0 MiB, 160.0 MiB (stage 36.0: task 109))<br>number of output rows: 4,988,053<br>number of sort fallback tasks: 0<br>avg hash probes per key (min, med, max (stageId: taskId)):<br>(1.6, 1.6, 1.7 (stage 36.0: task 111))"];
}
7 [labelType="html" label="<b>AQEShuffleRead</b><br><br>number of partitions: 9<br>partition data size total (min, med, max (stageId: taskId))<br>519.4 MiB (20.7 MiB, 62.3 MiB, 62.8 MiB (driver))<br>number of coalesced partitions: 9"];
8 [labelType="html" label="<b>Exchange</b><br><br>shuffle records written: 4,988,053<br>local merged chunks fetched: 0<br>shuffle write time total (min, med, max (stageId: taskId))<br>2.2 s (0 ms, 0 ms, 329 ms (stage 34.0: task 101))<br>remote merged bytes read total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>local merged blocks fetched: 0<br>corrupt merged block chunks: 0<br>remote merged reqs duration total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 36.0: task 109))<br>remote merged blocks fetched: 0<br>records read: 4,988,053<br>local bytes read total (min, med, max (stageId: taskId))<br>493.7 MiB (19.7 MiB, 59.2 MiB, 59.5 MiB (stage 36.0: task 113))<br>fetch wait time total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 36.0: task 109))<br>remote bytes read total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>merged fetch fallback count: 0<br>local blocks read: 72<br>remote merged chunks fetched: 0<br>remote blocks read: 0<br>data size total (min, med, max (stageId: taskId))<br>939.4 MiB (0.0 B, 0.0 B, 129.6 MiB (stage 34.0: task 102))<br>local merged bytes read total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>number of partitions: 200<br>remote reqs duration total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 36.0: task 109))<br>remote bytes read to disk total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 36.0: task 109))<br>shuffle bytes written total (min, med, max (stageId: taskId))<br>493.7 MiB (0.0 B, 0.0 B, 70.4 MiB (stage 34.0: task 101))"];
subgraph cluster9 {
isCluster="true";
label="WholeStageCodegen (1)\n \nduration: total (min, med, max (stageId: taskId))\n54.4 s (0 ms, 0 ms, 7.7 s (stage 34.0: task 104))";
10 [labelType="html" label="<b>HashAggregate</b><br><br>spill size total (min, med, max (stageId: taskId))<br>0.0 B (0.0 B, 0.0 B, 0.0 B (stage 34.0: task 101))<br>time in aggregation build total (min, med, max (stageId: taskId))<br>43.2 s (0 ms, 0 ms, 6.2 s (stage 34.0: task 102))<br>peak memory total (min, med, max (stageId: taskId))<br>1192.0 MiB (0.0 B, 0.0 B, 160.0 MiB (stage 34.0: task 101))<br>number of output rows: 4,988,053<br>number of sort fallback tasks: 0<br>avg hash probes per key (min, med, max (stageId: taskId)):<br>(1.6, 1.6, 1.6 (stage 34.0: task 101))"];
11 [labelType="html" label="<br><b>Project</b><br><br>"];
}
12 [labelType="html" label="<b>Scan csv </b><br><br>number of output rows: 4,988,053<br>number of files read: 5<br>metadata time total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 34.0: task 101))<br>size of files read total (min, med, max (stageId: taskId))<br>928.8 MiB (0.0 B, 0.0 B, 928.8 MiB (driver))"];
2->0;
3->2;
5->3;
6->5;
7->6;
8->7;
10->8;
11->10;
12->11;
}
13
AdaptiveSparkPlan isFinalPlan=true
HashAggregate(keys=[], functions=[count(1)])
WholeStageCodegen (3)
Exchange SinglePartition, ENSURE_REQUIREMENTS, [plan_id=480]
HashAggregate(keys=[], functions=[partial_count(1)])
HashAggregate(keys=[start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26], functions=[])
WholeStageCodegen (2)
AQEShuffleRead coalesced
Exchange hashpartitioning(start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26, 200), ENSURE_REQUIREMENTS, [plan_id=445]
HashAggregate(keys=[knownfloatingpointnormalized(normalizenanandzero(start_station_id#115)) AS start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, knownfloatingpointnormalized(normalizenanandzero(end_lat#27)) AS end_lat#27, ride_id#17, knownfloatingpointnormalized(normalizenanandzero(end_station_id#117)) AS end_station_id#117, ended_at#20, knownfloatingpointnormalized(normalizenanandzero(end_lng#28)) AS end_lng#28, started_at#19, rideable_type#18, knownfloatingpointnormalized(normalizenanandzero(start_lat#25)) AS start_lat#25, knownfloatingpointnormalized(normalizenanandzero(start_lng#26)) AS start_lng#26], functions=[])
Project [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, cast(start_station_id#22 as float) AS start_station_id#115, end_station_name#23, cast(end_station_id#24 as float) AS end_station_id#117, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
WholeStageCodegen (1)
FileScan csv [ride_id#17,rideable_type#18,started_at#19,ended_at#20,start_station_name#21,start_station_id#22,end_station_name#23,end_station_id#24,start_lat#25,start_lng#26,end_lat#27,end_lng#28,member_casual#29] Batched: false, DataFilters: [], Format: CSV, Location: InMemoryFileIndex(5 paths)[s3a://rzvde-g9-chernyshev-miron/raw/citibike_data/202507/202507-citibi..., PartitionFilters: [], PushedFilters: [], ReadSchema: struct<ride_id:string,rideable_type:string,started_at:timestamp,ended_at:timestamp,start_station_...
== Physical Plan ==
AdaptiveSparkPlan (19)
+- == Final Plan ==
* HashAggregate (11)
+- ShuffleQueryStage (10), Statistics(sizeInBytes=144.0 B, rowCount=9)
+- Exchange (9)
+- * HashAggregate (8)
+- * HashAggregate (7)
+- AQEShuffleRead (6)
+- ShuffleQueryStage (5), Statistics(sizeInBytes=939.4 MiB, rowCount=4.99E+6)
+- Exchange (4)
+- * HashAggregate (3)
+- * Project (2)
+- Scan csv (1)
+- == Initial Plan ==
HashAggregate (18)
+- Exchange (17)
+- HashAggregate (16)
+- HashAggregate (15)
+- Exchange (14)
+- HashAggregate (13)
+- Project (12)
+- Scan csv (1)
(1) Scan csv
Output [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, start_station_id#22, end_station_name#23, end_station_id#24, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
Batched: false
Location: InMemoryFileIndex [s3a://rzvde-g9-chernyshev-miron/raw/citibike_data/202507/202507-citibike-tripdata-part00.csv, ... 4 entries]
ReadSchema: struct<ride_id:string,rideable_type:string,started_at:timestamp,ended_at:timestamp,start_station_name:string,start_station_id:string,end_station_name:string,end_station_id:string,start_lat:double,start_lng:double,end_lat:double,end_lng:double,member_casual:string>
(2) Project [codegen id : 1]
Output [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, cast(start_station_id#22 as float) AS start_station_id#115, end_station_name#23, cast(end_station_id#24 as float) AS end_station_id#117, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
Input [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, start_station_id#22, end_station_name#23, end_station_id#24, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
(3) HashAggregate [codegen id : 1]
Input [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, start_station_id#115, end_station_name#23, end_station_id#117, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
Keys [13]: [knownfloatingpointnormalized(normalizenanandzero(start_station_id#115)) AS start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, knownfloatingpointnormalized(normalizenanandzero(end_lat#27)) AS end_lat#27, ride_id#17, knownfloatingpointnormalized(normalizenanandzero(end_station_id#117)) AS end_station_id#117, ended_at#20, knownfloatingpointnormalized(normalizenanandzero(end_lng#28)) AS end_lng#28, started_at#19, rideable_type#18, knownfloatingpointnormalized(normalizenanandzero(start_lat#25)) AS start_lat#25, knownfloatingpointnormalized(normalizenanandzero(start_lng#26)) AS start_lng#26]
Functions: []
Aggregate Attributes: []
Results [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
(4) Exchange
Input [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Arguments: hashpartitioning(start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26, 200), ENSURE_REQUIREMENTS, [plan_id=445]
(5) ShuffleQueryStage
Output [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Arguments: 0
(6) AQEShuffleRead
Input [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Arguments: coalesced
(7) HashAggregate [codegen id : 2]
Input [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Keys [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Functions: []
Aggregate Attributes: []
Results: []
(8) HashAggregate [codegen id : 2]
Input: []
Keys: []
Functions [1]: [partial_count(1)]
Aggregate Attributes [1]: [count#596L]
Results [1]: [count#597L]
(9) Exchange
Input [1]: [count#597L]
Arguments: SinglePartition, ENSURE_REQUIREMENTS, [plan_id=480]
(10) ShuffleQueryStage
Output [1]: [count#597L]
Arguments: 1
(11) HashAggregate [codegen id : 3]
Input [1]: [count#597L]
Keys: []
Functions [1]: [count(1)]
Aggregate Attributes [1]: [count(1)#593L]
Results [1]: [count(1)#593L AS count#594L]
(12) Project
Output [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, cast(start_station_id#22 as float) AS start_station_id#115, end_station_name#23, cast(end_station_id#24 as float) AS end_station_id#117, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
Input [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, start_station_id#22, end_station_name#23, end_station_id#24, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
(13) HashAggregate
Input [13]: [ride_id#17, rideable_type#18, started_at#19, ended_at#20, start_station_name#21, start_station_id#115, end_station_name#23, end_station_id#117, start_lat#25, start_lng#26, end_lat#27, end_lng#28, member_casual#29]
Keys [13]: [knownfloatingpointnormalized(normalizenanandzero(start_station_id#115)) AS start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, knownfloatingpointnormalized(normalizenanandzero(end_lat#27)) AS end_lat#27, ride_id#17, knownfloatingpointnormalized(normalizenanandzero(end_station_id#117)) AS end_station_id#117, ended_at#20, knownfloatingpointnormalized(normalizenanandzero(end_lng#28)) AS end_lng#28, started_at#19, rideable_type#18, knownfloatingpointnormalized(normalizenanandzero(start_lat#25)) AS start_lat#25, knownfloatingpointnormalized(normalizenanandzero(start_lng#26)) AS start_lng#26]
Functions: []
Aggregate Attributes: []
Results [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
(14) Exchange
Input [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Arguments: hashpartitioning(start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26, 200), ENSURE_REQUIREMENTS, [plan_id=423]
(15) HashAggregate
Input [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Keys [13]: [start_station_id#115, end_station_name#23, member_casual#29, start_station_name#21, end_lat#27, ride_id#17, end_station_id#117, ended_at#20, end_lng#28, started_at#19, rideable_type#18, start_lat#25, start_lng#26]
Functions: []
Aggregate Attributes: []
Results: []
(16) HashAggregate
Input: []
Keys: []
Functions [1]: [partial_count(1)]
Aggregate Attributes [1]: [count#596L]
Results [1]: [count#597L]
(17) Exchange
Input [1]: [count#597L]
Arguments: SinglePartition, ENSURE_REQUIREMENTS, [plan_id=427]
(18) HashAggregate
Input [1]: [count#597L]
Keys: []
Functions [1]: [count(1)]
Aggregate Attributes [1]: [count(1)#593L]
Results [1]: [count(1)#593L AS count#594L]
(19) AdaptiveSparkPlan
Output [1]: [count#594L]
Arguments: isFinalPlan=true