digraph G {
0 [id="node0" labelType="html" label="<br><b>DeserializeToObject</b><br><br>" tooltip="DeserializeToObject createexternalrow(invoke(shardId#477721.toString()), static_invoke(java.lang.Long.valueOf(worklistShardItemId#477722L)), static_invoke(java.lang.Double.valueOf(qty#477724)), invoke(demandChannel#477725.toString()), invoke(demandStream#477726.toString()), mapobjects(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1), if (isnull(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1))) null else createexternalrow(invoke(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1).label.toString()), static_invoke(DateTimeUtils.toJavaTimestamp(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1).dateTime)), static_invoke(java.lang.Double.valueOf(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1).value)), StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true)), kpis#477727, Some(class scala.collection.mutable.ArraySeq)), StructField(shardId,StringType,true), StructField(worklistShardItemId,LongType,true), StructField(qty,DoubleType,true), StructField(demandChannel,StringType,true), StructField(demandStream,StringType,true), StructField(kpis,ArrayType(StructType(StructField(label,StringType,true),StructField(dateTime,TimestampType,true),StructField(value,DoubleType,true)),true),true)), obj#582934: org.apache.spark.sql.Row"];
1 [id="node1" labelType="html" label="<b>Exchange</b><br><br>number of partitions: 20" tooltip="Exchange hashpartitioning(worklistShardItemId#477722L, 20), REPARTITION_BY_NUM, [plan_id=246599]"];
subgraph cluster2 {
isCluster="true";
id="cluster2";
label="WholeStageCodegen (2)";
tooltip="WholeStageCodegen (2)";
3 [id="node3" labelType="html" label="<br><b>HashAggregate</b><br><br>" tooltip="HashAggregate(keys=[demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726], functions=[])"];
}
4 [id="node4" labelType="html" label="<b>Exchange</b><br><br>number of partitions: 20" tooltip="Exchange hashpartitioning(demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726, 20), ENSURE_REQUIREMENTS, [plan_id=246595]"];
5 [id="node5" labelType="html" label="<br><b>HashAggregate</b><br><br>" tooltip="HashAggregate(keys=[demandChannel#477725, shardId#477721, knownfloatingpointnormalized(transform(kpis#477727, lambdafunction(knownfloatingpointnormalized(if (isnull(lambda arg#582935)) null else named_struct(label, lambda arg#582935.label, dateTime, lambda arg#582935.dateTime, value, knownfloatingpointnormalized(normalizenanandzero(lambda arg#582935.value)))), lambda arg#582935, false))) AS kpis#477727, version#477723, knownfloatingpointnormalized(normalizenanandzero(qty#477724)) AS qty#477724, worklistShardItemId#477722L, demandStream#477726], functions=[])"];
subgraph cluster6 {
isCluster="true";
id="cluster6";
label="WholeStageCodegen (1)";
tooltip="WholeStageCodegen (1)";
7 [id="node7" labelType="html" label="<br><b>Scan ExistingRDD</b><br><br>" tooltip="Scan ExistingRDD[shardId#477721,worklistShardItemId#477722L,version#477723,qty#477724,demandChannel#477725,demandStream#477726,kpis#477727]"];
}
1->0;
3->1;
4->3;
5->4;
7->5;
}
== Physical Plan ==
DeserializeToObject (6)
+- Exchange (5)
+- * HashAggregate (4)
+- Exchange (3)
+- HashAggregate (2)
+- * Scan ExistingRDD (1)
(1) Scan ExistingRDD [codegen id : 1]
Output [7]: [shardId#477721, worklistShardItemId#477722L, version#477723, qty#477724, demandChannel#477725, demandStream#477726, kpis#477727]
Arguments: [shardId#477721, worklistShardItemId#477722L, version#477723, qty#477724, demandChannel#477725, demandStream#477726, kpis#477727], SQLExecutionRDD[27856] at start at FileStorageAdapterImpl.java:592, ExistingRDD, UnknownPartitioning(0)
(2) HashAggregate
Input [7]: [shardId#477721, worklistShardItemId#477722L, version#477723, qty#477724, demandChannel#477725, demandStream#477726, kpis#477727]
Keys [7]: [demandChannel#477725, shardId#477721, knownfloatingpointnormalized(transform(kpis#477727, lambdafunction(knownfloatingpointnormalized(if (isnull(lambda arg#582935)) null else named_struct(label, lambda arg#582935.label, dateTime, lambda arg#582935.dateTime, value, knownfloatingpointnormalized(normalizenanandzero(lambda arg#582935.value)))), lambda arg#582935, false))) AS kpis#477727, version#477723, knownfloatingpointnormalized(normalizenanandzero(qty#477724)) AS qty#477724, worklistShardItemId#477722L, demandStream#477726]
Functions: []
Aggregate Attributes: []
Results [7]: [demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726]
(3) Exchange
Input [7]: [demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726]
Arguments: hashpartitioning(demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726, 20), ENSURE_REQUIREMENTS, [plan_id=246595]
(4) HashAggregate [codegen id : 2]
Input [7]: [demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726]
Keys [7]: [demandChannel#477725, shardId#477721, kpis#477727, version#477723, qty#477724, worklistShardItemId#477722L, demandStream#477726]
Functions: []
Aggregate Attributes: []
Results [6]: [shardId#477721, worklistShardItemId#477722L, qty#477724, demandChannel#477725, demandStream#477726, kpis#477727]
(5) Exchange
Input [6]: [shardId#477721, worklistShardItemId#477722L, qty#477724, demandChannel#477725, demandStream#477726, kpis#477727]
Arguments: hashpartitioning(worklistShardItemId#477722L, 20), REPARTITION_BY_NUM, [plan_id=246599]
(6) DeserializeToObject
Input [6]: [shardId#477721, worklistShardItemId#477722L, qty#477724, demandChannel#477725, demandStream#477726, kpis#477727]
Arguments: createexternalrow(invoke(shardId#477721.toString()), static_invoke(java.lang.Long.valueOf(worklistShardItemId#477722L)), static_invoke(java.lang.Double.valueOf(qty#477724)), invoke(demandChannel#477725.toString()), invoke(demandStream#477726.toString()), mapobjects(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1), if (isnull(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1))) null else createexternalrow(invoke(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1).label.toString()), static_invoke(DateTimeUtils.toJavaTimestamp(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1).dateTime)), static_invoke(java.lang.Double.valueOf(lambdavariable(MapObject, StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true), true, -1).value)), StructField(label,StringType,true), StructField(dateTime,TimestampType,true), StructField(value,DoubleType,true)), kpis#477727, Some(class scala.collection.mutable.ArraySeq)), StructField(shardId,StringType,true), StructField(worklistShardItemId,LongType,true), StructField(qty,DoubleType,true), StructField(demandChannel,StringType,true), StructField(demandStream,StringType,true), StructField(kpis,ArrayType(StructType(StructField(label,StringType,true),StructField(dateTime,TimestampType,true),StructField(value,DoubleType,true)),true),true)), obj#582934: org.apache.spark.sql.Row