This is an automated email from the ASF dual-hosted git repository.
dianfu pushed a commit to branch release-1.17
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/release-1.17 by this push:
new 44e6cfb87c1 [FLINK-30962][python] Improve the error message during
launching py4j gateway server
44e6cfb87c1 is described below
commit 44e6cfb87c1b2a5f4106df61cd06c4870d4802f8
Author: Juntao Hu <[email protected]>
AuthorDate: Wed Feb 8 18:05:30 2023 +0800
[FLINK-30962][python] Improve the error message during launching py4j
gateway server
This closes #21894.
---
flink-python/pyflink/java_gateway.py | 6 +++++-
flink-python/pyflink/pyflink_gateway_server.py | 2 +-
2 files changed, 6 insertions(+), 2 deletions(-)
diff --git a/flink-python/pyflink/java_gateway.py
b/flink-python/pyflink/java_gateway.py
index bdcb57e8094..4bd1c3c858e 100644
--- a/flink-python/pyflink/java_gateway.py
+++ b/flink-python/pyflink/java_gateway.py
@@ -109,7 +109,11 @@ def launch_gateway():
time.sleep(0.1)
if not os.path.isfile(conn_info_file):
- raise Exception("Java gateway process exited before sending its
port number")
+ stderr_info = p.stderr.read().decode('utf-8')
+ raise RuntimeError(
+ "Java gateway process exited before sending its port
number.\nStderr:\n"
+ + stderr_info
+ )
with open(conn_info_file, "rb") as info:
gateway_port = struct.unpack("!I", info.read(4))[0]
diff --git a/flink-python/pyflink/pyflink_gateway_server.py
b/flink-python/pyflink/pyflink_gateway_server.py
index f7947d7897c..fdb791120ef 100644
--- a/flink-python/pyflink/pyflink_gateway_server.py
+++ b/flink-python/pyflink/pyflink_gateway_server.py
@@ -261,7 +261,7 @@ def launch_gateway_server_process(env, args):
signal.signal(signal.SIGINT, signal.SIG_IGN)
preexec_fn = preexec_func
return Popen(list(filter(lambda c: len(c) != 0, command)),
- stdin=PIPE, preexec_fn=preexec_fn, env=env)
+ stdin=PIPE, stderr=PIPE, preexec_fn=preexec_fn, env=env)
if __name__ == "__main__":