We have a stateful multi-stage DStream application where each stage is mapped to a different Key. Beyond the first stage we are experiencing huge Shuffle read and write which consistently grows and becomes huge after about 20 windows or so.
val stageOnePartitionCount = 1000
val stageTwoPartitionCount = 1000
val stageThreePartitionCount = 1000
// Very little Shuffle for this stage consistent across all windows
val stageOneDataSet = dataSet
.map(data => (data.key1, data))
.reduceByKey(..., stageOnePartitionCount)
.mapWithState(
StateSpec
.function(...)
.initialState(...) // uses key1
.numPartitions(stageOnePartitionCount)
)
// Shuffle for this stage grows every window and becomes huge after about 20 windows or so..
val stageTwoDataSet = stageOneDataSet
.map(data => (data.key2, data))
.reduceByKey(..., stageTwoPartitionCount)
.mapWithState(
StateSpec
.function(...)
.initialState(...) // uses key2
.numPartitions(stageTwoPartitionCount)
)
// Shuffle for this stage grows every window and becomes huge after about 20 windows or so..
val stageThreeDataSet = stageTwoDataSet
.map(data => (data.key3, data))
.reduceByKey(..., stageThreePartitionCount)
.mapWithState(
StateSpec
.function(...)
.initialState(...) // uses key3
.numPartitions(stageThreePartitionCount)
)
stageThreeDataSet
.foreachRDD(...)
We are performing reduceByKey (to merge partial updates) operation before every mapWithState operation in every stage, where the partition-count matches between them, the intention being that at the end of the reduceByKey function over the incoming stream the data will already be in a node where it's state resides.
So what is triggering such exponential rise in Shuffle for the second and third stage with every passing window, but the first stage remains performant with very little Shuffle across all windows?
Note: Looking at the amount of data being Shuffled (read/write) for second and third stage, I was wondering whether if there is a chance that the stateful RDD (containing state data) is being Shuffled to co-locate them with the incoming stream which will be quite smaller in size than the accumulated state after 20 windows or so. If this is the case how do we solve the same.
Thanks for all the help!!