Prechádzať zdrojové kódy

HUE-9280 [flink] Send live results via check status

Romain 5 rokov pred
rodič
commit
043977fb17

+ 1 - 1
desktop/core/src/desktop/js/api/apiHelper.js

@@ -1991,7 +1991,7 @@ class ApiHelper {
     })
       .done(response => {
         if (response && response.query_status) {
-          deferred.resolve(response.query_status.status);
+          deferred.resolve(response.query_status);
         } else if (response && response.status === -3) {
           deferred.resolve(EXECUTION_STATUS.expired);
         } else {

+ 2 - 2
desktop/core/src/desktop/js/apps/notebook2/components/ko.snippetResults.js

@@ -115,8 +115,8 @@ const TEMPLATE = `
       <div data-bind="visible: executing" style="display: none;">
         <h1 class="empty"><i class="fa fa-spinner fa-spin"></i> ${ I18n('Executing...') }</h1>
       </div>
-      <div id="wsResult">
-      </div>
+      <ul id="wsResult">
+      </ul>
     </div>
   </div>
 </div>

+ 10 - 7
desktop/core/src/desktop/js/apps/notebook2/execution/executable.js

@@ -259,15 +259,15 @@ export default class Executable {
 
     this.cancellables.push(
       apiHelper.checkExecutionStatus({ executable: this }).done(async queryStatus => {
-        switch (queryStatus) {
+        switch (queryStatus.status) {
           case EXECUTION_STATUS.success:
             this.executeEnded = Date.now();
-            this.setStatus(queryStatus);
+            this.setStatus(queryStatus.status);
             this.setProgress(99); // TODO: why 99 here (from old code)?
             break;
           case EXECUTION_STATUS.available:
             this.executeEnded = Date.now();
-            this.setStatus(queryStatus);
+            this.setStatus(queryStatus.status);
             this.setProgress(100);
             if (!this.result && this.handle.has_result_set) {
               this.result = new ExecutionResult(this);
@@ -283,12 +283,15 @@ export default class Executable {
           case EXECUTION_STATUS.canceled:
           case EXECUTION_STATUS.expired:
             this.executeEnded = Date.now();
-            this.setStatus(queryStatus);
+            this.setStatus(queryStatus.status);
             break;
           case EXECUTION_STATUS.running:
+            if (queryStatus.data) {
+              huePubSub.publish('editor.ws.query.fetch_result', queryStatus);
+            }
           case EXECUTION_STATUS.starting:
           case EXECUTION_STATUS.waiting:
-            this.setStatus(queryStatus);
+            this.setStatus(queryStatus.status);
             checkStatusTimeout = window.setTimeout(
               () => {
                 this.checkStatus(statusCheckCount);
@@ -298,11 +301,11 @@ export default class Executable {
             break;
           case EXECUTION_STATUS.failed:
             this.executeEnded = Date.now();
-            this.setStatus(queryStatus);
+            this.setStatus(queryStatus.status);
             break;
           default:
             this.executeEnded = Date.now();
-            console.warn('Got unknown status ' + queryStatus);
+            console.warn('Got unknown status ' + queryStatus.status);
         }
       })
     );

+ 1 - 1
desktop/core/src/desktop/js/apps/notebook2/execution/executionResult.js

@@ -70,7 +70,7 @@ huePubSub.subscribe('editor.ws.query.fetch_result', executionResult => {
     result.fetchedOnce = true;
     result.handleResultResponse(executionResult);
     // eslint-disable-next-line no-undef
-    $('#wsResult').append('<div>' + executionResult.data + '</div>');
+    $('#wsResult').append('<li>' + executionResult.data + '</li>');
   }
 });
 

+ 5 - 3
desktop/libs/notebook/src/notebook/connectors/flink_sql.py

@@ -139,6 +139,7 @@ class FlinkSqlApi(Api):
   @query_error_handler
   def check_status(self, notebook, snippet):
     global n
+    response = {}
     session = self._get_session()
     statement_id = snippet['result']['handle']['guid']
 
@@ -152,8 +153,7 @@ class FlinkSqlApi(Api):
           resp = self.db.fetch_status(session['id'], statement_id)
           if resp.get('status') == 'RUNNING':
             status = 'running'
-            # if n >= 5:
-            print(self.fetch_result(notebook, snippet, n, False)['data'])
+            response.update(self.fetch_result(notebook, snippet, n, False))
           elif resp.get('status') == 'FINISHED':
             status = 'available'
           elif resp.get('status') == 'FAILED':
@@ -166,7 +166,9 @@ class FlinkSqlApi(Api):
           else:
             raise e
 
-    return {'status': status}
+    response['status'] = status
+
+    return response
 
 
   @query_error_handler