소스 검색

HUE-9280 [flink] Cancel and check status calls

Romain 5 년 전
부모
커밋
11a13829cd

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

@@ -432,7 +432,7 @@ export default class Executable {
     };
 
     return {
-      operationId: undefined,
+      operationId: this.operationId,
       snippet: JSON.stringify(snippet),
       notebook: JSON.stringify(notebook)
     };

+ 47 - 8
desktop/libs/notebook/src/notebook/connectors/flink_sql.py

@@ -37,6 +37,8 @@ _API_VERSION = 'v1'
 SESSIONS = {}
 SESSION_KEY = '%(username)s-%(connector_name)s'
 
+n = 0
+
 
 def query_error_handler(func):
   def decorator(*args, **kwargs):
@@ -55,6 +57,7 @@ def query_error_handler(func):
   return decorator
 
 
+
 class FlinkSqlApi(Api):
 
   def __init__(self, user, interpreter=None):
@@ -88,7 +91,7 @@ class FlinkSqlApi(Api):
       self.db.session_heartbeat(session_id=SESSIONS[session_key]['id'])
     except Exception as e:
       if 'Session: %(id)s does not exist' % SESSIONS[session_key] in str(e):
-        LOG.info('Session: %(id)s does not exist, opening a new one' % SESSIONS[session_key])
+        LOG.warn('Session: %(id)s does not exist, opening a new one' % SESSIONS[session_key])
         SESSIONS[session_key] = self.create_session()
       else:
         raise e
@@ -97,6 +100,8 @@ class FlinkSqlApi(Api):
 
   @query_error_handler
   def execute(self, notebook, snippet):
+    global n
+    n = 0
     session = self._get_session()
     session_id = session['id']
     job_id = None
@@ -106,6 +111,7 @@ class FlinkSqlApi(Api):
     if resp['statement_types'][0] == 'SELECT':
       job_id = resp['results'][0]['data'][0][0]
       data, description = [], []
+      # TODO: change_flags
     else:
       data, description = resp['results'][0]['data'], resp['results'][0]['columns']
 
@@ -132,6 +138,7 @@ class FlinkSqlApi(Api):
 
   @query_error_handler
   def check_status(self, notebook, snippet):
+    global n
     session = self._get_session()
     statement_id = snippet['result']['handle']['guid']
 
@@ -141,25 +148,42 @@ class FlinkSqlApi(Api):
       if not statement_id:  # Sync result
         status = 'available'
       else:
-        resp = self.db.fetch_status(session['id'], statement_id)
-        if resp.get('status') == 'RUNNING':
-          status = 'running'
-        elif resp.get('status') == 'FINISHED':
-          status = 'available'
+        try:
+          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'])
+          elif resp.get('status') == 'FINISHED':
+            status = 'available'
+          elif resp.get('status') == 'FAILED':
+            status = 'failed'
+          elif resp.get('status') == 'CANCELED':
+            status = 'expired'
+        except Exception as e:
+          if '%s does not exist in current session' % statement_id in str(e):
+            LOG.warn('Job: %s does not exist' % statement_id)
+          else:
+            raise e
 
     return {'status': status}
 
 
   @query_error_handler
   def fetch_result(self, notebook, snippet, rows, start_over):
+    global n
     session = self._get_session()
     statement_id = snippet['result']['handle']['guid']
-    token = 0
+    token = n #rows
 
     resp = self.db.fetch_results(session['id'], job_id=statement_id, token=token)
 
+    next_result = resp.get('next_result_uri')
+    if next_result:
+      n = int(next_result.rsplit('/', 1)[-1])
+
     return {
-        'has_more': bool(resp.get('next_result_uri')) and False,  # TODO: here we should increment the token
+        'has_more': bool(next_result),
         'data': resp['results'][0]['data'],  # No escaping...
         'meta': [{
             'name': column['name'],
@@ -221,6 +245,21 @@ class FlinkSqlApi(Api):
     return [table[0] for table in resp['results'][0]['data']]
 
 
+  def cancel(self, notebook, snippet):
+    session = self._get_session()
+    statement_id = snippet['result']['handle']['guid']
+
+    try:
+      self.db.close_statement(session_id=session['id'], job_id=statement_id)
+    except Exception as e:
+      if 'does not exist in current session:' in str(e):
+        return {'status': -1}  # skipped
+      else:
+        raise e
+
+    return {'status': 0}
+
+
   def close_session(self, session):
     session = self._get_session()
     self.db.close_session(session['id'])

+ 1 - 1
desktop/libs/notebook/src/notebook/decorators.py

@@ -130,7 +130,7 @@ def api_error_handler(f):
       response['status'] = -4
     except FilesystemException as e:
       response['status'] = 2
-      response['message'] = e.message
+      response['message'] = e.message or 'Query history not found'
     except QueryError as e:
       LOG.exception('Error running %s' % f.__name__)
       response['status'] = 1