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: 1<br>local merged chunks fetched: 0<br>shuffle write time total (min, med, max (stageId: taskId))<br>0 ms (0 ms, 0 ms, 0 ms (stage 60.0: task 175))<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: 1<br>local bytes read: 72.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: 1<br>remote merged chunks fetched: 0<br>remote blocks read: 0<br>data size total (min, med, max (stageId: taskId))<br>40.0 B (0.0 B, 40.0 B, 40.0 B (stage 60.0: task 175))<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>72.0 B (0.0 B, 72.0 B, 72.0 B (stage 60.0: task 175))"];
subgraph cluster4 {
isCluster="true";
label="WholeStageCodegen (2)\n \nduration: 59 ms";
5 [labelType="html" label="<b>HashAggregate</b><br><br>spill size: 0.0 B<br>time in aggregation build: 58 ms<br>peak memory: 0.0 B<br>number of output rows: 1<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: 0.0 B<br>time in aggregation build: 48 ms<br>peak memory: 32.2 MiB<br>number of output rows: 4,374<br>number of sort fallback tasks: 0<br>avg hash probes per key: 1.1"];
}
7 [labelType="html" label="<b>AQEShuffleRead</b><br><br>number of partitions: 1<br>partition data size: 827.7 KiB<br>number of coalesced partitions: 1"];
8 [labelType="html" label="<b>Exchange</b><br><br>shuffle records written: 25,055<br>local merged chunks fetched: 0<br>shuffle write time total (min, med, max (stageId: taskId))<br>228 ms (0 ms, 26 ms, 33 ms (stage 58.0: task 172))<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: 25,055<br>local bytes read: 788.8 KiB<br>fetch wait time: 0 ms<br>remote bytes read: 0.0 B<br>merged fetch fallback count: 0<br>local blocks read: 8<br>remote merged chunks fetched: 0<br>remote blocks read: 0<br>data size total (min, med, max (stageId: taskId))<br>1757.3 KiB (0.0 B, 196.6 KiB, 298.1 KiB (stage 58.0: task 169))<br>local merged bytes read: 0.0 B<br>number of partitions: 200<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>788.8 KiB (0.0 B, 89.9 KiB, 127.4 KiB (stage 58.0: task 169))"];
subgraph cluster9 {
isCluster="true";
label="WholeStageCodegen (1)\n \nduration: total (min, med, max (stageId: taskId))\n26.8 s (0 ms, 3.5 s, 3.9 s (stage 58.0: task 168))";
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 60.0: task 175))<br>time in aggregation build total (min, med, max (stageId: taskId))<br>26.5 s (0 ms, 3.5 s, 3.8 s (stage 58.0: task 168))<br>peak memory total (min, med, max (stageId: taskId))<br>258.0 MiB (0.0 B, 32.2 MiB, 32.2 MiB (stage 58.0: task 167))<br>number of output rows: 25,055<br>number of sort fallback tasks: 0<br>avg hash probes per key (min, med, max (stageId: taskId)):<br>(1, 1, 1 (stage 58.0: task 167))"];
11 [labelType="html" label="<b>Expand</b><br><br>number of output rows: 19,952,212"];
}
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 60.0: task 175))<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(rideable_type#1092), count(start_station_name#1093), count(end_station_name#1094), count(member_casual#1095)])
WholeStageCodegen (3)
Exchange SinglePartition, ENSURE_REQUIREMENTS, [plan_id=805]
HashAggregate(keys=[], functions=[partial_count(rideable_type#1092) FILTER (WHERE (gid#1091 = 1)), partial_count(start_station_name#1093) FILTER (WHERE (gid#1091 = 2)), partial_count(end_station_name#1094) FILTER (WHERE (gid#1091 = 3)), partial_count(member_casual#1095) FILTER (WHERE (gid#1091 = 4))])
HashAggregate(keys=[rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091], functions=[])
WholeStageCodegen (2)
AQEShuffleRead coalesced
Exchange hashpartitioning(rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091, 200), ENSURE_REQUIREMENTS, [plan_id=770]
HashAggregate(keys=[rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091], functions=[])
Expand [[rideable_type#18, null, null, null, 1], [null, start_station_name#21, null, null, 2], [null, null, end_station_name#23, null, 3], [null, null, null, member_casual#29, 4]], [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
WholeStageCodegen (1)
FileScan csv [rideable_type#18,start_station_name#21,end_station_name#23,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<rideable_type:string,start_station_name:string,end_station_name:string,member_casual:string>
== Physical Plan ==
AdaptiveSparkPlan (19)
+- == Final Plan ==
* HashAggregate (11)
+- ShuffleQueryStage (10), Statistics(sizeInBytes=40.0 B, rowCount=1)
+- Exchange (9)
+- * HashAggregate (8)
+- * HashAggregate (7)
+- AQEShuffleRead (6)
+- ShuffleQueryStage (5), Statistics(sizeInBytes=1757.3 KiB, rowCount=2.51E+4)
+- Exchange (4)
+- * HashAggregate (3)
+- * Expand (2)
+- Scan csv (1)
+- == Initial Plan ==
HashAggregate (18)
+- Exchange (17)
+- HashAggregate (16)
+- HashAggregate (15)
+- Exchange (14)
+- HashAggregate (13)
+- Expand (12)
+- Scan csv (1)
(1) Scan csv
Output [4]: [rideable_type#18, start_station_name#21, end_station_name#23, 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<rideable_type:string,start_station_name:string,end_station_name:string,member_casual:string>
(2) Expand [codegen id : 1]
Input [4]: [rideable_type#18, start_station_name#21, end_station_name#23, member_casual#29]
Arguments: [[rideable_type#18, null, null, null, 1], [null, start_station_name#21, null, null, 2], [null, null, end_station_name#23, null, 3], [null, null, null, member_casual#29, 4]], [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
(3) HashAggregate [codegen id : 1]
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Keys [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Functions: []
Aggregate Attributes: []
Results [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
(4) Exchange
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Arguments: hashpartitioning(rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091, 200), ENSURE_REQUIREMENTS, [plan_id=770]
(5) ShuffleQueryStage
Output [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Arguments: 0
(6) AQEShuffleRead
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Arguments: coalesced
(7) HashAggregate [codegen id : 2]
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Keys [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Functions: []
Aggregate Attributes: []
Results [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
(8) HashAggregate [codegen id : 2]
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Keys: []
Functions [4]: [partial_count(rideable_type#1092) FILTER (WHERE (gid#1091 = 1)), partial_count(start_station_name#1093) FILTER (WHERE (gid#1091 = 2)), partial_count(end_station_name#1094) FILTER (WHERE (gid#1091 = 3)), partial_count(member_casual#1095) FILTER (WHERE (gid#1091 = 4))]
Aggregate Attributes [4]: [count#1096L, count#1097L, count#1098L, count#1099L]
Results [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
(9) Exchange
Input [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
Arguments: SinglePartition, ENSURE_REQUIREMENTS, [plan_id=805]
(10) ShuffleQueryStage
Output [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
Arguments: 1
(11) HashAggregate [codegen id : 3]
Input [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
Keys: []
Functions [4]: [count(rideable_type#1092), count(start_station_name#1093), count(end_station_name#1094), count(member_casual#1095)]
Aggregate Attributes [4]: [count(rideable_type#1092)#1071L, count(start_station_name#1093)#1072L, count(end_station_name#1094)#1073L, count(member_casual#1095)#1074L]
Results [4]: [toprettystring(count(rideable_type#1092)#1071L, Some(Etc/UTC)) AS toprettystring(rideable_type)#1083, toprettystring(count(start_station_name#1093)#1072L, Some(Etc/UTC)) AS toprettystring(start_station_name)#1084, toprettystring(count(end_station_name#1094)#1073L, Some(Etc/UTC)) AS toprettystring(end_station_name)#1085, toprettystring(count(member_casual#1095)#1074L, Some(Etc/UTC)) AS toprettystring(member_casual)#1086]
(12) Expand
Input [4]: [rideable_type#18, start_station_name#21, end_station_name#23, member_casual#29]
Arguments: [[rideable_type#18, null, null, null, 1], [null, start_station_name#21, null, null, 2], [null, null, end_station_name#23, null, 3], [null, null, null, member_casual#29, 4]], [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
(13) HashAggregate
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Keys [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Functions: []
Aggregate Attributes: []
Results [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
(14) Exchange
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Arguments: hashpartitioning(rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091, 200), ENSURE_REQUIREMENTS, [plan_id=748]
(15) HashAggregate
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Keys [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Functions: []
Aggregate Attributes: []
Results [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
(16) HashAggregate
Input [5]: [rideable_type#1092, start_station_name#1093, end_station_name#1094, member_casual#1095, gid#1091]
Keys: []
Functions [4]: [partial_count(rideable_type#1092) FILTER (WHERE (gid#1091 = 1)), partial_count(start_station_name#1093) FILTER (WHERE (gid#1091 = 2)), partial_count(end_station_name#1094) FILTER (WHERE (gid#1091 = 3)), partial_count(member_casual#1095) FILTER (WHERE (gid#1091 = 4))]
Aggregate Attributes [4]: [count#1096L, count#1097L, count#1098L, count#1099L]
Results [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
(17) Exchange
Input [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
Arguments: SinglePartition, ENSURE_REQUIREMENTS, [plan_id=752]
(18) HashAggregate
Input [4]: [count#1100L, count#1101L, count#1102L, count#1103L]
Keys: []
Functions [4]: [count(rideable_type#1092), count(start_station_name#1093), count(end_station_name#1094), count(member_casual#1095)]
Aggregate Attributes [4]: [count(rideable_type#1092)#1071L, count(start_station_name#1093)#1072L, count(end_station_name#1094)#1073L, count(member_casual#1095)#1074L]
Results [4]: [toprettystring(count(rideable_type#1092)#1071L, Some(Etc/UTC)) AS toprettystring(rideable_type)#1083, toprettystring(count(start_station_name#1093)#1072L, Some(Etc/UTC)) AS toprettystring(start_station_name)#1084, toprettystring(count(end_station_name#1094)#1073L, Some(Etc/UTC)) AS toprettystring(end_station_name)#1085, toprettystring(count(member_casual#1095)#1074L, Some(Etc/UTC)) AS toprettystring(member_casual)#1086]
(19) AdaptiveSparkPlan
Output [4]: [toprettystring(rideable_type)#1083, toprettystring(start_station_name)#1084, toprettystring(end_station_name)#1085, toprettystring(member_casual)#1086]
Arguments: isFinalPlan=true