We read every piece of feedback, and take your input very seriously.
To see all available qualifiers, see our documentation.
There was an error while loading. Please reload this page.
1 parent 6325fc1 commit 8911f44Copy full SHA for 8911f44
core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala
@@ -303,7 +303,9 @@ private[spark] object PythonRDD extends Logging {
303
// remember the broadcasts sent to each worker
304
private val workerBroadcasts = new mutable.WeakHashMap[Socket, mutable.Set[Long]]()
305
private def getWorkerBroadcasts(worker: Socket) = {
306
- workerBroadcasts.getOrElseUpdate(worker, new mutable.HashSet[Long]())
+ synchronized {
307
+ workerBroadcasts.getOrElseUpdate(worker, new mutable.HashSet[Long]())
308
+ }
309
}
310
311
/**
0 commit comments