@@ -1840,20 +1840,20 @@ class LifecycleManager(val appUniqueId: String, val conf: CelebornConf) extends
1840
1840
cancelShuffleCallback = Some (callback)
1841
1841
}
1842
1842
1843
- @ volatile private var broadcastGetReducerFileGroupResponse
1843
+ @ volatile private var broadcastGetReducerFileGroupResponseCallback
1844
1844
: Option [java.util.function.BiFunction [Integer , GetReducerFileGroupResponse , Array [Byte ]]] =
1845
1845
None
1846
- def registerBroadcastGetReducerFileGroupResponse (call : java.util.function.BiFunction [
1846
+ def registerBroadcastGetReducerFileGroupResponseCallback (call : java.util.function.BiFunction [
1847
1847
Integer ,
1848
1848
GetReducerFileGroupResponse ,
1849
1849
Array [Byte ]]): Unit = {
1850
- broadcastGetReducerFileGroupResponse = Some (call)
1850
+ broadcastGetReducerFileGroupResponseCallback = Some (call)
1851
1851
}
1852
1852
1853
- @ volatile private var invalidatedBroadcastGetReducerFileGroupResponse : Option [Consumer [Integer ]] =
1853
+ @ volatile private var invalidatedBroadcastCallback : Option [Consumer [Integer ]] =
1854
1854
None
1855
- def registerInvalidatedBroadcastGetReducerFileGroupResponse (call : Consumer [Integer ]): Unit = {
1856
- invalidatedBroadcastGetReducerFileGroupResponse = Some (call)
1855
+ def registerInvalidatedBroadcastCallback (call : Consumer [Integer ]): Unit = {
1856
+ invalidatedBroadcastCallback = Some (call)
1857
1857
}
1858
1858
1859
1859
def invalidateLatestMaxLocsCache (shuffleId : Int ): Unit = {
@@ -1893,14 +1893,14 @@ class LifecycleManager(val appUniqueId: String, val conf: CelebornConf) extends
1893
1893
def broadcastGetReducerFileGroupResponse (
1894
1894
shuffleId : Int ,
1895
1895
response : GetReducerFileGroupResponse ): Option [Array [Byte ]] = {
1896
- broadcastGetReducerFileGroupResponse match {
1896
+ broadcastGetReducerFileGroupResponseCallback match {
1897
1897
case Some (c) => Option (c.apply(shuffleId, response))
1898
1898
case _ => None
1899
1899
}
1900
1900
}
1901
1901
1902
1902
private def invalidatedBroadcastGetReducerFileGroupResponse (shuffleId : Int ): Unit = {
1903
- invalidatedBroadcastGetReducerFileGroupResponse match {
1903
+ invalidatedBroadcastCallback match {
1904
1904
case Some (c) => c.accept(shuffleId)
1905
1905
case _ =>
1906
1906
}
0 commit comments