diff options
author | Matei Zaharia <matei@eecs.berkeley.edu> | 2013-09-02 18:38:12 -0700 |
---|---|---|
committer | Matei Zaharia <matei@eecs.berkeley.edu> | 2013-09-02 18:38:12 -0700 |
commit | a106ed8b97e707b36818c11d1d7211fa28636178 (patch) | |
tree | 5ce12b04c710bd8e776c31bc3c8cef63f3313622 /python/pyspark/worker.py | |
parent | 2ce200bf7f7a38afbcacf3303ca2418e49bdbe2a (diff) | |
parent | 59218bdd4996a13116009e3669b1b875be23a694 (diff) | |
download | spark-a106ed8b97e707b36818c11d1d7211fa28636178.tar.gz spark-a106ed8b97e707b36818c11d1d7211fa28636178.tar.bz2 spark-a106ed8b97e707b36818c11d1d7211fa28636178.zip |
Merge remote-tracking branch 'old/master'
Diffstat (limited to 'python/pyspark/worker.py')
-rw-r--r-- | python/pyspark/worker.py | 11 |
1 files changed, 7 insertions, 4 deletions
diff --git a/python/pyspark/worker.py b/python/pyspark/worker.py index 695f6dfb84..d63c2aaef7 100644 --- a/python/pyspark/worker.py +++ b/python/pyspark/worker.py @@ -21,6 +21,7 @@ Worker that receives input from Piped RDD. import os import sys import time +import socket import traceback from base64 import standard_b64decode # CloudPickler needs to be imported so that depicklers are registered using the @@ -94,7 +95,9 @@ def main(infile, outfile): if __name__ == '__main__': - # Redirect stdout to stderr so that users must return values from functions. - old_stdout = os.fdopen(os.dup(1), 'w') - os.dup2(2, 1) - main(sys.stdin, old_stdout) + # Read a local port to connect to from stdin + java_port = int(sys.stdin.readline()) + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.connect(("127.0.0.1", java_port)) + sock_file = sock.makefile("a+", 65536) + main(sock_file, sock_file) |