소스 검색

[flink] Handle better expired queries (#2032)

- [flink] Handle completely empty results
- [meta] Avoid error message when no top columns are recommended
- [flink] Avoid false positive error when closing empty job
Romain Rigaux 4 년 전
부모
커밋
be0a6ec893

+ 5 - 2
desktop/libs/dashboard/src/dashboard/models.py

@@ -1003,8 +1003,11 @@ def extract_solr_exception_message(e):
 
   try:
     message = json.loads(e.message)
-    msg = message['error'].get('msg')
-    response['error'] = msg if msg else message['error']['trace']
+    if 'error' in message:
+      msg = message['error'].get('msg')
+      response['error'] = msg if msg else message['error']['trace']
+    else:
+      response['error'] = message['errors'][0] if message['errors'] else 'Empty errors'
   except ValueError as e:
     LOG.warning('Failed to parse json response: %s' % force_unicode(e))
     response['error'] = force_unicode(e)

+ 4 - 0
desktop/libs/kafka/src/kafka/ksql_client.py

@@ -189,6 +189,10 @@ class KSqlApi(object):
     return data, metadata
 
 
+  def cancel(self, notebook, snippet):
+    return {'status': -1}
+
+
   def _decode_result(self, result):
     columns = []
     data = []

+ 3 - 5
desktop/libs/metadata/src/metadata/optimizer_api.py

@@ -369,11 +369,9 @@ def top_columns(request):
 
   data = api.top_columns(db_tables=db_tables, connector=connector)
 
-  if data:
-    response['status'] = 0
-    response['values'] = data
-  else:
-    response['message'] = 'Optimizer: %s' % data
+  response['status'] = 0
+  response['values'] = data or []
+  response['message'] = 'Optimizer: %s' % data
 
   return JsonResponse(response)
 

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

@@ -198,13 +198,13 @@ class FlinkSqlApi(Api):
 
     return {
         'has_more': bool(next_result),
-        'data': resp['results'][0]['data'],  # No escaping...
+        'data': resp and resp['results'][0]['data'] or [],  # No escaping...
         'meta': [{
             'name': column['name'],
             'type': column['type'],
             'comment': ''
           }
-          for column in resp['results'][0]['columns']
+          for column in resp['results'][0]['columns'] if resp
         ],
         'type': 'table'
     }
@@ -256,7 +256,10 @@ class FlinkSqlApi(Api):
     statement_id = snippet['result']['handle']['guid']
 
     try:
-      self.db.close_statement(session_id=session['id'], job_id=statement_id)
+      if session and statement_id:
+        self.db.close_statement(session_id=session['id'], job_id=statement_id)
+      else:
+        return {'status': -1} # missing operation ids
     except Exception as e:
       if 'does not exist in current session:' in str(e):
         return {'status': -1}  # skipped