aboutsummaryrefslogtreecommitdiff
path: root/core/src/main/scala/spark/network/netty/ShuffleSender.scala
blob: cdf88b03a0840ce884d65543ccf01cecd35b4ef7 (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
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
/*
 * Licensed to the Apache Software Foundation (ASF) under one or more
 * contributor license agreements.  See the NOTICE file distributed with
 * this work for additional information regarding copyright ownership.
 * The ASF licenses this file to You under the Apache License, Version 2.0
 * (the "License"); you may not use this file except in compliance with
 * the License.  You may obtain a copy of the License at
 *
 *    http://www.apache.org/licenses/LICENSE-2.0
 *
 * Unless required by applicable law or agreed to in writing, software
 * distributed under the License is distributed on an "AS IS" BASIS,
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 * See the License for the specific language governing permissions and
 * limitations under the License.
 */

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)
  }
}