Skip to content

Commit b9118fd

Browse files
authored
[ZEPPELIN-5417] Unable to set conda env in pyspark (#4147)
* [ZEPPELIN-5417] Unable to set conda env in pyspark
1 parent e34cf0f commit b9118fd

4 files changed

Lines changed: 29 additions & 17 deletions

File tree

‎python/src/main/java/org/apache/zeppelin/python/PythonInterpreter.java‎

Lines changed: 18 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -80,16 +80,24 @@ public PythonInterpreter(Properties property) {
8080
public void open() throws InterpreterException {
8181
// try IPythonInterpreter first
8282
iPythonInterpreter = getIPythonInterpreter();
83-
if (getProperty("zeppelin.python.useIPython", "true").equals("true") &&
84-
StringUtils.isEmpty(
85-
iPythonInterpreter.checkKernelPrerequisite(getPythonExec()))) {
86-
try {
87-
iPythonInterpreter.open();
88-
LOGGER.info("IPython is available, Use IPythonInterpreter to replace PythonInterpreter");
89-
return;
90-
} catch (Exception e) {
91-
iPythonInterpreter = null;
92-
LOGGER.warn("Fail to open IPythonInterpreter", e);
83+
boolean useIPython = Boolean.parseBoolean(getProperty("zeppelin.python.useIPython", "true"));
84+
85+
LOGGER.info("zeppelin.python.useIPython: {}", useIPython);
86+
if (useIPython) {
87+
String checkKernelPrerequisiteResult = iPythonInterpreter.checkKernelPrerequisite(
88+
getPythonExec());
89+
if (StringUtils.isEmpty(checkKernelPrerequisiteResult)) {
90+
try {
91+
iPythonInterpreter.open();
92+
LOGGER.info("IPython is available, Use IPythonInterpreter to replace PythonInterpreter");
93+
return;
94+
} catch (Exception e) {
95+
iPythonInterpreter = null;
96+
LOGGER.warn("Fail to open IPythonInterpreter", e);
97+
}
98+
} else {
99+
LOGGER.info("IPython requirement is not met, checkKernelPrerequisiteResult: {}",
100+
checkKernelPrerequisiteResult);
93101
}
94102
}
95103

‎spark/interpreter/src/main/java/org/apache/zeppelin/spark/PySparkInterpreter.java‎

Lines changed: 9 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -171,18 +171,20 @@ public void setInterpreterContextInPython() {
171171
// spark.pyspark.driver.python > spark.pyspark.python > PYSPARK_DRIVER_PYTHON > PYSPARK_PYTHON
172172
@Override
173173
protected String getPythonExec() {
174-
if (!StringUtils.isBlank(getProperty("spark.pyspark.driver.python", ""))) {
175-
return properties.getProperty("spark.pyspark.driver.python");
174+
SparkConf sparkConf = getSparkConf();
175+
if (StringUtils.isNotBlank(sparkConf.get("spark.pyspark.driver.python", ""))) {
176+
return sparkConf.get("spark.pyspark.driver.python");
176177
}
177-
if (!StringUtils.isBlank(getProperty("spark.pyspark.python", ""))) {
178-
return properties.getProperty("spark.pyspark.python");
179-
}
180-
if (System.getenv("PYSPARK_PYTHON") != null) {
181-
return System.getenv("PYSPARK_PYTHON");
178+
if (StringUtils.isNotBlank(sparkConf.get("spark.pyspark.python", ""))) {
179+
return sparkConf.get("spark.pyspark.python");
182180
}
183181
if (System.getenv("PYSPARK_DRIVER_PYTHON") != null) {
184182
return System.getenv("PYSPARK_DRIVER_PYTHON");
185183
}
184+
if (System.getenv("PYSPARK_PYTHON") != null) {
185+
return System.getenv("PYSPARK_PYTHON");
186+
}
187+
186188
return "python";
187189
}
188190

‎spark/interpreter/src/test/java/org/apache/zeppelin/spark/IPySparkInterpreterTest.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,7 @@ protected void startInterpreter(Properties properties) throws InterpreterExcepti
9090
intpGroup.get("session_1").add(interpreter);
9191
interpreter.setInterpreterGroup(intpGroup);
9292

93+
pySparkInterpreter.open();
9394
interpreter.open();
9495
}
9596

‎zeppelin-jupyter-interpreter/src/main/java/org/apache/zeppelin/jupyter/JupyterKernelInterpreter.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -153,6 +153,7 @@ public void open() throws InterpreterException {
153153
* @return check result of checking kernel prerequisite.
154154
*/
155155
public String checkKernelPrerequisite(String pythonExec) {
156+
LOGGER.info("checkKernelPrerequisite using python executable: {}", pythonExec);
156157
ProcessBuilder processBuilder = new ProcessBuilder(pythonExec, "-m", "pip", "freeze");
157158
File stderrFile = null;
158159
File stdoutFile = null;

0 commit comments

Comments
 (0)