How is duration being calculated in a Spark Structured Streaming UI?

Viewed 170

We have a Spark SQL job that we would like to optimize. We are trying to figure out which part of our pipeline is slower/faster.

In the attached SQL query graph, there are 3 WholeStageCodegen boxes, all with the same duration: 2.9s, 2.9s, 2.9s. See the below picture:

enter image description here

But if we check the Stage graph, it shows 3 seconds for the total stage. See the below picture:

enter image description here

So the durations in the WholeStageCodegen boxes do not add up, it seems that these durations refer to the sum of the whole stage. Do we miss something here? Is there a way to figure out the duration for the individual boxes?

Sometimes there is some difference in the duration, but not more than 0.1s, examples:

  • 18.3s, 18.3s, 18.4s
  • 968ms, 967ms, 1.0s The Stage duration is always as much as one of the WholeStageCodegen's duration, or at most 0.1-0.3sec larger.

How can one figure out the duration for each of the WholeStageCodegen parts, and is that actually measured? I suspect that Spark would have to trace individual operations as units of generated functions. Is that measurement actually performed there, or are these numbers more like a placeholder for a feature that does not exist?

0 Answers
Related