Repository: spark
Updated Branches:
  refs/heads/master 5c506cecb -> d3f07fd23


[SPARK-4688] Have a single shared network timeout in Spark

[SPARK-4688] Have a single shared network timeout in Spark

Author: Varun Saxena <[email protected]>
Author: varunsaxena <[email protected]>

Closes #3562 from varunsaxena/SPARK-4688 and squashes the following commits:

6e97f72 [Varun Saxena] [SPARK-4688] Single shared network timeout
cd783a2 [Varun Saxena] SPARK-4688
d6f8c29 [Varun Saxena] SCALA-4688
9562b15 [Varun Saxena] SPARK-4688
a75f014 [varunsaxena] SPARK-4688
594226c [varunsaxena] SPARK-4688


Project: http://git-wip-us.apache.org/repos/asf/spark/repo
Commit: http://git-wip-us.apache.org/repos/asf/spark/commit/d3f07fd2
Tree: http://git-wip-us.apache.org/repos/asf/spark/tree/d3f07fd2
Diff: http://git-wip-us.apache.org/repos/asf/spark/diff/d3f07fd2

Branch: refs/heads/master
Commit: d3f07fd23cc26a70f44c52e24445974d4885d58a
Parents: 5c506ce
Author: Varun Saxena <[email protected]>
Authored: Mon Jan 5 10:32:37 2015 -0800
Committer: Reynold Xin <[email protected]>
Committed: Mon Jan 5 10:32:37 2015 -0800

----------------------------------------------------------------------
 .../org/apache/spark/network/nio/ConnectionManager.scala  |  3 ++-
 .../apache/spark/storage/BlockManagerMasterActor.scala    |  7 +++++--
 core/src/main/scala/org/apache/spark/util/AkkaUtils.scala |  2 +-
 docs/configuration.md                                     | 10 ++++++++++
 .../java/org/apache/spark/network/util/TransportConf.java |  4 +++-
 5 files changed, 21 insertions(+), 5 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/spark/blob/d3f07fd2/core/src/main/scala/org/apache/spark/network/nio/ConnectionManager.scala
----------------------------------------------------------------------
diff --git 
a/core/src/main/scala/org/apache/spark/network/nio/ConnectionManager.scala 
b/core/src/main/scala/org/apache/spark/network/nio/ConnectionManager.scala
index 243b71c..98455c0 100644
--- a/core/src/main/scala/org/apache/spark/network/nio/ConnectionManager.scala
+++ b/core/src/main/scala/org/apache/spark/network/nio/ConnectionManager.scala
@@ -81,7 +81,8 @@ private[nio] class ConnectionManager(
   private val ackTimeoutMonitor =
     new HashedWheelTimer(Utils.namedThreadFactory("AckTimeoutMonitor"))
 
-  private val ackTimeout = 
conf.getInt("spark.core.connection.ack.wait.timeout", 60)
+  private val ackTimeout =
+    conf.getInt("spark.core.connection.ack.wait.timeout", 
conf.getInt("spark.network.timeout", 100))
 
   // Get the thread counts from the Spark Configuration.
   // 

http://git-wip-us.apache.org/repos/asf/spark/blob/d3f07fd2/core/src/main/scala/org/apache/spark/storage/BlockManagerMasterActor.scala
----------------------------------------------------------------------
diff --git 
a/core/src/main/scala/org/apache/spark/storage/BlockManagerMasterActor.scala 
b/core/src/main/scala/org/apache/spark/storage/BlockManagerMasterActor.scala
index 9cbda41..9d77cf2 100644
--- a/core/src/main/scala/org/apache/spark/storage/BlockManagerMasterActor.scala
+++ b/core/src/main/scala/org/apache/spark/storage/BlockManagerMasterActor.scala
@@ -52,8 +52,11 @@ class BlockManagerMasterActor(val isLocal: Boolean, conf: 
SparkConf, listenerBus
 
   private val akkaTimeout = AkkaUtils.askTimeout(conf)
 
-  val slaveTimeout = conf.getLong("spark.storage.blockManagerSlaveTimeoutMs",
-    math.max(conf.getInt("spark.executor.heartbeatInterval", 10000) * 3, 
45000))
+  val slaveTimeout = {
+    val defaultMs = math.max(conf.getInt("spark.executor.heartbeatInterval", 
10000) * 3, 45000)
+    val networkTimeout = conf.getInt("spark.network.timeout", defaultMs / 1000)
+    conf.getLong("spark.storage.blockManagerSlaveTimeoutMs", networkTimeout * 
1000)
+  }
 
   val checkTimeoutInterval = 
conf.getLong("spark.storage.blockManagerTimeoutIntervalMs", 60000)
 

http://git-wip-us.apache.org/repos/asf/spark/blob/d3f07fd2/core/src/main/scala/org/apache/spark/util/AkkaUtils.scala
----------------------------------------------------------------------
diff --git a/core/src/main/scala/org/apache/spark/util/AkkaUtils.scala 
b/core/src/main/scala/org/apache/spark/util/AkkaUtils.scala
index 8c2457f..64e3a54 100644
--- a/core/src/main/scala/org/apache/spark/util/AkkaUtils.scala
+++ b/core/src/main/scala/org/apache/spark/util/AkkaUtils.scala
@@ -65,7 +65,7 @@ private[spark] object AkkaUtils extends Logging {
 
     val akkaThreads   = conf.getInt("spark.akka.threads", 4)
     val akkaBatchSize = conf.getInt("spark.akka.batchSize", 15)
-    val akkaTimeout = conf.getInt("spark.akka.timeout", 100)
+    val akkaTimeout = conf.getInt("spark.akka.timeout", 
conf.getInt("spark.network.timeout", 100))
     val akkaFrameSize = maxFrameSizeBytes(conf)
     val akkaLogLifecycleEvents = 
conf.getBoolean("spark.akka.logLifecycleEvents", false)
     val lifecycleEvents = if (akkaLogLifecycleEvents) "on" else "off"

http://git-wip-us.apache.org/repos/asf/spark/blob/d3f07fd2/docs/configuration.md
----------------------------------------------------------------------
diff --git a/docs/configuration.md b/docs/configuration.md
index 9bb6499..7ada67f 100644
--- a/docs/configuration.md
+++ b/docs/configuration.md
@@ -819,6 +819,16 @@ Apart from these, the following properties are also 
available, and may be useful
   </td>
 </tr>
 <tr>
+  <td><code>spark.network.timeout</code></td>
+  <td>100</td>
+  <td>
+    Default timeout for all network interactions, in seconds. This config will 
be used in 
+    place of <code>spark.core.connection.ack.wait.timeout</code>, 
<code>spark.akka.timeout</code>, 
+    <code>spark.storage.blockManagerSlaveTimeoutMs</code> or 
<code>spark.shuffle.io.connectionTimeout</code>, 
+    if they are not configured.  
+  </td>
+</tr>
+<tr>
   <td><code>spark.akka.heartbeat.pauses</code></td>
   <td>6000</td>
   <td>

http://git-wip-us.apache.org/repos/asf/spark/blob/d3f07fd2/network/common/src/main/java/org/apache/spark/network/util/TransportConf.java
----------------------------------------------------------------------
diff --git 
a/network/common/src/main/java/org/apache/spark/network/util/TransportConf.java 
b/network/common/src/main/java/org/apache/spark/network/util/TransportConf.java
index 7c9adf5..e34382d 100644
--- 
a/network/common/src/main/java/org/apache/spark/network/util/TransportConf.java
+++ 
b/network/common/src/main/java/org/apache/spark/network/util/TransportConf.java
@@ -37,7 +37,9 @@ public class TransportConf {
 
   /** Connect timeout in milliseconds. Default 120 secs. */
   public int connectionTimeoutMs() {
-    return conf.getInt("spark.shuffle.io.connectionTimeout", 120) * 1000;
+    int timeout = 
+      conf.getInt("spark.shuffle.io.connectionTimeout", 
conf.getInt("spark.network.timeout", 100));
+    return timeout * 1000;
   }
 
   /** Number of concurrent connections between two nodes for fetching data. */


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to