From f16422f297a37cf9283bbe878366cc7418ef1120 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Sun, 20 Nov 2016 20:30:40 +0530 Subject: [PATCH 01/22] Ability to view spark job urls in each paragraph --- .../zeppelin/spark/SparkInterpreter.java | 48 +++++++++++++++++-- .../zeppelin/spark/SparkRInterpreter.java | 7 ++- .../zeppelin/spark/SparkSqlInterpreter.java | 11 ++--- .../zeppelin/spark/ZeppelinContext.java | 22 ++++++++- .../interpreter/remote/RemoteEventClient.java | 10 ++++ .../remote/RemoteEventClientWrapper.java | 3 ++ .../remote/RemoteInterpreterEventClient.java | 4 ++ .../remote/RemoteInterpreterEventPoller.java | 20 ++++++-- .../RemoteInterpreterProcessListener.java | 2 + .../thrift/RemoteInterpreterEventType.java | 5 +- .../thrift/RemoteInterpreterService.thrift | 3 +- .../RemoteInterpreterOutputTestStream.java | 6 +++ .../scheduler/RemoteSchedulerTest.java | 4 ++ .../zeppelin/rest/InterpreterRestApi.java | 10 +++- .../zeppelin/server/ZeppelinServer.java | 2 +- .../zeppelin/socket/NotebookServer.java | 39 +++++++++++++++ .../notebook/paragraph/paragraph-control.html | 14 ++++++ .../paragraph/paragraph.controller.js | 5 +- .../interpreter/InterpreterSetting.java | 19 ++++++++ .../apache/zeppelin/notebook/Notebook.java | 1 + .../apache/zeppelin/notebook/Paragraph.java | 22 +++++++++ 21 files changed, 234 insertions(+), 23 deletions(-) diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java index 0584a302c55..7d3caf8a85d 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java @@ -42,6 +42,8 @@ import org.apache.spark.scheduler.ActiveJob; import org.apache.spark.scheduler.DAGScheduler; import org.apache.spark.scheduler.Pool; +import org.apache.spark.scheduler.SparkListenerApplicationEnd; +import org.apache.spark.scheduler.SparkListenerJobStart; import org.apache.spark.sql.SQLContext; import org.apache.spark.ui.SparkUI; import org.apache.spark.ui.jobs.JobProgressListener; @@ -57,6 +59,7 @@ import org.apache.zeppelin.interpreter.util.InterpreterOutputStream; import org.apache.zeppelin.resource.ResourcePool; import org.apache.zeppelin.resource.WellKnownResourceName; +import org.apache.zeppelin.interpreter.remote.RemoteEventClientWrapper; import org.apache.zeppelin.interpreter.thrift.InterpreterCompletion; import org.apache.zeppelin.scheduler.Scheduler; import org.apache.zeppelin.scheduler.SchedulerFactory; @@ -112,7 +115,7 @@ public class SparkInterpreter extends Interpreter { private InterpreterOutputStream out; private SparkDependencyResolver dep; - private String sparkUrl; + private static String sparkUrl; /** * completer - org.apache.spark.repl.SparkJLineCompletion (scala 2.10) @@ -156,7 +159,43 @@ public boolean isSparkContextInitialized() { } static JobProgressListener setupListeners(SparkContext context) { - JobProgressListener pl = new JobProgressListener(context.getConf()); + JobProgressListener pl = new JobProgressListener(context.getConf()) { + @Override + public synchronized void onJobStart(SparkListenerJobStart jobStart) { + super.onJobStart(jobStart); + int jobId = jobStart.jobId(); + String jobGroupId = jobStart.properties().getProperty("spark.jobGroup.id"); + String jobUrl = getJobUrl(jobId); + String noteId = getNoteId(jobGroupId); + String paragraphId = getParagraphId(jobGroupId); + if (jobUrl != null && noteId != null && paragraphId != null) { + RemoteEventClientWrapper eventClient = ZeppelinContext.getEventClient(); + Map infos = new java.util.HashMap<>(); + infos.put("jobUrl", jobUrl); + eventClient.onParaInfosReceived(noteId, paragraphId, infos); + } + } + + private String getJobUrl(int jobId) { + String jobUrl = null; + if (sparkUrl != null) { + jobUrl = sparkUrl + "/jobs/job?id=" + jobId; + } + return jobUrl; + } + + private String getNoteId(String jobgroupId) { + int indexOf = jobgroupId.indexOf("-"); + int secondIndex = jobgroupId.indexOf("-", indexOf + 1); + return jobgroupId.substring(indexOf + 1, secondIndex); + } + + private String getParagraphId(String jobgroupId) { + int indexOf = jobgroupId.indexOf("-"); + int secondIndex = jobgroupId.indexOf("-", indexOf + 1); + return jobgroupId.substring(secondIndex + 1, jobgroupId.length()); + } + }; try { Object listenerBus = context.getClass().getMethod("listenerBus").invoke(context); @@ -973,6 +1012,7 @@ public void populateSparkWebUrl(InterpreterContext ctx) { infos.put("url", sparkUrl); logger.info("Sending metainfos to Zeppelin server: {}", infos.toString()); if (ctx != null && ctx.getClient() != null) { + getZeppelinContext().setEventClient(ctx.getClient()); ctx.getClient().onMetaInfosReceived(infos); } } @@ -1105,8 +1145,8 @@ public Object getLastObject() { return obj; } - String getJobGroup(InterpreterContext context){ - return "zeppelin-" + context.getParagraphId(); + String getJobGroup(InterpreterContext context) { + return "zeppelin-" + context.getNoteId() + "-" + context.getParagraphId(); } /** diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java index 8f3e93c024c..53beea854a5 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java @@ -102,7 +102,12 @@ String getJobGroup(InterpreterContext context){ @Override public InterpreterResult interpret(String lines, InterpreterContext interpreterContext) { - getSparkInterpreter().populateSparkWebUrl(interpreterContext); + SparkInterpreter sparkInterpreter = getSparkInterpreter(); + sparkInterpreter.populateSparkWebUrl(interpreterContext); + + String jobGroup = sparkInterpreter.getJobGroup(interpreterContext); + sparkInterpreter.getSparkContext().setJobGroup(jobGroup, "Zeppelin", false); + String imageWidth = getProperty("zeppelin.R.image.width"); String[] sl = lines.split("\n"); diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java index e6fe137273a..81a02234381 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java @@ -45,10 +45,6 @@ public class SparkSqlInterpreter extends Interpreter { Logger logger = LoggerFactory.getLogger(SparkSqlInterpreter.class); AtomicInteger num = new AtomicInteger(0); - private String getJobGroup(InterpreterContext context){ - return "zeppelin-" + context.getParagraphId(); - } - private int maxResult; public SparkSqlInterpreter(Properties property) { @@ -105,7 +101,7 @@ public InterpreterResult interpret(String st, InterpreterContext context) { sc.setLocalProperty("spark.scheduler.pool", null); } - sc.setJobGroup(getJobGroup(context), "Zeppelin", false); + sc.setJobGroup(sparkInterpreter.getJobGroup(context), "Zeppelin", false); Object rdd = null; try { // method signature of sqlc.sql() is changed @@ -134,10 +130,11 @@ public InterpreterResult interpret(String st, InterpreterContext context) { @Override public void cancel(InterpreterContext context) { - SQLContext sqlc = getSparkInterpreter().getSQLContext(); + SparkInterpreter sparkInterpreter = getSparkInterpreter(); + SQLContext sqlc = sparkInterpreter.getSQLContext(); SparkContext sc = sqlc.sparkContext(); - sc.cancelJobGroup(getJobGroup(context)); + sc.cancelJobGroup(sparkInterpreter.getJobGroup(context)); } @Override diff --git a/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java b/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java index d1234dfd987..620543d5829 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java @@ -46,6 +46,7 @@ import org.apache.zeppelin.interpreter.InterpreterException; import org.apache.zeppelin.interpreter.InterpreterHookRegistry; import org.apache.zeppelin.interpreter.RemoteWorksController; +import org.apache.zeppelin.interpreter.remote.RemoteEventClientWrapper; import org.apache.zeppelin.spark.dep.SparkDependencyResolver; import org.apache.zeppelin.resource.Resource; import org.apache.zeppelin.resource.ResourcePool; @@ -61,6 +62,7 @@ public class ZeppelinContext { // Map interpreter class name (to be used by hook registry) from // given replName in parapgraph private static final Map interpreterClassMap; + private static RemoteEventClientWrapper eventClient; static { interpreterClassMap = new HashMap<>(); interpreterClassMap.put("spark", "org.apache.zeppelin.spark.SparkInterpreter"); @@ -221,7 +223,8 @@ public static String showDF(SparkContext sc, Object df, int maxResult) { Object[] rows = null; Method take; - String jobGroup = "zeppelin-" + interpreterContext.getParagraphId(); + String jobGroup = "zeppelin-" + interpreterContext.getNoteId() + "-" + + interpreterContext.getParagraphId(); sc.setJobGroup(jobGroup, "Zeppelin", false); try { @@ -930,4 +933,21 @@ public ResourceSet getAll() { return resourcePool.getAll(); } + /** + * Get the event client + */ + @ZeppelinApi + public static RemoteEventClientWrapper getEventClient() { + return eventClient; + } + + /** + * Set event client + */ + @ZeppelinApi + public void setEventClient(RemoteEventClientWrapper eventClient) { + if (ZeppelinContext.eventClient == null) { + ZeppelinContext.eventClient = eventClient; + } + } } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClient.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClient.java index 3585a59eae4..5015a3f27df 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClient.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClient.java @@ -1,5 +1,6 @@ package org.apache.zeppelin.interpreter.remote; +import java.util.HashMap; import java.util.Map; /** @@ -21,4 +22,13 @@ public void onMetaInfosReceived(Map infos) { client.onMetaInfosReceived(infos); } + @Override + public void onParaInfosReceived(String noteId, String paragraphId, Map infos) { + Map paraInfos = new HashMap(infos); + paraInfos.put("noteId", noteId); + paraInfos.put("paraId", paragraphId); + client.onParaInfosReceived(paraInfos); + } + + } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClientWrapper.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClientWrapper.java index 339f7714a21..bf36cd6354f 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClientWrapper.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteEventClientWrapper.java @@ -12,4 +12,7 @@ public interface RemoteEventClientWrapper { public void onMetaInfosReceived(Map infos); + public void onParaInfosReceived(String noteId, String paragraphId, + Map infos); + } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java index 606d35f60ee..4b721f5ad60 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java @@ -470,6 +470,10 @@ public void onMetaInfosReceived(Map infos) { gson.toJson(infos))); } + public void onParaInfosReceived(Map infos) { + sendEvent(new RemoteInterpreterEvent(RemoteInterpreterEventType.PARA_INFOS, + gson.toJson(infos))); + } /** * Wait for eventQueue becomes empty */ diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java index e794140e8bd..0ecb2ab9e28 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java @@ -243,10 +243,18 @@ public void run() { Map metaInfos = gson.fromJson(event.getData(), new TypeToken>() { }.getType()); - String id = interpreterGroup.getId(); - int indexOfColon = id.indexOf(":"); - String settingId = id.substring(0, indexOfColon); + String settingId = getInterpreterSettingId(); listener.onMetaInfosReceived(settingId, metaInfos); + } else if (event.getType() == RemoteInterpreterEventType.PARA_INFOS) { + Map paraInfos = gson.fromJson(event.getData(), + new TypeToken>() { + }.getType()); + String noteId = paraInfos.get("noteId"); + String paraId = paraInfos.get("paraId"); + String settingId = getInterpreterSettingId(); + if (noteId != null && paraId != null && settingId != null) { + listener.onParaInfosReceived(noteId, paraId, settingId, paraInfos); + } } logger.debug("Event from remote process {}", event.getType()); } catch (Exception e) { @@ -334,6 +342,12 @@ public void onError() { } } + private String getInterpreterSettingId() { + String id = interpreterGroup.getId(); + int indexOfColon = id.indexOf(":"); + return id.substring(0, indexOfColon); + } + private void sendResourcePoolResponseGetAll(ResourceSet resourceSet) { Client client = null; boolean broken = false; diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java index 66b08c95a1d..0e9dc5128dc 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java @@ -40,4 +40,6 @@ public interface RemoteWorksEventListener { public void onFinished(Object resultObject); public void onError(); } + public void onParaInfosReceived(String noteId, String paragraphId, + String interpreterSettingId, Map metaInfos); } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java index 7ca406c6709..9e5a0b49f07 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java @@ -43,7 +43,8 @@ public enum RemoteInterpreterEventType implements org.apache.thrift.TEnum { APP_STATUS_UPDATE(12), META_INFOS(13), REMOTE_ZEPPELIN_SERVER_RESOURCE(14), - RESOURCE_INVOKE_METHOD(15); + RESOURCE_INVOKE_METHOD(15), + PARA_INFOS(16); private final int value; @@ -94,6 +95,8 @@ public static RemoteInterpreterEventType findByValue(int value) { return REMOTE_ZEPPELIN_SERVER_RESOURCE; case 15: return RESOURCE_INVOKE_METHOD; + case 16: + return PARA_INFOS; default: return null; } diff --git a/zeppelin-interpreter/src/main/thrift/RemoteInterpreterService.thrift b/zeppelin-interpreter/src/main/thrift/RemoteInterpreterService.thrift index 08a15ad9702..fc09adea5d4 100644 --- a/zeppelin-interpreter/src/main/thrift/RemoteInterpreterService.thrift +++ b/zeppelin-interpreter/src/main/thrift/RemoteInterpreterService.thrift @@ -56,7 +56,8 @@ enum RemoteInterpreterEventType { APP_STATUS_UPDATE = 12, META_INFOS = 13, REMOTE_ZEPPELIN_SERVER_RESOURCE = 14, - RESOURCE_INVOKE_METHOD = 15 + RESOURCE_INVOKE_METHOD = 15, + PARA_INFOS = 16 } diff --git a/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterOutputTestStream.java b/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterOutputTestStream.java index e3dc6b4c1b3..3f865cb370d 100644 --- a/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterOutputTestStream.java +++ b/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterOutputTestStream.java @@ -182,4 +182,10 @@ public void onGetParagraphRunners(String noteId, String paragraphId, RemoteWorks public void onRemoteRunParagraph(String noteId, String ParagraphID) throws Exception { } + + @Override + public void onParaInfosReceived(String noteId, String paragraphId, + String interpreterSettingId, Map metaInfos) { + } + } diff --git a/zeppelin-interpreter/src/test/java/org/apache/zeppelin/scheduler/RemoteSchedulerTest.java b/zeppelin-interpreter/src/test/java/org/apache/zeppelin/scheduler/RemoteSchedulerTest.java index d7b2007e73c..ebb51004285 100644 --- a/zeppelin-interpreter/src/test/java/org/apache/zeppelin/scheduler/RemoteSchedulerTest.java +++ b/zeppelin-interpreter/src/test/java/org/apache/zeppelin/scheduler/RemoteSchedulerTest.java @@ -355,6 +355,10 @@ public void onGetParagraphRunners(String noteId, String paragraphId, RemoteWorks @Override public void onRemoteRunParagraph(String noteId, String PsaragraphID) throws Exception { + } + @Override + public void onParaInfosReceived(String noteId, String paragraphId, + String interpreterSettingId, Map metaInfos) { } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java index 09280074375..57eb8514a31 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java @@ -51,6 +51,7 @@ import org.apache.zeppelin.rest.message.NewInterpreterSettingRequest; import org.apache.zeppelin.rest.message.UpdateInterpreterSettingRequest; import org.apache.zeppelin.server.JsonResponse; +import org.apache.zeppelin.socket.NotebookServer; /** * Interpreter Rest API @@ -61,14 +62,17 @@ public class InterpreterRestApi { private static final Logger logger = LoggerFactory.getLogger(InterpreterRestApi.class); private InterpreterFactory interpreterFactory; + private NotebookServer notebookServer; Gson gson = new Gson(); public InterpreterRestApi() { } - public InterpreterRestApi(InterpreterFactory interpreterFactory) { + public InterpreterRestApi(InterpreterFactory interpreterFactory, + NotebookServer notebookWsServer) { this.interpreterFactory = interpreterFactory; + this.notebookServer = notebookWsServer; } /** @@ -179,18 +183,20 @@ public Response removeSetting(@PathParam("settingId") String settingId) throws I @ZeppelinApi public Response restartSetting(String message, @PathParam("settingId") String settingId) { logger.info("Restart interpreterSetting {}, msg={}", settingId, message); + + InterpreterSetting setting = interpreterFactory.get(settingId); try { RestartInterpreterRequest request = gson.fromJson(message, RestartInterpreterRequest.class); String noteId = request == null ? null : request.getNoteId(); interpreterFactory.restart(settingId, noteId, SecurityUtils.getPrincipal()); + notebookServer.clearParagraphRuntimeInfo(setting); } catch (InterpreterException e) { logger.error("Exception in InterpreterRestApi while restartSetting ", e); return new JsonResponse<>(Status.NOT_FOUND, e.getMessage(), ExceptionUtils.getStackTrace(e)) .build(); } - InterpreterSetting setting = interpreterFactory.get(settingId); if (setting == null) { return new JsonResponse<>(Status.NOT_FOUND, "", settingId).build(); } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java index 371d0a131d0..045558f45a4 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java @@ -374,7 +374,7 @@ public Set getSingletons() { HeliumRestApi heliumApi = new HeliumRestApi(helium, notebook); singletons.add(heliumApi); - InterpreterRestApi interpreterApi = new InterpreterRestApi(replFactory); + InterpreterRestApi interpreterApi = new InterpreterRestApi(replFactory, notebookWsServer); singletons.add(interpreterApi); CredentialRestApi credentialApi = new CredentialRestApi(credentials); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 68b015d0bac..4b783b9d610 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -2314,4 +2314,43 @@ private void broadcastToWatchers(String noteId, String subject, Message message) } } } + + @Override + public void onParaInfosReceived(String noteId, String paragraphId, + String interpreterSettingId, Map metaInfos) { + Note note = notebook().getNote(noteId); + if (note != null) { + Paragraph paragraph = note.getParagraph(paragraphId); + if (paragraph != null) { + InterpreterSetting setting = notebook().getInterpreterFactory() + .get(interpreterSettingId); + setting.addNoteToPara(noteId, paragraphId); + metaInfos.remove("noteId"); + metaInfos.remove("paraId"); + paragraph.updateRuntimeInfos(metaInfos); + broadcast(note.getId(), new Message(OP.PARAGRAPH).put("paragraph", paragraph)); + } + } + } + + public void clearParagraphRuntimeInfo(InterpreterSetting setting) { + Map> noteIdAndParaMap = setting.getNoteIdAndParaMap(); + if (noteIdAndParaMap != null && !noteIdAndParaMap.isEmpty()) { + for (String noteId : noteIdAndParaMap.keySet()) { + Set paraIdSet = noteIdAndParaMap.get(noteId); + if (paraIdSet != null && !paraIdSet.isEmpty()) { + for (String paraId : paraIdSet) { + Note note = notebook().getNote(noteId); + if (note != null) { + Paragraph paragraph = note.getParagraph(paraId); + paragraph.clearRuntimeInfo(); + if (paragraph != null) { + broadcast(noteId, new Message(OP.PARAGRAPH).put("paragraph", paragraph)); + } + } + } + } + } + } + } } diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html index 351fb5ffecc..91d094f3682 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html @@ -31,6 +31,20 @@ tooltip="Cancel (Ctrl+{{ (isMac ? 'Option' : 'Alt') }}+c)" ng-click="cancelParagraph(paragraph)" ng-show="paragraph.status=='RUNNING' || paragraph.status=='PENDING'"> + + Spark job + + + Spark Jobs + + + diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js index 342d41f37d2..5742ab1e03b 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js @@ -1083,7 +1083,8 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat isEmpty(newPara.results) !== isEmpty(oldPara.results) || newPara.errorMessage !== oldPara.errorMessage || !angular.equals(newPara.settings, oldPara.settings) || - !angular.equals(newPara.config, oldPara.config))) + !angular.equals(newPara.config, oldPara.config) || + !angular.equals(newPara.runtimeInfos, oldPara.runtimeInfos))) } $scope.updateAllScopeTexts = function(oldPara, newPara) { @@ -1126,6 +1127,7 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat $scope.paragraph.results = newPara.results; } $scope.paragraph.settings = newPara.settings; + $scope.paragraph.runtimeInfos = newPara.runtimeInfos; if ($scope.editor) { $scope.editor.setReadOnly($scope.isRunning(newPara)); } @@ -1139,7 +1141,6 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat $scope.paragraph.config = newPara.config; } }; - $scope.updateParagraph = function(oldPara, newPara, updateCallback) { // 1. get status, refreshed const statusChanged = (newPara.status !== oldPara.status); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index bd7d664841e..07b98e36991 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -25,6 +25,7 @@ import java.util.List; import java.util.Map; import java.util.Properties; +import java.util.Set; import java.util.concurrent.locks.ReentrantReadWriteLock; import com.google.gson.annotations.SerializedName; @@ -47,6 +48,7 @@ public class InterpreterSetting { // always be null in case of InterpreterSettingRef private String group; private transient Map infos; + private transient Map> noteIdToParaIdsetMap; /** * properties can be either Properties or Map @@ -394,11 +396,28 @@ public Map getInfos() { return infos; } +<<<<<<< b1e06919fada4240f5440430f2c27c02e4e5626a public InterpreterRunner getInterpreterRunner() { return interpreterRunner; } public void setInterpreterRunner(InterpreterRunner interpreterRunner) { this.interpreterRunner = interpreterRunner; +======= + public void addNoteToPara(String noteId, String paraId) { + if(noteIdToParaIdsetMap == null) { + noteIdToParaIdsetMap = new HashMap<>(); + } + Set paraIdSet = noteIdToParaIdsetMap.get(noteId); + if(paraIdSet == null) { + paraIdSet = new HashSet<>(); + noteIdToParaIdsetMap.put(noteId, paraIdSet); + } + paraIdSet.add(paraId); + } + + public Map> getNoteIdAndParaMap() { + return noteIdToParaIdsetMap; +>>>>>>> Ability to view spark job urls in each paragraph } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java index 8b946f2c344..f94e5dd571d 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java @@ -480,6 +480,7 @@ public Note loadNoteFromRepo(String id, AuthenticationInfo subject) { if (p.getDateFinished() != null && lastUpdatedDate.before(p.getDateFinished())) { lastUpdatedDate = p.getDateFinished(); } + p.clearRuntimeInfo(); } Map> savedObjects = note.getAngularObjects(); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java index 27a707137a1..c7c5fd5e58d 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java @@ -70,6 +70,7 @@ public class Paragraph extends Job implements Serializable, Cloneable { // For backward compatibility of note.json format after ZEPPELIN-212 Object result; + private Map> runtimeInfos; /** * Applicaiton states in this paragraph @@ -676,4 +677,25 @@ private boolean isValidInterpreter(String replName) { return false; } } + + public void updateRuntimeInfos(Map infos) { + if (this.runtimeInfos == null) { + this.runtimeInfos = new HashMap>(); + } + + if (infos != null) { + for (String key : infos.keySet()) { + Set values = this.runtimeInfos.get(key); + if (values == null) { + values = new HashSet<>(); + this.runtimeInfos.put(key, values); + } + values.add(infos.get(key)); + } + } + } + + public void clearRuntimeInfo() { + this.runtimeInfos = null; + } } From 9b3a3e28e732c836592c58eabbc52725d75f9b36 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Mon, 21 Nov 2016 20:18:59 +0530 Subject: [PATCH 02/22] Fix checkstyle --- .../org/apache/zeppelin/interpreter/InterpreterSetting.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index 07b98e36991..1f258c4f82c 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -405,11 +405,11 @@ public void setInterpreterRunner(InterpreterRunner interpreterRunner) { this.interpreterRunner = interpreterRunner; ======= public void addNoteToPara(String noteId, String paraId) { - if(noteIdToParaIdsetMap == null) { - noteIdToParaIdsetMap = new HashMap<>(); + if (noteIdToParaIdsetMap == null) { + noteIdToParaIdsetMap = new HashMap<>(); } Set paraIdSet = noteIdToParaIdsetMap.get(noteId); - if(paraIdSet == null) { + if (paraIdSet == null) { paraIdSet = new HashSet<>(); noteIdToParaIdsetMap.put(noteId, paraIdSet); } From 3d9a5733c08db4b394f2a9988510f4d5b8647186 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Mon, 21 Nov 2016 21:14:12 +0530 Subject: [PATCH 03/22] Fix NPE and some refactoring --- .../main/java/org/apache/zeppelin/socket/NotebookServer.java | 3 ++- .../org/apache/zeppelin/interpreter/InterpreterSetting.java | 4 ++++ 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 4b783b9d610..85b2a8b5396 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -2343,8 +2343,8 @@ public void clearParagraphRuntimeInfo(InterpreterSetting setting) { Note note = notebook().getNote(noteId); if (note != null) { Paragraph paragraph = note.getParagraph(paraId); - paragraph.clearRuntimeInfo(); if (paragraph != null) { + paragraph.clearRuntimeInfo(); broadcast(noteId, new Message(OP.PARAGRAPH).put("paragraph", paragraph)); } } @@ -2352,5 +2352,6 @@ public void clearParagraphRuntimeInfo(InterpreterSetting setting) { } } } + setting.clearNoteIdAndParaMap(); } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index 1f258c4f82c..d7f0c479bf9 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -420,4 +420,8 @@ public Map> getNoteIdAndParaMap() { return noteIdToParaIdsetMap; >>>>>>> Ability to view spark job urls in each paragraph } + + public void clearNoteIdAndParaMap() { + noteIdToParaIdsetMap = null; + } } From e2cd4db0211bc26d7bf8d3ea40f75ef13217e9e7 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Tue, 22 Nov 2016 10:42:37 +0530 Subject: [PATCH 04/22] Fix NPE in tests Signed-off-by: karuppayya --- .../main/java/org/apache/zeppelin/spark/SparkInterpreter.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java index 7d3caf8a85d..f4833f6f867 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java @@ -172,7 +172,9 @@ public synchronized void onJobStart(SparkListenerJobStart jobStart) { RemoteEventClientWrapper eventClient = ZeppelinContext.getEventClient(); Map infos = new java.util.HashMap<>(); infos.put("jobUrl", jobUrl); - eventClient.onParaInfosReceived(noteId, paragraphId, infos); + if (eventClient != null) { + eventClient.onParaInfosReceived(noteId, paragraphId, infos); + } } } From 7383c0a5c45b0eaa954433528944e2efa68875a9 Mon Sep 17 00:00:00 2001 From: Karup Date: Sat, 26 Nov 2016 08:32:47 +0530 Subject: [PATCH 05/22] Address review comments Signed-off-by: Karup --- .../zeppelin/spark/PySparkInterpreter.java | 2 +- .../zeppelin/spark/SparkInterpreter.java | 11 +++---- .../zeppelin/spark/SparkRInterpreter.java | 2 +- .../zeppelin/spark/SparkSqlInterpreter.java | 4 +-- .../java/org/apache/zeppelin/spark/Utils.java | 5 ++++ .../zeppelin/spark/ZeppelinContext.java | 3 +- .../zeppelin/socket/NotebookServer.java | 4 ++- .../notebook/paragraph/paragraph-control.html | 19 +++++++++++- .../interpreter/InterpreterSetting.java | 2 +- .../apache/zeppelin/notebook/Paragraph.java | 15 +++++----- .../notebook/ParagraphRuntimeInfos.java | 29 +++++++++++++++++++ 11 files changed, 73 insertions(+), 23 deletions(-) create mode 100644 zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java diff --git a/spark/src/main/java/org/apache/zeppelin/spark/PySparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/PySparkInterpreter.java index 5a8e0407158..0679fcc7b12 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/PySparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/PySparkInterpreter.java @@ -377,7 +377,7 @@ public InterpreterResult interpret(String st, InterpreterContext context) { "pyspark " + sparkInterpreter.getSparkContext().version() + " is not supported")); return new InterpreterResult(Code.ERROR, errorMessage); } - String jobGroup = sparkInterpreter.getJobGroup(context); + String jobGroup = Utils.buildJobGroupId(context); ZeppelinContext z = sparkInterpreter.getZeppelinContext(); z.setInterpreterContext(context); z.setGui(context.getGui()); diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java index f4833f6f867..4611ddf2c4a 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java @@ -172,6 +172,7 @@ public synchronized void onJobStart(SparkListenerJobStart jobStart) { RemoteEventClientWrapper eventClient = ZeppelinContext.getEventClient(); Map infos = new java.util.HashMap<>(); infos.put("jobUrl", jobUrl); + infos.put("label", "SPARK JOB"); if (eventClient != null) { eventClient.onParaInfosReceived(noteId, paragraphId, infos); } @@ -1147,10 +1148,6 @@ public Object getLastObject() { return obj; } - String getJobGroup(InterpreterContext context) { - return "zeppelin-" + context.getNoteId() + "-" + context.getParagraphId(); - } - /** * Interpret a single line. */ @@ -1171,7 +1168,7 @@ public InterpreterResult interpret(String line, InterpreterContext context) { public InterpreterResult interpret(String[] lines, InterpreterContext context) { synchronized (this) { z.setGui(context.getGui()); - sc.setJobGroup(getJobGroup(context), "Zeppelin", false); + sc.setJobGroup(Utils.buildJobGroupId(context), "Zeppelin", false); InterpreterResult r = interpretInput(lines, context); sc.clearJobGroup(); return r; @@ -1294,12 +1291,12 @@ private void putLatestVarInResourcePool(InterpreterContext context) { @Override public void cancel(InterpreterContext context) { - sc.cancelJobGroup(getJobGroup(context)); + sc.cancelJobGroup(Utils.buildJobGroupId(context)); } @Override public int getProgress(InterpreterContext context) { - String jobGroup = getJobGroup(context); + String jobGroup = Utils.buildJobGroupId(context); int completedTasks = 0; int totalTasks = 0; diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java index 53beea854a5..97c01360776 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java @@ -105,7 +105,7 @@ public InterpreterResult interpret(String lines, InterpreterContext interpreterC SparkInterpreter sparkInterpreter = getSparkInterpreter(); sparkInterpreter.populateSparkWebUrl(interpreterContext); - String jobGroup = sparkInterpreter.getJobGroup(interpreterContext); + String jobGroup = Utils.buildJobGroupId(interpreterContext); sparkInterpreter.getSparkContext().setJobGroup(jobGroup, "Zeppelin", false); String imageWidth = getProperty("zeppelin.R.image.width"); diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java index 81a02234381..1d5282f132e 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java @@ -101,7 +101,7 @@ public InterpreterResult interpret(String st, InterpreterContext context) { sc.setLocalProperty("spark.scheduler.pool", null); } - sc.setJobGroup(sparkInterpreter.getJobGroup(context), "Zeppelin", false); + sc.setJobGroup(Utils.buildJobGroupId(context), "Zeppelin", false); Object rdd = null; try { // method signature of sqlc.sql() is changed @@ -134,7 +134,7 @@ public void cancel(InterpreterContext context) { SQLContext sqlc = sparkInterpreter.getSQLContext(); SparkContext sc = sqlc.sparkContext(); - sc.cancelJobGroup(sparkInterpreter.getJobGroup(context)); + sc.cancelJobGroup(Utils.buildJobGroupId(context)); } @Override diff --git a/spark/src/main/java/org/apache/zeppelin/spark/Utils.java b/spark/src/main/java/org/apache/zeppelin/spark/Utils.java index 78304fd8713..1a7fa667b32 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/Utils.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/Utils.java @@ -17,6 +17,7 @@ package org.apache.zeppelin.spark; +import org.apache.zeppelin.interpreter.InterpreterContext; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -106,4 +107,8 @@ static boolean isSpark2() { return false; } } + + public static String buildJobGroupId(InterpreterContext context) { + return "zeppelin-" + context.getNoteId() + "-" + context.getParagraphId(); + } } diff --git a/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java b/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java index 620543d5829..d62b68e75ff 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/ZeppelinContext.java @@ -223,8 +223,7 @@ public static String showDF(SparkContext sc, Object df, int maxResult) { Object[] rows = null; Method take; - String jobGroup = "zeppelin-" + interpreterContext.getNoteId() + "-" - + interpreterContext.getParagraphId(); + String jobGroup = Utils.buildJobGroupId(interpreterContext); sc.setJobGroup(jobGroup, "Zeppelin", false); try { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 85b2a8b5396..8b5c93929d4 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -2327,7 +2327,9 @@ public void onParaInfosReceived(String noteId, String paragraphId, setting.addNoteToPara(noteId, paragraphId); metaInfos.remove("noteId"); metaInfos.remove("paraId"); - paragraph.updateRuntimeInfos(metaInfos); + String label = metaInfos.get("label"); + metaInfos.remove("label"); + paragraph.updateRuntimeInfos(label, metaInfos, setting.getGroup()); broadcast(note.getId(), new Message(OP.PARAGRAPH).put("paragraph", paragraph)); } } diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html index 91d094f3682..98381a09c63 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html @@ -13,7 +13,24 @@ -->
- + + + + {{paragraph.runtimeInfos.jobUrl.label}} + + + + + + {{paragraph.runtimeInfos.jobUrl.label}}S + + + {{paragraph.status}} diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index d7f0c479bf9..3e36d6e0ad0 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -127,7 +127,7 @@ public String getName() { return name; } - String getGroup() { + public String getGroup() { return group; } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java index c7c5fd5e58d..4a5637d6fe7 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java @@ -71,6 +71,7 @@ public class Paragraph extends Job implements Serializable, Cloneable { // For backward compatibility of note.json format after ZEPPELIN-212 Object result; private Map> runtimeInfos; + private Map runtimeInfos; /** * Applicaiton states in this paragraph @@ -678,19 +679,19 @@ private boolean isValidInterpreter(String replName) { } } - public void updateRuntimeInfos(Map infos) { + public void updateRuntimeInfos(String label, Map infos, String group) { if (this.runtimeInfos == null) { - this.runtimeInfos = new HashMap>(); + this.runtimeInfos = new HashMap(); } if (infos != null) { for (String key : infos.keySet()) { - Set values = this.runtimeInfos.get(key); - if (values == null) { - values = new HashSet<>(); - this.runtimeInfos.put(key, values); + ParagraphRuntimeInfos info = this.runtimeInfos.get(key); + if (info == null) { + info = new ParagraphRuntimeInfos(key, label, group); + this.runtimeInfos.put(key, info); } - values.add(infos.get(key)); + info.addValue(infos.get(key)); } } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java new file mode 100644 index 00000000000..b48cf2bef5b --- /dev/null +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java @@ -0,0 +1,29 @@ +package org.apache.zeppelin.notebook; + +import java.util.ArrayList; +import java.util.List; + +/** + * + * + */ +public class ParagraphRuntimeInfos { + + String propertyName; + String label; + String group; + List values; + + public ParagraphRuntimeInfos(String propertyName, String label, String group) { + this.propertyName = propertyName; + this.label = label; + this.group = group; + } + + public void addValue(String value) { + if (values == null) { + values = new ArrayList<>(); + } + values.add(value); + } +} From 25379aa2f60b7c946e1c37ce4cc9b29e5dd660d5 Mon Sep 17 00:00:00 2001 From: Karup Date: Sat, 26 Nov 2016 09:08:21 +0530 Subject: [PATCH 06/22] Clear job urls when we clear output Signed-off-by: Karup --- .../src/main/java/org/apache/zeppelin/notebook/Note.java | 1 + 1 file changed, 1 insertion(+) diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java index 26f4e1a9032..060dc7b5ad1 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java @@ -379,6 +379,7 @@ public Paragraph clearParagraphOutput(String paragraphId) { for (Paragraph p : paragraphs) { if (p.getId().equals(paragraphId)) { p.setReturn(null, null); + p.clearRuntimeInfo(); return p; } } From 717eedf030d6117021a62584a0adedbbd13931c4 Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 27 Nov 2016 14:30:35 +0530 Subject: [PATCH 07/22] Add tests , refactor Signed-off-by: Karup --- .../zeppelin/spark/SparkInterpreter.java | 5 +- .../zeppelin/spark/SparkInterpreterTest.java | 61 ++++++++++++++++++- 2 files changed, 64 insertions(+), 2 deletions(-) diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java index 4611ddf2c4a..36dc605c74f 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java @@ -992,7 +992,10 @@ public void open() { numReferenceOfSparkContext.incrementAndGet(); } - private String getSparkUIUrl() { + public String getSparkUIUrl() { + if (sparkUrl != null) { + return sparkUrl; + } Option sparkUiOption = (Option) Utils.invokeMethod(sc, "ui"); SparkUI sparkUi = sparkUiOption.get(); String sparkWebUrl = sparkUi.appUIAddress(); diff --git a/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java b/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java index 14108901ce7..4f457aa8aa8 100644 --- a/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java +++ b/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java @@ -23,11 +23,13 @@ import java.util.HashMap; import java.util.LinkedList; import java.util.List; +import java.util.Map; import java.util.Properties; import org.apache.spark.SparkConf; import org.apache.spark.SparkContext; import org.apache.zeppelin.display.AngularObjectRegistry; +import org.apache.zeppelin.interpreter.remote.RemoteEventClientWrapper; import org.apache.zeppelin.interpreter.thrift.InterpreterCompletion; import org.apache.zeppelin.resource.LocalResourcePool; import org.apache.zeppelin.resource.WellKnownResourceName; @@ -54,6 +56,8 @@ public class SparkInterpreterTest { public static InterpreterGroup intpGroup; private InterpreterContext context; public static Logger LOGGER = LoggerFactory.getLogger(SparkInterpreterTest.class); + private Map> paraIdToInfosMap = + new HashMap<>(); /** * Get spark version number as a numerical value. @@ -92,6 +96,20 @@ public void setUp() throws Exception { repl.open(); } + final RemoteEventClientWrapper remoteEventClientWrapper = new RemoteEventClientWrapper() { + + @Override + public void onParaInfosReceived(String noteId, String paragraphId, + Map infos) { + if (infos != null) { + paraIdToInfosMap.put(paragraphId, infos); + } + } + + @Override + public void onMetaInfosReceived(Map infos) { + } + }; context = new InterpreterContext("note", "id", null, "title", "text", new AuthenticationInfo(), new HashMap(), @@ -99,7 +117,15 @@ public void setUp() throws Exception { new AngularObjectRegistry(intpGroup.getId(), null), new LocalResourcePool("id"), new LinkedList(), - new InterpreterOutput(null)); + new InterpreterOutput(null)) { + + @Override + public RemoteEventClientWrapper getClient() { + return remoteEventClientWrapper; + } + }; + //first paragraph initializes sparkurl, run a dummy para + repl.interpret("sc", context); } @Test @@ -273,4 +299,37 @@ public void testCompletion() { List completions = repl.completion("sc.", "sc.".length()); assertTrue(completions.size() > 0); } + + @Test + public void testParagraphUrls() { + String paraId = "test_para_job_url"; + InterpreterContext intpCtx = new InterpreterContext("note", paraId, null, "title", "text", + new AuthenticationInfo(), + new HashMap(), + new GUI(), + new AngularObjectRegistry(intpGroup.getId(), null), + new LocalResourcePool("id"), + new LinkedList(), + new InterpreterOutput(new InterpreterOutputListener() { + @Override + public void onAppend(InterpreterOutput out, byte[] line) { + + } + + @Override + public void onUpdate(InterpreterOutput out, byte[] output) { + + } + })); + repl.interpret("sc.parallelize(1 to 10).map(x => {x}).collect", intpCtx); + Map paraInfos = paraIdToInfosMap.get(intpCtx.getParagraphId()); + String jobUrl = null; + if (paraInfos != null) { + jobUrl = paraInfos.get("jobUrl"); + } + String sparkUIUrl = repl.getSparkUIUrl(); + assertNotNull(jobUrl); + assertTrue(jobUrl.startsWith(sparkUIUrl + "/jobs/job?id=")); + + } } From 1a4528405c54751cf6c7f645a2af41d35652a5cf Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 27 Nov 2016 16:27:00 +0530 Subject: [PATCH 08/22] Fix test Signed-off-by: Karup --- .../org/apache/zeppelin/spark/SparkInterpreterTest.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java b/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java index 4f457aa8aa8..b9cf90c48f6 100644 --- a/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java +++ b/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java @@ -56,7 +56,7 @@ public class SparkInterpreterTest { public static InterpreterGroup intpGroup; private InterpreterContext context; public static Logger LOGGER = LoggerFactory.getLogger(SparkInterpreterTest.class); - private Map> paraIdToInfosMap = + private static Map> paraIdToInfosMap = new HashMap<>(); /** @@ -124,7 +124,11 @@ public RemoteEventClientWrapper getClient() { return remoteEventClientWrapper; } }; - //first paragraph initializes sparkurl, run a dummy para + // The first para interpretdr will set the Eventclient wrapper + //SparkInterpreter.interpret(String, InterpreterContext) -> + //SparkInterpreter.populateSparkWebUrl(InterpreterContext) -> + //ZeppelinContext.setEventClient(RemoteEventClientWrapper) + //running a dummy to ensure that we dont have any race conditions among tests repl.interpret("sc", context); } From d4e54e8cad60710f14e635f941f13f006a33754d Mon Sep 17 00:00:00 2001 From: karuppayya Date: Thu, 1 Dec 2016 21:34:00 +0530 Subject: [PATCH 09/22] Address review feedbacks Signed-off-by: karuppayya --- .../zeppelin/spark/SparkInterpreter.java | 17 ++------- .../java/org/apache/zeppelin/spark/Utils.java | 12 ++++++ .../remote/RemoteInterpreterEventPoller.java | 12 ++---- .../remote/RemoteInterpreterUtils.java | 10 +++++ .../zeppelin/socket/NotebookServer.java | 4 +- .../org/apache/zeppelin/notebook/Note.java | 3 +- .../apache/zeppelin/notebook/Notebook.java | 2 +- .../apache/zeppelin/notebook/Paragraph.java | 35 +++++++++++++++--- .../notebook/ParagraphRuntimeInfo.java | 37 +++++++++++++++++++ .../notebook/ParagraphRuntimeInfos.java | 29 --------------- 10 files changed, 100 insertions(+), 61 deletions(-) create mode 100644 zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java delete mode 100644 zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java index 36dc605c74f..2d18510993d 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java @@ -166,8 +166,8 @@ public synchronized void onJobStart(SparkListenerJobStart jobStart) { int jobId = jobStart.jobId(); String jobGroupId = jobStart.properties().getProperty("spark.jobGroup.id"); String jobUrl = getJobUrl(jobId); - String noteId = getNoteId(jobGroupId); - String paragraphId = getParagraphId(jobGroupId); + String noteId = Utils.getNoteId(jobGroupId); + String paragraphId = Utils.getParagraphId(jobGroupId); if (jobUrl != null && noteId != null && paragraphId != null) { RemoteEventClientWrapper eventClient = ZeppelinContext.getEventClient(); Map infos = new java.util.HashMap<>(); @@ -187,17 +187,6 @@ private String getJobUrl(int jobId) { return jobUrl; } - private String getNoteId(String jobgroupId) { - int indexOf = jobgroupId.indexOf("-"); - int secondIndex = jobgroupId.indexOf("-", indexOf + 1); - return jobgroupId.substring(indexOf + 1, secondIndex); - } - - private String getParagraphId(String jobgroupId) { - int indexOf = jobgroupId.indexOf("-"); - int secondIndex = jobgroupId.indexOf("-", indexOf + 1); - return jobgroupId.substring(secondIndex + 1, jobgroupId.length()); - } }; try { Object listenerBus = context.getClass().getMethod("listenerBus").invoke(context); @@ -1016,8 +1005,8 @@ public void populateSparkWebUrl(InterpreterContext ctx) { Map infos = new java.util.HashMap<>(); if (sparkUrl != null) { infos.put("url", sparkUrl); - logger.info("Sending metainfos to Zeppelin server: {}", infos.toString()); if (ctx != null && ctx.getClient() != null) { + logger.info("Sending metainfos to Zeppelin server: {}", infos.toString()); getZeppelinContext().setEventClient(ctx.getClient()); ctx.getClient().onMetaInfosReceived(infos); } diff --git a/spark/src/main/java/org/apache/zeppelin/spark/Utils.java b/spark/src/main/java/org/apache/zeppelin/spark/Utils.java index 1a7fa667b32..17edb0d3257 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/Utils.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/Utils.java @@ -111,4 +111,16 @@ static boolean isSpark2() { public static String buildJobGroupId(InterpreterContext context) { return "zeppelin-" + context.getNoteId() + "-" + context.getParagraphId(); } + + public static String getNoteId(String jobgroupId) { + int indexOf = jobgroupId.indexOf("-"); + int secondIndex = jobgroupId.indexOf("-", indexOf + 1); + return jobgroupId.substring(indexOf + 1, secondIndex); + } + + public static String getParagraphId(String jobgroupId) { + int indexOf = jobgroupId.indexOf("-"); + int secondIndex = jobgroupId.indexOf("-", indexOf + 1); + return jobgroupId.substring(secondIndex + 1, jobgroupId.length()); + } } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java index 0ecb2ab9e28..f46d31af6c2 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventPoller.java @@ -243,7 +243,8 @@ public void run() { Map metaInfos = gson.fromJson(event.getData(), new TypeToken>() { }.getType()); - String settingId = getInterpreterSettingId(); + String settingId = RemoteInterpreterUtils. + getInterpreterSettingId(interpreterGroup.getId()); listener.onMetaInfosReceived(settingId, metaInfos); } else if (event.getType() == RemoteInterpreterEventType.PARA_INFOS) { Map paraInfos = gson.fromJson(event.getData(), @@ -251,7 +252,8 @@ public void run() { }.getType()); String noteId = paraInfos.get("noteId"); String paraId = paraInfos.get("paraId"); - String settingId = getInterpreterSettingId(); + String settingId = RemoteInterpreterUtils. + getInterpreterSettingId(interpreterGroup.getId()); if (noteId != null && paraId != null && settingId != null) { listener.onParaInfosReceived(noteId, paraId, settingId, paraInfos); } @@ -342,12 +344,6 @@ public void onError() { } } - private String getInterpreterSettingId() { - String id = interpreterGroup.getId(); - int indexOfColon = id.indexOf(":"); - return id.substring(0, indexOfColon); - } - private void sendResourcePoolResponseGetAll(ResourceSet resourceSet) { Client client = null; boolean broken = false; diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterUtils.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterUtils.java index 2937e2d4c09..8308222d9a8 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterUtils.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterUtils.java @@ -17,6 +17,7 @@ package org.apache.zeppelin.interpreter.remote; +import org.apache.zeppelin.interpreter.InterpreterGroup; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -63,4 +64,13 @@ public static boolean checkIfRemoteEndpointAccessible(String host, int port) { return false; } } + + public static String getInterpreterSettingId(String intpGrpId) { + String settingId = null; + if (intpGrpId != null) { + int indexOfColon = intpGrpId.indexOf(":"); + settingId = intpGrpId.substring(0, indexOfColon); + } + return settingId; + } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 8b5c93929d4..e6ea91bb73a 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -2329,7 +2329,7 @@ public void onParaInfosReceived(String noteId, String paragraphId, metaInfos.remove("paraId"); String label = metaInfos.get("label"); metaInfos.remove("label"); - paragraph.updateRuntimeInfos(label, metaInfos, setting.getGroup()); + paragraph.updateRuntimeInfos(label, metaInfos, setting.getGroup(), setting.getId()); broadcast(note.getId(), new Message(OP.PARAGRAPH).put("paragraph", paragraph)); } } @@ -2346,7 +2346,7 @@ public void clearParagraphRuntimeInfo(InterpreterSetting setting) { if (note != null) { Paragraph paragraph = note.getParagraph(paraId); if (paragraph != null) { - paragraph.clearRuntimeInfo(); + paragraph.clearRuntimeInfo(setting.getId()); broadcast(noteId, new Message(OP.PARAGRAPH).put("paragraph", paragraph)); } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java index 060dc7b5ad1..224dd4b7697 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Note.java @@ -379,7 +379,7 @@ public Paragraph clearParagraphOutput(String paragraphId) { for (Paragraph p : paragraphs) { if (p.getId().equals(paragraphId)) { p.setReturn(null, null); - p.clearRuntimeInfo(); + p.clearRuntimeInfo(null); return p; } } @@ -564,6 +564,7 @@ public void run(String paragraphId) { return; } + p.clearRuntimeInfo(null); String requiredReplName = p.getRequiredReplName(); Interpreter intp = factory.getInterpreter(p.getUser(), getId(), requiredReplName); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java index f94e5dd571d..b853e070431 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Notebook.java @@ -480,7 +480,7 @@ public Note loadNoteFromRepo(String id, AuthenticationInfo subject) { if (p.getDateFinished() != null && lastUpdatedDate.before(p.getDateFinished())) { lastUpdatedDate = p.getDateFinished(); } - p.clearRuntimeInfo(); + p.clearRuntimeInfo(null); } Map> savedObjects = note.getAngularObjects(); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java index 4a5637d6fe7..5423a3e9141 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java @@ -679,16 +679,17 @@ private boolean isValidInterpreter(String replName) { } } - public void updateRuntimeInfos(String label, Map infos, String group) { + public void updateRuntimeInfos(String label, Map infos, String group, + String intpSettingId) { if (this.runtimeInfos == null) { - this.runtimeInfos = new HashMap(); + this.runtimeInfos = new HashMap(); } if (infos != null) { for (String key : infos.keySet()) { - ParagraphRuntimeInfos info = this.runtimeInfos.get(key); + ParagraphRuntimeInfo info = this.runtimeInfos.get(key); if (info == null) { - info = new ParagraphRuntimeInfos(key, label, group); + info = new ParagraphRuntimeInfo(key, label, group, intpSettingId); this.runtimeInfos.put(key, info); } info.addValue(infos.get(key)); @@ -696,7 +697,29 @@ public void updateRuntimeInfos(String label, Map infos, String g } } - public void clearRuntimeInfo() { - this.runtimeInfos = null; + /** + * Remove runtimeinfo taht were got from the setting with id settingId + * @param settingId + */ + public void clearRuntimeInfo(String settingId) { + if (settingId != null) { + Set keys = runtimeInfos.keySet(); + if (keys.size() > 0) { + List infosToRemove = new ArrayList<>(); + for (String key : keys) { + ParagraphRuntimeInfo paragraphRuntimeInfo = runtimeInfos.get(key); + if (paragraphRuntimeInfo.getInterpreterSettingId().equals(settingId)) { + infosToRemove.add(key); + } + } + if (infosToRemove.size() > 0) { + for (String info : infosToRemove) { + runtimeInfos.remove(info); + } + } + } + } else { + this.runtimeInfos = null; + } } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java new file mode 100644 index 00000000000..22fbc8924b1 --- /dev/null +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java @@ -0,0 +1,37 @@ +package org.apache.zeppelin.notebook; + +import java.util.ArrayList; +import java.util.List; + +/** + * Store runtime infos of each para + * + */ +public class ParagraphRuntimeInfo { + + private String propertyName; //Name of the property + private String label; //Label to be used in UI + private String group; //The interpretergroup from which the info was derived + private List values; // values for the property + private String interpreterSettingId; + + public ParagraphRuntimeInfo(String propertyName, String label, + String group, String intpSettingId) { + if (interpreterSettingId == null) { + throw new IllegalArgumentException("Interpreter setting Id cannot be null"); + } + this.propertyName = propertyName; + this.label = label; + this.group = group; + this.interpreterSettingId = intpSettingId; + this.values = new ArrayList<>(); + } + + public void addValue(String value) { + values.add(value); + } + + public String getInterpreterSettingId() { + return interpreterSettingId; + } +} diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java deleted file mode 100644 index b48cf2bef5b..00000000000 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfos.java +++ /dev/null @@ -1,29 +0,0 @@ -package org.apache.zeppelin.notebook; - -import java.util.ArrayList; -import java.util.List; - -/** - * - * - */ -public class ParagraphRuntimeInfos { - - String propertyName; - String label; - String group; - List values; - - public ParagraphRuntimeInfos(String propertyName, String label, String group) { - this.propertyName = propertyName; - this.label = label; - this.group = group; - } - - public void addValue(String value) { - if (values == null) { - values = new ArrayList<>(); - } - values.add(value); - } -} From 42d92ac9649ac39f27c8e6edbcf6859461897b69 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Fri, 2 Dec 2016 09:36:23 +0530 Subject: [PATCH 10/22] Fix test Signed-off-by: karuppayya --- .../apache/zeppelin/spark/SparkInterpreterTest.java | 12 +----------- 1 file changed, 1 insertion(+), 11 deletions(-) diff --git a/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java b/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java index b9cf90c48f6..8c78b663b2e 100644 --- a/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java +++ b/spark/src/test/java/org/apache/zeppelin/spark/SparkInterpreterTest.java @@ -314,17 +314,7 @@ public void testParagraphUrls() { new AngularObjectRegistry(intpGroup.getId(), null), new LocalResourcePool("id"), new LinkedList(), - new InterpreterOutput(new InterpreterOutputListener() { - @Override - public void onAppend(InterpreterOutput out, byte[] line) { - - } - - @Override - public void onUpdate(InterpreterOutput out, byte[] output) { - - } - })); + new InterpreterOutput(null)); repl.interpret("sc.parallelize(1 to 10).map(x => {x}).collect", intpCtx); Map paraInfos = paraIdToInfosMap.get(intpCtx.getParagraphId()); String jobUrl = null; From b837c6cb0be209c68249c95795510edc1e1c8be3 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Thu, 8 Dec 2016 20:25:19 +0530 Subject: [PATCH 11/22] Fix incorrect variable used Signed-off-by: karuppayya --- .../java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java index 22fbc8924b1..964afa9d254 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java @@ -17,7 +17,7 @@ public class ParagraphRuntimeInfo { public ParagraphRuntimeInfo(String propertyName, String label, String group, String intpSettingId) { - if (interpreterSettingId == null) { + if (intpSettingId == null) { throw new IllegalArgumentException("Interpreter setting Id cannot be null"); } this.propertyName = propertyName; From fc44d9b479f031f77a30c7d92db56c56307abf11 Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 18 Dec 2016 09:30:37 +0530 Subject: [PATCH 12/22] Address review comments Signed-off-by: Karup --- .../zeppelin/spark/SparkInterpreter.java | 1 + .../zeppelin/socket/NotebookServer.java | 17 +++++--- .../notebook/paragraph/paragraph-control.html | 6 +-- .../paragraph/paragraph.controller.js | 41 ++++--------------- .../websocketEvents.factory.js | 2 + .../interpreter/InterpreterSetting.java | 20 +++++---- .../apache/zeppelin/notebook/Paragraph.java | 10 +++-- .../notebook/ParagraphRuntimeInfo.java | 10 +++-- .../zeppelin/notebook/socket/Message.java | 3 +- 9 files changed, 51 insertions(+), 59 deletions(-) diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java index 2d18510993d..3c1288e1afc 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java @@ -173,6 +173,7 @@ public synchronized void onJobStart(SparkListenerJobStart jobStart) { Map infos = new java.util.HashMap<>(); infos.put("jobUrl", jobUrl); infos.put("label", "SPARK JOB"); + infos.put("tooltip", "View in Spark web UI"); if (eventClient != null) { eventClient.onParaInfosReceived(noteId, paragraphId, infos); } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index e6ea91bb73a..11702a4e47a 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -20,6 +20,7 @@ import java.net.URISyntaxException; import java.net.UnknownHostException; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashMap; import java.util.HashSet; import java.util.LinkedList; @@ -2325,12 +2326,18 @@ public void onParaInfosReceived(String noteId, String paragraphId, InterpreterSetting setting = notebook().getInterpreterFactory() .get(interpreterSettingId); setting.addNoteToPara(noteId, paragraphId); - metaInfos.remove("noteId"); - metaInfos.remove("paraId"); String label = metaInfos.get("label"); - metaInfos.remove("label"); - paragraph.updateRuntimeInfos(label, metaInfos, setting.getGroup(), setting.getId()); - broadcast(note.getId(), new Message(OP.PARAGRAPH).put("paragraph", paragraph)); + String tooltip = metaInfos.get("tooltip"); + List keysToRemove = Arrays.asList("noteId", "paraId", "label", "tooltip"); + for (String removeKey : keysToRemove) { + metaInfos.remove(removeKey); + } + paragraph + .updateRuntimeInfos(label, tooltip, metaInfos, setting.getGroup(), setting.getId()); + broadcast( + note.getId(), + new Message(OP.PARAS_INFO).put("id", paragraphId).put("infos", + paragraph.getRuntimeInfos())); } } } diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html index 98381a09c63..e1bff080b4f 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html @@ -13,14 +13,14 @@ -->
- + {{paragraph.runtimeInfos.jobUrl.label}} - - + {{paragraph.runtimeInfos.jobUrl.label}}S diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js index 5742ab1e03b..94d30f5ac96 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js @@ -1036,6 +1036,12 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat } }); + $scope.$on('updateParaInfos', function(event, data) { + if (data.id === $scope.paragraph.id) { + $scope.paragraph.runtimeInfos = data.infos; + } + }); + $scope.$on('angularObjectRemove', function(event, data) { var noteId = $route.current.pathParams.noteId; if (!data.noteId || data.noteId === noteId) { @@ -1084,7 +1090,7 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat newPara.errorMessage !== oldPara.errorMessage || !angular.equals(newPara.settings, oldPara.settings) || !angular.equals(newPara.config, oldPara.config) || - !angular.equals(newPara.runtimeInfos, oldPara.runtimeInfos))) + !angular.equals(newPara.runtimeInfos, oldPara.runtimeInfos))) } $scope.updateAllScopeTexts = function(oldPara, newPara) { @@ -1132,39 +1138,6 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat $scope.editor.setReadOnly($scope.isRunning(newPara)); } - if (!$scope.asIframe) { - $scope.paragraph.config = newPara.config; - initializeDefault(newPara.config); - } else { - newPara.config.editorHide = true; - newPara.config.tableHide = false; - $scope.paragraph.config = newPara.config; - } - }; - $scope.updateParagraph = function(oldPara, newPara, updateCallback) { - // 1. get status, refreshed - const statusChanged = (newPara.status !== oldPara.status); - const resultRefreshed = (newPara.dateFinished !== oldPara.dateFinished) || - isEmpty(newPara.results) !== isEmpty(oldPara.results) || - newPara.status === 'ERROR' || (newPara.status === 'FINISHED' && statusChanged); - - // 2. update texts managed by $scope - $scope.updateAllScopeTexts(oldPara, newPara); - - // 3. execute callback to update result - updateCallback(); - - // 4. update remaining paragraph objects - $scope.updateParagraphObjectWhenUpdated(newPara); - - // 5. handle scroll down by key properly if new paragraph is added - if (statusChanged || resultRefreshed) { - // when last paragraph runs, zeppelin automatically appends new paragraph. - // this broadcast will focus to the newly inserted paragraph - const paragraphs = angular.element('div[id$="_paragraphColumn_main"]'); - if (paragraphs.length >= 2 && paragraphs[paragraphs.length - 2].id.indexOf($scope.paragraph.id) === 0) { - // rendering output can took some time. So delay scrolling event firing for sometime. - setTimeout(() => { $rootScope.$broadcast('scrollToCursor'); }, 500); } } }; diff --git a/zeppelin-web/src/components/websocketEvents/websocketEvents.factory.js b/zeppelin-web/src/components/websocketEvents/websocketEvents.factory.js index aceffbbe30f..df0ea4806fb 100644 --- a/zeppelin-web/src/components/websocketEvents/websocketEvents.factory.js +++ b/zeppelin-web/src/components/websocketEvents/websocketEvents.factory.js @@ -164,6 +164,8 @@ function websocketEvents($rootScope, $websocket, $location, baseUrlSrv) { $rootScope.$broadcast('updateNote', data.name, data.config, data.info); } else if (op === 'SET_NOTE_REVISION') { $rootScope.$broadcast('setNoteRevisionResult', data); + } else if (op === 'PARAS_INFO') { + $rootScope.$broadcast('updateParaInfos', data); } }); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index 3e36d6e0ad0..e94cc9a6867 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -48,7 +48,10 @@ public class InterpreterSetting { // always be null in case of InterpreterSettingRef private String group; private transient Map infos; - private transient Map> noteIdToParaIdsetMap; + + // Map of the note and paragraphs which has runtime infos generated by this interpreter setting. + // This map is used to clear the infos in paragraph when the interpretersetting is restarted + private transient Map> runtimeInfosToBeCleared; /** * properties can be either Properties or Map @@ -396,32 +399,31 @@ public Map getInfos() { return infos; } -<<<<<<< b1e06919fada4240f5440430f2c27c02e4e5626a public InterpreterRunner getInterpreterRunner() { return interpreterRunner; } public void setInterpreterRunner(InterpreterRunner interpreterRunner) { this.interpreterRunner = interpreterRunner; -======= + } + public void addNoteToPara(String noteId, String paraId) { - if (noteIdToParaIdsetMap == null) { - noteIdToParaIdsetMap = new HashMap<>(); + if (runtimeInfosToBeCleared == null) { + runtimeInfosToBeCleared = new HashMap<>(); } - Set paraIdSet = noteIdToParaIdsetMap.get(noteId); + Set paraIdSet = runtimeInfosToBeCleared.get(noteId); if (paraIdSet == null) { paraIdSet = new HashSet<>(); - noteIdToParaIdsetMap.put(noteId, paraIdSet); + runtimeInfosToBeCleared.put(noteId, paraIdSet); } paraIdSet.add(paraId); } public Map> getNoteIdAndParaMap() { return noteIdToParaIdsetMap; ->>>>>>> Ability to view spark job urls in each paragraph } public void clearNoteIdAndParaMap() { - noteIdToParaIdsetMap = null; + runtimeInfosToBeCleared = null; } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java index 5423a3e9141..11754be0c8a 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java @@ -679,8 +679,8 @@ private boolean isValidInterpreter(String replName) { } } - public void updateRuntimeInfos(String label, Map infos, String group, - String intpSettingId) { + public void updateRuntimeInfos(String label, String tooltip, Map infos, + String group, String intpSettingId) { if (this.runtimeInfos == null) { this.runtimeInfos = new HashMap(); } @@ -689,7 +689,7 @@ public void updateRuntimeInfos(String label, Map infos, String g for (String key : infos.keySet()) { ParagraphRuntimeInfo info = this.runtimeInfos.get(key); if (info == null) { - info = new ParagraphRuntimeInfo(key, label, group, intpSettingId); + info = new ParagraphRuntimeInfo(key, label, tooltip, group, intpSettingId); this.runtimeInfos.put(key, info); } info.addValue(infos.get(key)); @@ -722,4 +722,8 @@ public void clearRuntimeInfo(String settingId) { this.runtimeInfos = null; } } + + public Map getRuntimeInfos() { + return runtimeInfos; + } } diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java index 964afa9d254..0042023d0b1 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/ParagraphRuntimeInfo.java @@ -9,19 +9,21 @@ */ public class ParagraphRuntimeInfo { - private String propertyName; //Name of the property - private String label; //Label to be used in UI - private String group; //The interpretergroup from which the info was derived + private String propertyName; // Name of the property + private String label; // Label to be used in UI + private String tooltip; // Tooltip text toshow in UI + private String group; // The interpretergroup from which the info was derived private List values; // values for the property private String interpreterSettingId; public ParagraphRuntimeInfo(String propertyName, String label, - String group, String intpSettingId) { + String tooltip, String group, String intpSettingId) { if (intpSettingId == null) { throw new IllegalArgumentException("Interpreter setting Id cannot be null"); } this.propertyName = propertyName; this.label = label; + this.tooltip = tooltip; this.group = group; this.interpreterSettingId = intpSettingId; this.values = new ArrayList<>(); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/socket/Message.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/socket/Message.java index a6d15462152..d40afc22ffb 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/socket/Message.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/socket/Message.java @@ -174,7 +174,8 @@ public static enum OP { NOTE_UPDATED, // [s-c] paragraph updated(name, config) RUN_ALL_PARAGRAPHS, // [c-s] run all paragraphs PARAGRAPH_EXECUTED_BY_SPELL, // [c-s] paragraph was executed by spell - RUN_PARAGRAPH_USING_SPELL // [s-c] run paragraph using spell + RUN_PARAGRAPH_USING_SPELL, // [s-c] run paragraph using spell + PARAS_INFO // [s-c] paragraph runtime infos } public static final Message EMPTY = new Message(null); From 09fc0e223ca33b25c07ae96cf6fc800d3cfba0e9 Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 18 Dec 2016 10:03:50 +0530 Subject: [PATCH 13/22] Fix compilation Signed-off-by: Karup --- .../main/java/org/apache/zeppelin/spark/SparkRInterpreter.java | 1 - 1 file changed, 1 deletion(-) diff --git a/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java b/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java index 97c01360776..16b1a21144f 100644 --- a/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java +++ b/spark/src/main/java/org/apache/zeppelin/spark/SparkRInterpreter.java @@ -127,7 +127,6 @@ public InterpreterResult interpret(String lines, InterpreterContext interpreterC } } - String jobGroup = getJobGroup(interpreterContext); String setJobGroup = ""; // assign setJobGroup to dummy__, otherwise it would print NULL for this statement if (Utils.isSpark2()) { From 19513a66048a7040043611f26d9e3cbf9da08d63 Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 1 Jan 2017 22:15:54 +0530 Subject: [PATCH 14/22] Send para runtimeinfos via websocker, but dont persist in json --- .../json/NotebookTypeAdapterFactory.java | 68 +++++++++++++++++++ .../zeppelin/socket/NotebookServer.java | 21 +++++- .../apache/zeppelin/notebook/Paragraph.java | 2 +- 3 files changed, 89 insertions(+), 2 deletions(-) create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java b/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java new file mode 100644 index 00000000000..922a5c1b127 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java @@ -0,0 +1,68 @@ +package org.apache.zeppelin.json; + +import java.io.IOException; + +import org.apache.zeppelin.socket.NotebookServer; + +import com.google.gson.Gson; +import com.google.gson.JsonElement; +import com.google.gson.TypeAdapter; +import com.google.gson.TypeAdapterFactory; +import com.google.gson.reflect.TypeToken; +import com.google.gson.stream.JsonReader; +import com.google.gson.stream.JsonWriter; + +/** + * Custom adapter type factory + * Modify the jsonObject before serailaization/deserialization + * Check sample implementation at {@link NotebookServer} + * @param the type whose json is to be customized for serialization/deserialization + */ +public class NotebookTypeAdapterFactory implements TypeAdapterFactory { + private final Class customizedClass; + + public NotebookTypeAdapterFactory(Class customizedClass) { + this.customizedClass = customizedClass; + } + + @SuppressWarnings("unchecked") + // we use a runtime check to guarantee that 'C' and 'T' are equal + public final TypeAdapter create(Gson gson, TypeToken type) { + return type.getRawType() == customizedClass ? (TypeAdapter) customizeTypeAdapter(gson, + (TypeToken) type) : null; + } + + private TypeAdapter customizeTypeAdapter(Gson gson, TypeToken type) { + final TypeAdapter delegate = gson.getDelegateAdapter(this, type); + final TypeAdapter elementAdapter = gson.getAdapter(JsonElement.class); + return new TypeAdapter() { + @Override + public void write(JsonWriter out, C value) throws IOException { + JsonElement tree = delegate.toJsonTree(value); + beforeWrite(value, tree); + elementAdapter.write(out, tree); + } + + @Override + public C read(JsonReader in) throws IOException { + JsonElement tree = elementAdapter.read(in); + afterRead(tree); + return delegate.fromJsonTree(tree); + } + }; + } + + /** + * Override this to change {@code toSerialize} before it is written to the + * outgoing JSON stream. + */ + protected void beforeWrite(C source, JsonElement toSerialize) { + } + + /** + * Override this to change {@code deserialized} before it parsed into the + * application type. + */ + protected void afterRead(JsonElement deserialized) { + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 11702a4e47a..e692b12fadf 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -55,6 +55,7 @@ import org.apache.zeppelin.interpreter.remote.RemoteAngularObjectRegistry; import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcessListener; import org.apache.zeppelin.interpreter.thrift.InterpreterCompletion; +import org.apache.zeppelin.json.NotebookTypeAdapterFactory; import org.apache.zeppelin.notebook.JobListenerFactory; import org.apache.zeppelin.notebook.Folder; import org.apache.zeppelin.notebook.Note; @@ -63,6 +64,7 @@ import org.apache.zeppelin.notebook.NotebookEventListener; import org.apache.zeppelin.notebook.Paragraph; import org.apache.zeppelin.notebook.ParagraphJobListener; +import org.apache.zeppelin.notebook.ParagraphRuntimeInfo; import org.apache.zeppelin.notebook.repo.NotebookRepo.Revision; import org.apache.zeppelin.notebook.socket.Message; import org.apache.zeppelin.notebook.socket.Message.OP; @@ -87,6 +89,10 @@ import org.slf4j.LoggerFactory; import com.google.common.collect.Queues; +import com.google.gson.Gson; +import com.google.gson.GsonBuilder; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; import com.google.gson.reflect.TypeToken; /** @@ -114,7 +120,20 @@ String getKey() { private static final Logger LOG = LoggerFactory.getLogger(NotebookServer.class); - Gson gson = new GsonBuilder().setDateFormat("yyyy-MM-dd'T'HH:mm:ssZ").create(); + Gson gson = new GsonBuilder() + .registerTypeAdapterFactory(new NotebookTypeAdapterFactory(Paragraph.class) { + @Override + protected void beforeWrite(Paragraph source, JsonElement toSerialize) { + Map runtimeInfos = source.getRuntimeInfos(); + if (runtimeInfos != null) { + JsonElement jsonTree = gson.toJsonTree(runtimeInfos); + if (toSerialize instanceof JsonObject) { + JsonObject jsonObj = (JsonObject) toSerialize; + jsonObj.add("runtimeInfos", jsonTree); + } + } + } + }).setDateFormat("yyyy-MM-dd'T'HH:mm:ssZ").create(); final Map> noteSocketMap = new HashMap<>(); final Queue connectedSockets = new ConcurrentLinkedQueue<>(); final Map> userConnectedSockets = new ConcurrentHashMap<>(); diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java index 11754be0c8a..15b7836815e 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java @@ -74,7 +74,7 @@ public class Paragraph extends Job implements Serializable, Cloneable { private Map runtimeInfos; /** - * Applicaiton states in this paragraph + * Application states in this paragraph */ private final List apps = new LinkedList<>(); From 87214a7fefe9c182267e465d14c5dfe8bc465a67 Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 1 Jan 2017 22:54:57 +0530 Subject: [PATCH 15/22] Fix incorrect rebase --- .../src/main/java/org/apache/zeppelin/notebook/Paragraph.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java index 15b7836815e..12501946730 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/notebook/Paragraph.java @@ -70,8 +70,7 @@ public class Paragraph extends Job implements Serializable, Cloneable { // For backward compatibility of note.json format after ZEPPELIN-212 Object result; - private Map> runtimeInfos; - private Map runtimeInfos; + private Map runtimeInfos; /** * Application states in this paragraph From d27221df3181338736f44cb59886bdbf6892385a Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 1 Jan 2017 23:33:53 +0530 Subject: [PATCH 16/22] Adding license header --- .../json/NotebookTypeAdapterFactory.java | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java b/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java index 922a5c1b127..a22c03b1518 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/json/NotebookTypeAdapterFactory.java @@ -1,3 +1,20 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + package org.apache.zeppelin.json; import java.io.IOException; From ed4685cf3ff0531e8616846f3758a54a2c1c7dc3 Mon Sep 17 00:00:00 2001 From: Karup Date: Fri, 27 Jan 2017 21:03:36 +0530 Subject: [PATCH 17/22] Fix tooltip --- zeppelin-web/src/app/notebook/paragraph/paragraph-control.html | 3 ++- .../org/apache/zeppelin/interpreter/InterpreterSetting.java | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html index e1bff080b4f..9a5916d1eec 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html @@ -14,7 +14,8 @@
- + {{paragraph.runtimeInfos.jobUrl.label}} diff --git a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index e94cc9a6867..74424303daf 100644 --- a/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-zengine/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -420,7 +420,7 @@ public void addNoteToPara(String noteId, String paraId) { } public Map> getNoteIdAndParaMap() { - return noteIdToParaIdsetMap; + return runtimeInfosToBeCleared; } public void clearNoteIdAndParaMap() { From 890107dc5c6c2331f4be4046fc4272b64ae431e9 Mon Sep 17 00:00:00 2001 From: Karup Date: Sun, 29 Jan 2017 14:57:51 +0530 Subject: [PATCH 18/22] Fix test - tryout --- .../apache/zeppelin/AbstractZeppelinIT.java | 2 +- .../notebook/paragraph/paragraph-control.html | 36 ++++++++++--------- 2 files changed, 20 insertions(+), 18 deletions(-) diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java b/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java index f8498685580..5575f04952a 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java @@ -64,7 +64,7 @@ protected String getParagraphXPath(int paragraphNo) { protected boolean waitForParagraph(final int paragraphNo, final String state) { By locator = By.xpath(getParagraphXPath(paragraphNo) - + "//div[contains(@class, 'control')]//span[1][contains(.,'" + state + "')]"); + + "//div[contains(@class, 'control')]//span[2][contains(.,'" + state + "')]"); WebElement element = pollingWait(locator, MAX_PARAGRAPH_TIMEOUT_SEC); return element.isDisplayed(); } diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html index 9a5916d1eec..5f9c4620ae0 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph-control.html @@ -13,24 +13,26 @@ -->
- - - - {{paragraph.runtimeInfos.jobUrl.label}} - - - - - - {{paragraph.runtimeInfos.jobUrl.label}}S + + + + + {{paragraph.runtimeInfos.jobUrl.label}} + + + + + + {{paragraph.runtimeInfos.jobUrl.label}}S + + - {{paragraph.status}} From 732b0a41cb9211757cd8fde8a5107260c72e8ce8 Mon Sep 17 00:00:00 2001 From: karuppayya Date: Mon, 30 Jan 2017 22:14:00 +0530 Subject: [PATCH 19/22] Fix test --- .../src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java b/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java index 5575f04952a..b0da7647042 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java @@ -71,7 +71,7 @@ protected boolean waitForParagraph(final int paragraphNo, final String state) { protected String getParagraphStatus(final int paragraphNo) { By locator = By.xpath(getParagraphXPath(paragraphNo) - + "//div[contains(@class, 'control')]//span[1]"); + + "//div[contains(@class, 'control')]//span[2]"); return driver.findElement(locator).getText(); } From 8e2cd853792bfee7426b2156e65e261652dfb8d9 Mon Sep 17 00:00:00 2001 From: Karup Date: Thu, 2 Feb 2017 12:39:28 +0530 Subject: [PATCH 20/22] tryout: fix selenium tests based on moons suggstion --- .../src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java b/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java index b0da7647042..41f4cec4841 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/AbstractZeppelinIT.java @@ -71,7 +71,7 @@ protected boolean waitForParagraph(final int paragraphNo, final String state) { protected String getParagraphStatus(final int paragraphNo) { By locator = By.xpath(getParagraphXPath(paragraphNo) - + "//div[contains(@class, 'control')]//span[2]"); + + "//div[contains(@class, 'control')]/span[2]"); return driver.findElement(locator).getText(); } From d7eb3b6bfa011a517ecc6d69c65f98700a38c795 Mon Sep 17 00:00:00 2001 From: Karup Date: Fri, 3 Feb 2017 17:58:11 +0530 Subject: [PATCH 21/22] Fix paragraph.js --- .../src/app/notebook/paragraph/paragraph.controller.js | 3 --- 1 file changed, 3 deletions(-) diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js index 94d30f5ac96..b7c3653f77d 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js @@ -1137,9 +1137,6 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat if ($scope.editor) { $scope.editor.setReadOnly($scope.isRunning(newPara)); } - - } - } }; $scope.$on('runParagraphUsingSpell', function(event, data) { From 4253d0bfa72217055889b975bd69b9de62f27bca Mon Sep 17 00:00:00 2001 From: Karup Date: Fri, 3 Feb 2017 18:05:01 +0530 Subject: [PATCH 22/22] Fix bad rebase --- .../paragraph/paragraph.controller.js | 37 +++++++++++++++++++ 1 file changed, 37 insertions(+) diff --git a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js index b7c3653f77d..e6f3244eee4 100644 --- a/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js +++ b/zeppelin-web/src/app/notebook/paragraph/paragraph.controller.js @@ -1137,6 +1137,43 @@ function ParagraphCtrl($scope, $rootScope, $route, $window, $routeParams, $locat if ($scope.editor) { $scope.editor.setReadOnly($scope.isRunning(newPara)); } + + if (!$scope.asIframe) { + $scope.paragraph.config = newPara.config; + initializeDefault(newPara.config); + } else { + newPara.config.editorHide = true; + newPara.config.tableHide = false; + $scope.paragraph.config = newPara.config; + } + }; + + $scope.updateParagraph = function(oldPara, newPara, updateCallback) { + // 1. get status, refreshed + const statusChanged = (newPara.status !== oldPara.status); + const resultRefreshed = (newPara.dateFinished !== oldPara.dateFinished) || + isEmpty(newPara.results) !== isEmpty(oldPara.results) || + newPara.status === 'ERROR' || (newPara.status === 'FINISHED' && statusChanged); + + // 2. update texts managed by $scope + $scope.updateAllScopeTexts(oldPara, newPara); + + // 3. execute callback to update result + updateCallback(); + + // 4. update remaining paragraph objects + $scope.updateParagraphObjectWhenUpdated(newPara); + + // 5. handle scroll down by key properly if new paragraph is added + if (statusChanged || resultRefreshed) { + // when last paragraph runs, zeppelin automatically appends new paragraph. + // this broadcast will focus to the newly inserted paragraph + const paragraphs = angular.element('div[id$="_paragraphColumn_main"]'); + if (paragraphs.length >= 2 && paragraphs[paragraphs.length - 2].id.indexOf($scope.paragraph.id) === 0) { + // rendering output can took some time. So delay scrolling event firing for sometime. + setTimeout(() => { $rootScope.$broadcast('scrollToCursor'); }, 500); + } + } }; $scope.$on('runParagraphUsingSpell', function(event, data) {