Эх сурвалжийг харах

[beeswax] Support getting a specific session by ID, and add a close_session API endpoint

Beeswax/Impala API Updates:

GET /beeswax/api/sessions: returns latest open session if found
GET /beeswax/api/sessions/<id>: return session for given session ID
POST /beeswax/api/session/<id>/close: attempt to close session with given ID
Jenny Kim 10 жил өмнө
parent
commit
796319a

+ 37 - 5
apps/beeswax/src/beeswax/api.py

@@ -23,6 +23,7 @@ from django.contrib.auth.models import User
 from django.core.urlresolvers import reverse
 from django.http import Http404
 from django.utils.translation import ugettext as _
+from django.views.decorators.http import require_POST
 
 from thrift.transport.TTransport import TTransportException
 from desktop.context_processors import get_app_name
@@ -802,22 +803,53 @@ def get_top_terms(request, database, table, column, prefix=None):
 
 
 @error_handler
-def get_session(request):
+def get_session(request, session_id=None):
   app_name = get_app_name(request)
   query_server = get_query_server_config(app_name)
 
-  session = Session.objects.get_session(request.user, query_server['server_name'])
+  response = {'status': -1, 'message': ''}
 
-  if session:
+  if session_id:
+    session = Session.objects.get(id=session_id, owner=request.user, application=query_server['server_name'])
+  else:  # get the latest session for given user and server type
+    session = Session.objects.get_session(request.user, query_server['server_name'])
+
+  if session is not None:
     properties = json.loads(session.properties)
     # Redact passwords
     for key, value in properties.items():
       if 'password' in key.lower():
         properties[key] = '*' * len(value)
+
+    response['status'] = 0
+    response['session'] = {'id': session.id, 'application': session.application, 'status': session.status_code}
+    response['properties'] = properties
   else:
-    properties = {}
+    response['message'] = _('Could not find session or no open sessions found.')
+
+  return JsonResponse(response)
+
+
+@require_POST
+@error_handler
+def close_session(request, session_id):
+  app_name = get_app_name(request)
+  query_server = get_query_server_config(app_name)
+
+  response = {'status': -1, 'message': ''}
+
+  try:
+    session = Session.objects.get(id=session_id, owner=request.user, application=query_server['server_name'])
+  except Session.DoesNotExist:
+    response['message'] = _('Session does not exist or you do not have permissions to close the session.')
 
-  return JsonResponse({'properties': properties})
+  if session:
+    session = dbms.get(request.user, query_server).close_session(session)
+    response['status'] = 0
+    response['message'] = _('Session successfully closed.')
+    response['session'] = {'id': session_id, 'application': session.application, 'status': session.status_code}
+
+  return JsonResponse(response)
 
 
 """

+ 11 - 16
apps/beeswax/src/beeswax/management/commands/close_sessions.py

@@ -45,17 +45,19 @@ class Command(BaseCommand):
 
     self.stdout.write('Closing (all=%s) HiveServer2 sessions older than %s days...\n' % (close_all, days))
 
-    sessions = Session.objects.all()
+    sessions = Session.objects.filter(status_code=0)
 
     if not close_all:
       sessions = sessions.filter(application='beeswax')
 
     sessions = sessions.filter(last_used__lte=datetime.today() - timedelta(days=days))
 
+    self.stdout.write('Found %d open HiveServer2 sessions to close' % len(sessions))
+
     import os
     import beeswax
-    from beeswax import conf
     from beeswax import hive_site
+
     try:
       beeswax.conf.HIVE_CONF_DIR.set_for_testing(os.environ['HIVE_CONF_DIR'])
     except:
@@ -66,21 +68,14 @@ class Command(BaseCommand):
     hive_site.reset()
     hive_site.get_conf()
 
-    closed_sessions = 0
-    already_closed_sessions = 0
-
+    closed = 0
+    skipped = 0
     for session in sessions:
       try:
-        resp = dbms.get(user=session.owner).close_session(session)
-        if not 'Session does not exist!' in str(resp):
-          self.stdout.write('Info: %s\n' % resp)
-          closed_sessions += 1
-        else:
-          already_closed_sessions += 1
+        session = dbms.get(user=session.owner).close_session(session)
+        closed += 1
       except Exception, e:
-        if 'Session does not exist!' in str(e):
-          already_closed_sessions += 1
-        else:
-          self.stdout.write('Info: %s\n' % e)
+        skipped += 1
+        self.stdout.write('Session with ID %d could not be closed: %s' % (session.id, str(e)))
 
-    self.stdout.write('%s sessions closed. %s sessions already closed.\n' % (closed_sessions, already_closed_sessions))
+    self.stdout.write('%d sessions closed.\n%d sessions skipped because already closed.' % (closed, skipped))

+ 9 - 5
apps/beeswax/src/beeswax/models.py

@@ -385,11 +385,15 @@ class SavedQuery(models.Model):
 
 
 class SessionManager(models.Manager):
-  def get_session(self, user, application='beeswax'):
+
+  def get_session(self, user, application='beeswax', open_sessions=True):
     try:
-      return self.filter(owner=user, application=application).latest("last_used")
-    except Session.DoesNotExist:
-      pass
+      q = self.filter(owner=user, application=application)
+      if open_sessions:
+        q = q.filter(status_code=0)
+      return q.latest("last_used")
+    except Session.DoesNotExist, e:
+      return None
 
 
 class Session(models.Model):
@@ -397,7 +401,7 @@ class Session(models.Model):
   A sessions is bound to a user and an application (e.g. Bob with the Impala application).
   """
   owner = models.ForeignKey(User, db_index=True)
-  status_code = models.PositiveSmallIntegerField()
+  status_code = models.PositiveSmallIntegerField()  # ttypes.TStatusCode
   secret = models.TextField(max_length='100')
   guid = models.TextField(max_length='100')
   server_protocol_version = models.SmallIntegerField(default=0)

+ 14 - 1
apps/beeswax/src/beeswax/server/dbms.py

@@ -293,11 +293,24 @@ class HiveServer2Dbms(object):
   def close_operation(self, query_handle):
     return self.client.close_operation(query_handle)
 
+
   def open_session(self, user):
     return self.client.open_session(user)
 
+
   def close_session(self, session):
-    return self.client.close_session(session)
+    resp = self.client.close_session(session)
+
+    if resp.status.statusCode != 0:
+      session.status_code = resp.status.statusCode
+      session.save()
+      raise QueryServerException(_('Failed to close session, session handle may already be closed or timed out.'))
+    else:
+      session.status_code = 4  # Set to ttypes.TStatusCode.INVALID_HANDLE_STATUS
+      session.save()
+
+    return session
+
 
   def cancel_operation(self, query_handle):
     resp = self.client.cancel_operation(query_handle)

+ 3 - 1
apps/beeswax/src/beeswax/urls.py

@@ -57,7 +57,9 @@ urlpatterns += patterns(
 urlpatterns += patterns(
   'beeswax.api',
 
-  url(r'^api/session/$', 'get_session', name='api_get_session'),
+  url(r'^api/session/?$', 'get_session', name='api_get_session'),
+  url(r'^api/session/(?P<session_id>\d+)/?$', 'get_session', name='api_get_session'),
+  url(r'^api/session/(?P<session_id>\d+)/close/?$', 'close_session', name='api_close_session'),
   url(r'^api/functions/?$', 'get_functions', name='get_functions'),
 
   # Deprecated by Notebook API