Pārlūkot izejas kodu

[notebook] Integrate hiveserver2 create_session to actual HS2 API and get or create open session

Jenny Kim 10 gadi atpakaļ
vecāks
revīzija
2ec886d

+ 3 - 2
apps/beeswax/src/beeswax/models.py

@@ -420,12 +420,13 @@ class Session(models.Model):
   def get_properties(self):
     return json.loads(self.properties)
 
+  def get_formatted_properties(self):
+    return [dict({'key': key, 'value': value}) for key, value in self.get_properties().items()]
+
   def __str__(self):
     return '%s %s' % (self.owner, self.last_used)
 
 
-
-
 class QueryHandle(object):
   def __init__(self, secret=None, guid=None, operation_type=None, has_result_set=None, modified_row_count=None, log_context=None):
     self.secret = secret

+ 54 - 14
desktop/libs/notebook/src/notebook/connectors/hiveserver2.py

@@ -34,7 +34,7 @@ try:
   from beeswax.api import _autocomplete
   from beeswax.design import hql_query, strip_trailing_semicolon, split_statements
   from beeswax import conf as beeswax_conf
-  from beeswax.models import QUERY_TYPES, HiveServerQueryHandle, QueryHistory, HiveServerQueryHistory
+  from beeswax.models import QUERY_TYPES, HiveServerQueryHandle, HiveServerQueryHistory, QueryHistory, Session
   from beeswax.server import dbms
   from beeswax.server.dbms import get_query_server_config, QueryServerException
   from beeswax.views import _parse_out_hadoop_jobs
@@ -57,22 +57,30 @@ def query_error_handler(func):
 
 class HS2Api(Api):
 
-  def _get_handle(self, snippet):
-    snippet['result']['handle']['secret'], snippet['result']['handle']['guid'] = HiveServerQueryHandle.get_decoded(snippet['result']['handle']['secret'], snippet['result']['handle']['guid'])
-    snippet['result']['handle'].pop('statement_id')
-    snippet['result']['handle'].pop('has_more_statements')
-    return HiveServerQueryHandle(**snippet['result']['handle'])
+  SUPPORTED_TYPES = ('hive', 'impala', 'spark-sql')
 
-  def _get_db(self, snippet):
-    if snippet['type'] == 'hive':
-      name = 'beeswax'
-    elif snippet['type'] == 'impala':
-      name = 'impala'
-    else:
-      name = 'spark-sql'
 
-    return dbms.get(self.user, query_server=get_query_server_config(name=name))
+  @query_error_handler
+  def create_session(self, lang='hive', properties=None):
+    if lang.lower() not in self.SUPPORTED_TYPES:
+      raise PopupException(_('Invalid HS2Api session lang.'))
 
+    if lang == 'hive':
+      lang = 'beeswax'
+
+    session = Session.objects.get_session(self.user, application=lang)
+
+    if session is None:
+      session = dbms.get(self.user, query_server=get_query_server_config(name=lang)).open_session(self.user)
+
+    return {
+        'type': lang,
+        'id': session.id,
+        'properties': session.get_formatted_properties()
+    }
+
+
+  @query_error_handler
   def execute(self, notebook, snippet):
     db = self._get_db(snippet)
 
@@ -119,6 +127,7 @@ class HS2Api(Api):
     hql_query = strip_trailing_semicolon(hql_query)
     return [strip_trailing_semicolon(statement.strip()) for statement in split_statements(hql_query)]
 
+
   @query_error_handler
   def check_status(self, notebook, snippet):
     response = {}
@@ -135,6 +144,7 @@ class HS2Api(Api):
 
     return response
 
+
   @query_error_handler
   def fetch_result(self, notebook, snippet, rows, start_over):
     db = self._get_db(snippet)
@@ -154,10 +164,12 @@ class HS2Api(Api):
         'type': 'table'
     }
 
+
   @query_error_handler
   def fetch_result_metadata(self):
     pass
 
+
   @query_error_handler
   def cancel(self, notebook, snippet):
     db = self._get_db(snippet)
@@ -166,6 +178,7 @@ class HS2Api(Api):
     db.cancel_operation(handle)
     return {'status': 0}
 
+
   @query_error_handler
   def get_log(self, notebook, snippet, startFrom=None, size=None):
     db = self._get_db(snippet)
@@ -173,6 +186,7 @@ class HS2Api(Api):
     handle = self._get_handle(snippet)
     return db.get_log(handle, start_over=startFrom == 0)
 
+
   @query_error_handler
   def close_statement(self, snippet):
     if snippet['type'] == 'impala':
@@ -187,6 +201,8 @@ class HS2Api(Api):
     else:
       return {'status': -1}  # skipped
 
+
+  @query_error_handler
   def download(self, notebook, snippet, format):
     try:
       db = self._get_db(snippet)
@@ -201,6 +217,8 @@ class HS2Api(Api):
         message = e.message
       raise PopupException(message, detail='')
 
+
+  @query_error_handler
   def progress(self, snippet, logs):
     if snippet['type'] == 'hive':
       match = re.search('Total jobs = (\d+)', logs, re.MULTILINE)
@@ -217,6 +235,8 @@ class HS2Api(Api):
     else:
       return 50
 
+
+  @query_error_handler
   def get_jobs(self, notebook, snippet, logs):
     job_ids = _parse_out_hadoop_jobs(logs)
 
@@ -227,12 +247,32 @@ class HS2Api(Api):
 
     return jobs
 
+
   @query_error_handler
   def autocomplete(self, snippet, database=None, table=None, column=None, nested=None):
     db = self._get_db(snippet)
     return _autocomplete(db, database, table, column, nested)
 
+
   def get_select_star_query(self, snippet, database, table):
     db = self._get_db(snippet)
     table = db.get_table(database, table)
     return db.get_select_star_query(database, table)
+
+
+  def _get_handle(self, snippet):
+    snippet['result']['handle']['secret'], snippet['result']['handle']['guid'] = HiveServerQueryHandle.get_decoded(snippet['result']['handle']['secret'], snippet['result']['handle']['guid'])
+    snippet['result']['handle'].pop('statement_id')
+    snippet['result']['handle'].pop('has_more_statements')
+    return HiveServerQueryHandle(**snippet['result']['handle'])
+
+
+  def _get_db(self, snippet):
+    if snippet['type'] == 'hive':
+      name = 'beeswax'
+    elif snippet['type'] == 'impala':
+      name = 'impala'
+    else:
+      name = 'spark-sql'
+
+    return dbms.get(self.user, query_server=get_query_server_config(name=name))

+ 32 - 35
desktop/libs/notebook/src/notebook/urls.py

@@ -32,54 +32,51 @@ import notebook.monkey_patches
 # Views
 urlpatterns = patterns('notebook.views',
   url(r'^$', 'notebook', name='index'),
-  url(r'^notebook$', 'notebook', name='notebook'),
-  url(r'^notebooks$', 'notebooks', name='notebooks'),
-  url(r'^new$', 'new', name='new'),
-  url(r'^download$', 'download', name='download'),
-  url(r'^install_examples$', 'install_examples', name='install_examples'),
-  url(r'^delete$', 'delete', name='delete'),
-  url(r'^copy$', 'copy', name='copy'),
+  url(r'^notebook/?$', 'notebook', name='notebook'),
+  url(r'^notebooks/?$', 'notebooks', name='notebooks'),
+  url(r'^new/?$', 'new', name='new'),
+  url(r'^download/?$', 'download', name='download'),
+  url(r'^install_examples/?$', 'install_examples', name='install_examples'),
+  url(r'^delete/?$', 'delete', name='delete'),
+  url(r'^copy/?$', 'copy', name='copy'),
 
-  url(r'^editor$', 'editor', name='editor'),
-  url(r'^browse/(?P<database>\w+)/(?P<table>\w+)$', 'browse', name='browse'),
+  url(r'^editor/?$', 'editor', name='editor'),
+  url(r'^browse/(?P<database>\w+)/(?P<table>\w+)/?$', 'browse', name='browse'),
 )
 
 # APIs
 urlpatterns += patterns('notebook.api',
-  url(r'^api/create_session$', 'create_session', name='create_session'),
-  url(r'^api/close_session$', 'close_session', name='close_session'),
-  url(r'^api/execute$', 'execute', name='execute'),
-  url(r'^api/check_status$', 'check_status', name='check_status'),
-  url(r'^api/fetch_result_data$', 'fetch_result_data', name='fetch_result_data'),
-  url(r'^api/fetch_result_metadata$', 'fetch_result_metadata', name='fetch_result_metadata'),
-  url(r'^api/cancel_statement$', 'cancel_statement', name='cancel_statement'),
-  url(r'^api/close_statement$', 'close_statement', name='close_statement'),
-  url(r'^api/get_logs$', 'get_logs', name='get_logs'),
+  url(r'^api/create_session/?$', 'create_session', name='create_session'),
+  url(r'^api/close_session/?$', 'close_session', name='close_session'),
+  url(r'^api/execute/?$', 'execute', name='execute'),
+  url(r'^api/check_status/?$', 'check_status', name='check_status'),
+  url(r'^api/fetch_result_data/?$', 'fetch_result_data', name='fetch_result_data'),
+  url(r'^api/fetch_result_metadata/?$', 'fetch_result_metadata', name='fetch_result_metadata'),
+  url(r'^api/cancel_statement/?$', 'cancel_statement', name='cancel_statement'),
+  url(r'^api/close_statement/?$', 'close_statement', name='close_statement'),
+  url(r'^api/get_logs/?$', 'get_logs', name='get_logs'),
 
-  url(r'^api/historify$', 'historify', name='historify'),
-  url(r'^api/get_history', 'get_history', name='get_history'),
-  url(r'^api/clear_history', 'clear_history', name='clear_history'),
+  url(r'^api/historify/?$', 'historify', name='historify'),
+  url(r'^api/get_history/?', 'get_history', name='get_history'),
+  url(r'^api/clear_history/?', 'clear_history', name='clear_history'),
 
-  url(r'^api/notebook/save$', 'save_notebook', name='save_notebook'),
-  url(r'^api/notebook/open$', 'open_notebook', name='open_notebook'),
-  url(r'^api/notebook/close$', 'close_notebook', name='close_notebook'),
+  url(r'^api/notebook/save/?$', 'save_notebook', name='save_notebook'),
+  url(r'^api/notebook/open/?$', 'open_notebook', name='open_notebook'),
+  url(r'^api/notebook/close/?$', 'close_notebook', name='close_notebook'),
 )
 
 # Github
 urlpatterns += patterns('notebook.api',
-  url(r'^api/github/fetch$', 'github_fetch', name='github_fetch'),
-  url(r'^api/github/authorize', 'github_authorize', name='github_authorize'),
-  url(r'^api/github/callback', 'github_callback', name='github_callback'),
+  url(r'^api/github/fetch/?$', 'github_fetch', name='github_fetch'),
+  url(r'^api/github/authorize/?$', 'github_authorize', name='github_authorize'),
+  url(r'^api/github/callback/?$', 'github_callback', name='github_callback'),
 )
 
 # Assist API
 urlpatterns += patterns('notebook.api',
-  url(r'^api/autocomplete/$', 'autocomplete', name='api_autocomplete_databases'),
-  url(r'^api/autocomplete/(?P<database>\w+)$', 'autocomplete', name='api_autocomplete_tables'),
-  url(r'^api/autocomplete/(?P<database>\w+)/$', 'autocomplete', name='api_autocomplete_tables'),
-  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)$', 'autocomplete', name='api_autocomplete_columns'),
-  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/$', 'autocomplete', name='api_autocomplete_columns'),
-  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/(?P<column>\w+)$', 'autocomplete', name='api_autocomplete_column'),
-  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/(?P<column>\w+)/$', 'autocomplete', name='api_autocomplete_column'),
-  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/(?P<column>\w+)/(?P<nested>.+)$', 'autocomplete', name='api_autocomplete_nested'),
+  url(r'^api/autocomplete/?$', 'autocomplete', name='api_autocomplete_databases'),
+  url(r'^api/autocomplete/(?P<database>\w+)/?$', 'autocomplete', name='api_autocomplete_tables'),
+  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/?$', 'autocomplete', name='api_autocomplete_columns'),
+  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/(?P<column>\w+)/?$', 'autocomplete', name='api_autocomplete_column'),
+  url(r'^api/autocomplete/(?P<database>\w+)/(?P<table>\w+)/(?P<column>\w+)/(?P<nested>.+)/?$', 'autocomplete', name='api_autocomplete_nested'),
 )