浏览代码

HUE-6245 [dataeng] Integrate the API in the editor as a new interpreter

Romain Rigaux 8 年之前
父节点
当前提交
f3a37ed1eb

+ 1 - 1
apps/jobbrowser/src/jobbrowser/apis/data_eng_api.py

@@ -24,7 +24,7 @@ from django.utils.translation import ugettext as _
 
 from jobbrowser.apis.base_api import Api, MockDjangoRequest, _extract_query_params
 from liboozie.oozie_api import get_oozie
-from notebook.connectors.dataeng_batch import DataEng, DATE_FORMAT
+from notebook.connectors.dataeng import DataEng, DATE_FORMAT
 
 
 LOG = logging.getLogger(__name__)

+ 3 - 0
desktop/libs/notebook/src/notebook/conf.py

@@ -202,6 +202,9 @@ def _default_interpreters(user):
       ('pyspark', {
           'name': 'PySpark', 'interface': 'livy', 'options': {}
       }),
+      ('dataeng', {
+          'name': 'DataEng', 'interface': 'dataeng', 'options': {}
+      }),
       ('r', {
           'name': 'R', 'interface': 'livy', 'options': {}
       }),

+ 3 - 3
desktop/libs/notebook/src/notebook/connectors/base.py

@@ -191,7 +191,7 @@ class Notebook(object):
 
 
 def get_api(request, snippet):
-  from notebook.connectors.dataeng_batch import DataEngBatchApi
+  from notebook.connectors.dataeng import DataEngApi
   from notebook.connectors.hiveserver2 import HS2Api
   from notebook.connectors.jdbc import JdbcApi
   from notebook.connectors.rdbms import RdbmsApi
@@ -222,8 +222,8 @@ def get_api(request, snippet):
     return TextApi(request.user)
   elif interface == 'rdbms':
     return RdbmsApi(request.user, interpreter=snippet['type'])
-  elif interface == 'dataeng-batch':
-    return DataEngBatchApi(user=request.user, request=request)  
+  elif interface == 'dataeng':
+    return DataEngApi(user=request.user, request=request)  
   elif interface == 'jdbc':
     return JdbcApi(request.user, interpreter=interpreter)
   elif interface == 'solr':

+ 16 - 17
desktop/libs/notebook/src/notebook/connectors/dataeng_batch.py → desktop/libs/notebook/src/notebook/connectors/dataeng.py

@@ -44,32 +44,31 @@ def _exec(args):
     raise PopupException(e, title=_('Error accessing'))
 
   response = json.loads(data)
+  # Chck data['status'] == 'success'
   response['status'] = 'success'
 
   return response
 
 DATE_FORMAT = "%Y-%m-%d"
+RUNNING_STATES = ('QUEUED', 'RUNNING')
 
 
-class DataEngBatchApi(Api):
+class DataEngApi(Api):
 
 
   def execute(self, notebook, snippet):
-    db = self._get_db(snippet)
+    statement = snippet['statement']
+    cluster_name = 'romain-cluster'
 
-    statement = self._get_current_statement(db, snippet)
-    session = self._get_session(notebook, snippet['type'])
+    handle = DataEng(self.user).submit_hive_job(cluster_name, statement, params=None, job_xml=None)
+    job = handle['job']
 
-    query = self._prepare_hql_query(snippet, statement['statement'], session)
-
-    handle = DataEng().submit_hive_job(cluster_name, query, params=None, job_xml=None)
-
-    if handle['status'] not in ('QUEUED', 'RUNNING'):
-      raise QueryError('Submission failure', handle=statement)
+    if job['status'] not in RUNNING_STATES:
+      raise QueryError('Submission failure', handle=job['status'])
 
     return {
-      'id': handle['jobType'],
-      'crn': handle['crn'],
+      'id': job['jobId'],
+      'crn': job['crn'],
       'has_result_set': False,
     }
 
@@ -79,11 +78,11 @@ class DataEngBatchApi(Api):
 
     job_id = snippet['result']['handle']['id']
 
-    handle = DataEng().list_jobs(job_ids=[job_id])
+    handle = DataEng(self.user).list_jobs(job_ids=[job_id])
 
-    if handle['status'] in ('QUEUED', 'RUNNING'):
+    if handle['status'] in RUNNING_STATES:
       return response
-    elif handle['status'] in ('KILLED', 'FAILED'):
+    elif handle['status'] in ('failed', 'terminated'):
       raise QueryError(_('Job was %s') % handle['status'])
     else:
       response['status'] = 'available'
@@ -103,7 +102,7 @@ class DataEngBatchApi(Api):
   def cancel(self, notebook, snippet):
     job_id = snippet['result']['handle']['id']
 
-    DataEng().terminate_jobs(job_ids=[job_id])
+    DataEng(self.user).terminate_jobs(job_ids=[job_id])
 
     return {'status': 0}
 
@@ -168,7 +167,7 @@ class DataEng():
     if job_statuses:
       args.extend(['--job-statuses', job_statuses])
     if job_ids:
-      args.extend(['--job-ids', job_ids])
+      args.extend(['--job-ids'] + job_ids)
     if job_types:
       args.extend(['--job-types', job_types])
     if creation_date_before: