aboutsummaryrefslogtreecommitdiff
path: root/core/src/main/scala/spark/network/netty/ShuffleSender.scala
blob: d6fa4b1e807444bf573460b715df6e38d60a647b (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
package spark.network.netty

import java.io.File

import spark.Logging


private[spark] class ShuffleSender(portIn: Int, val pResolver: PathResolver) extends Logging {

  val server = new FileServer(pResolver, portIn)
  server.start()

  def stop() {
    server.stop()
  }

  def port: Int = server.getPort()
}


/**
 * An application for testing the shuffle sender as a standalone program.
 */
private[spark] object ShuffleSender {

  def main(args: Array[String]) {
    if (args.length < 3) {
      System.err.println(
        "Usage: ShuffleSender <port> <subDirsPerLocalDir> <list of shuffle_block_directories>")
      System.exit(1)
    }

    val port = args(0).toInt
    val subDirsPerLocalDir = args(1).toInt
    val localDirs = args.drop(2).map(new File(_))

    val pResovler = new PathResolver {
      override def getAbsolutePath(blockId: String): String = {
        if (!blockId.startsWith("shuffle_")) {
          throw new Exception("Block " + blockId + " is not a shuffle block")
        }
        // Figure out which local directory it hashes to, and which subdirectory in that
        val hash = math.abs(blockId.hashCode)
        val dirId = hash % localDirs.length
        val subDirId = (hash / localDirs.length) % subDirsPerLocalDir
        val subDir = new File(localDirs(dirId), "%02x".format(subDirId))
        val file = new File(subDir, blockId)
        return file.getAbsolutePath
      }
    }
    val sender = new ShuffleSender(port, pResovler)
  }
}