acvictor commented on code in PR #13072: URL: https://github.com/apache/gluten/pull/13072#discussion_r4070767870
########## gluten-substrait/src/main/scala/org/apache/spark/shuffle/sort/ColumnarShuffleManager.scala: ########## @@ -51,11 +51,22 @@ class ColumnarShuffleManager(conf: SparkConf) import ColumnarShuffleManager._ private lazy val shuffleExecutorComponents = loadShuffleExecutorComponents(conf) - override val shuffleBlockResolver = new IndexShuffleBlockResolver(conf) - /** A mapping from shuffle ids to the number of mappers producing output for those shuffles. */ + /** + * A mapping from shuffle ids to the task ids of mappers producing output for those shuffles. + * + * Must be declared before `shuffleBlockResolver`: Scala initializes vals in declaration order, so + * the resolver would otherwise capture `null`. + */ private[this] val taskIdMapsForShuffle = new ConcurrentHashMap[Int, OpenHashSet[Long]]() + // The resolver must share this map rather than allocate its own. It records blocks migrated in + // during executor decommissioning, and `unregisterShuffle` reads the same map to delete the + // corresponding map output. Passed positionally because Spark 4.0+ gives `_blockManager` no + // default value. + override val shuffleBlockResolver = + new IndexShuffleBlockResolver(conf, null, taskIdMapsForShuffle) Review Comment: Thanks @marin-ma! Added a comment. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
