diff options
author | Davies Liu <davies@databricks.com> | 2015-04-02 12:18:33 -0700 |
---|---|---|
committer | Josh Rosen <joshrosen@databricks.com> | 2015-04-02 12:18:33 -0700 |
commit | 0cce5451adfc6bf4661bcf67aca3db26376455fe (patch) | |
tree | 145c47ad89bf382ce722fbfdfc10e16dae0a55ef | |
parent | 424e987dfebbbaa37f4496d44090d469a931ce76 (diff) | |
download | spark-0cce5451adfc6bf4661bcf67aca3db26376455fe.tar.gz spark-0cce5451adfc6bf4661bcf67aca3db26376455fe.tar.bz2 spark-0cce5451adfc6bf4661bcf67aca3db26376455fe.zip |
[SPARK-6667] [PySpark] remove setReuseAddress
The reused address on server side had caused the server can not acknowledge the connected connections, remove it.
This PR will retry once after timeout, it also add a timeout at client side.
Author: Davies Liu <davies@databricks.com>
Closes #5324 from davies/collect_hang and squashes the following commits:
e5a51a2 [Davies Liu] remove setReuseAddress
7977c2f [Davies Liu] do retry on client side
b838f35 [Davies Liu] retry after timeout
-rw-r--r-- | core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala | 1 | ||||
-rw-r--r-- | python/pyspark/rdd.py | 1 |
2 files changed, 1 insertions, 1 deletions
diff --git a/core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala b/core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala index 19f4c95fca..36cf2af085 100644 --- a/core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala +++ b/core/src/main/scala/org/apache/spark/api/python/PythonRDD.scala @@ -605,7 +605,6 @@ private[spark] object PythonRDD extends Logging { */ private def serveIterator[T](items: Iterator[T], threadName: String): Int = { val serverSocket = new ServerSocket(0, 1) - serverSocket.setReuseAddress(true) // Close the socket if no connection in 3 seconds serverSocket.setSoTimeout(3000) diff --git a/python/pyspark/rdd.py b/python/pyspark/rdd.py index c337a43c8a..2d05611321 100644 --- a/python/pyspark/rdd.py +++ b/python/pyspark/rdd.py @@ -113,6 +113,7 @@ def _parse_memory(s): def _load_from_socket(port, serializer): sock = socket.socket() + sock.settimeout(3) try: sock.connect(("localhost", port)) rf = sock.makefile("rb", 65536) |