Skip to content

Commit

Permalink
track broadcasts for each worker
Browse files Browse the repository at this point in the history
  • Loading branch information
davies committed Sep 4, 2014
1 parent 8d2f08c commit 6123d0f
Show file tree
Hide file tree
Showing 2 changed files with 31 additions and 6 deletions.
30 changes: 26 additions & 4 deletions core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import java.nio.charset.Charset
import java.util.{List => JList, ArrayList => JArrayList, Map => JMap, Collections}

import scala.collection.JavaConversions._
import scala.collection.mutable
import scala.language.existentials
import scala.reflect.ClassTag
import scala.util.{Try, Success, Failure}
Expand Down Expand Up @@ -193,11 +194,26 @@ private[spark] class PythonRDD(
PythonRDD.writeUTF(include, dataOut)
}
// Broadcast variables
dataOut.writeInt(broadcastVars.length)
val bids = PythonRDD.getWorkerBroadcasts(worker)
val nbids = broadcastVars.map(_.id).toSet
// number of different broadcasts
val cnt = bids.diff(nbids).size + nbids.diff(bids).size
dataOut.writeInt(cnt)
for (bid <- bids) {
if (!nbids.contains(bid)) {
// remove the broadcast from worker
dataOut.writeLong(-bid)
bids.remove(bid)
}
}
for (broadcast <- broadcastVars) {
dataOut.writeLong(broadcast.id)
dataOut.writeInt(broadcast.value.length)
dataOut.write(broadcast.value)
if (!bids.contains(broadcast.id)) {
// send new broadcast
dataOut.writeLong(broadcast.id)
dataOut.writeInt(broadcast.value.length)
dataOut.write(broadcast.value)
bids.add(broadcast.id)
}
}
dataOut.flush()
// Serialized command:
Expand Down Expand Up @@ -275,6 +291,12 @@ private object SpecialLengths {
private[spark] object PythonRDD extends Logging {
val UTF8 = Charset.forName("UTF-8")

// remember the broadcasts sent to each worker
private val workerBroadcasts = new mutable.WeakHashMap[Socket, mutable.Set[Long]]()
private def getWorkerBroadcasts(worker: Socket) = {
workerBroadcasts.getOrElseUpdate(worker, new mutable.HashSet[Long]())
}

/**
* Adapter for calling SparkContext#runJob from Python.
*
Expand Down
7 changes: 5 additions & 2 deletions python/pyspark/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,8 +69,11 @@ def main(infile, outfile):
ser = CompressedSerializer(pickleSer)
for _ in range(num_broadcast_variables):
bid = read_long(infile)
value = ser._read_with_length(infile)
_broadcastRegistry[bid] = Broadcast(bid, value)
if bid > 0:
value = ser._read_with_length(infile)
_broadcastRegistry[bid] = Broadcast(bid, value)
else:
_broadcastRegistry.pop(-bid, None)

command = pickleSer._read_with_length(infile)
(func, deserializer, serializer) = command
Expand Down

0 comments on commit 6123d0f

Please sign in to comment.