浏览代码

HUE-2936 [jb] Support delegation token for YARN kill button

Romain Rigaux 10 年之前
父节点
当前提交
a691ed721f

+ 2 - 2
apps/jobbrowser/src/jobbrowser/api.py

@@ -213,8 +213,8 @@ class YarnApi(JobBrowserApi):
   """
   def __init__(self, user):
     self.user = user
-    self.resource_manager_api = resource_manager_api.get_resource_manager(user)
-    self.mapreduce_api = mapreduce_api.get_mapreduce_api()
+    self.resource_manager_api = resource_manager_api.get_resource_manager(user.username)
+    self.mapreduce_api = mapreduce_api.get_mapreduce_api(user.username)
     self.history_server_api = history_server_api.get_history_server_api()
 
   def get_job_link(self, job_id):

+ 4 - 4
apps/jobbrowser/src/jobbrowser/tests.py

@@ -395,8 +395,8 @@ class TestMapReduce2NoHadoop:
     grant_access("test2", "test2", "jobbrowser")
     self.user2 = User.objects.get(username='test2')
 
-    resource_manager_api.get_resource_manager = lambda user: MockResourceManagerApi(user)
-    mapreduce_api.get_mapreduce_api = lambda: MockMapreduceApi()
+    resource_manager_api.get_resource_manager = lambda user: MockResourceManagerApi(user.username)
+    mapreduce_api.get_mapreduce_api = lambda user: MockMapreduceApi(user.username)
     history_server_api.get_history_server_api = lambda: HistoryServerApi()
 
     self.finish = [
@@ -606,7 +606,7 @@ class MockMapreduce2Api(object):
   MockMapreduceApi and HistoryServerApi are very similar and inherit from it.
   """
 
-  def __init__(self, oozie_url=None): pass
+  def __init__(self, mr_url=None): pass
 
   def tasks(self, job_id):
     return {
@@ -746,7 +746,7 @@ class MockMapreduceApi(MockMapreduce2Api):
 
 class HistoryServerApi(MockMapreduce2Api):
 
-  def __init__(self, oozie_url=None): pass
+  def __init__(self, hs_url=None): pass
 
   def job(self, user, job_id):
     if '1356251510842_0054' == job_id:

+ 5 - 1
apps/jobbrowser/src/jobbrowser/views.py

@@ -242,7 +242,11 @@ def kill_job(request, job):
     access_warn(request, _('Insufficient permission'))
     raise MessageException(_("Permission denied.  User %(username)s cannot delete user %(user)s's job.") % {'username': request.user.username, 'user': job.user})
 
-  job.kill()
+  try:
+    job.kill()
+  except Exception, e:
+    LOGGER.exception('Killing job')
+    raise PopupException(e)
 
   cur_time = time.time()
   api = get_api(request.user, request.jt)

+ 3 - 2
desktop/core/src/desktop/lib/rest/resource.py

@@ -97,15 +97,16 @@ class Resource(object):
     return self.invoke("GET", relpath, params, headers=headers, allow_redirects=True)
 
 
-  def delete(self, relpath=None, params=None):
+  def delete(self, relpath=None, params=None, headers=None):
     """
     Invoke the DELETE method on a resource.
     @param relpath: Optional. A relative path to this resource's path.
     @param params: Key-value data.
+    @param headers: Optional. Base set of headers.
 
     @return: A dictionary of the JSON result.
     """
-    return self.invoke("DELETE", relpath, params)
+    return self.invoke("DELETE", relpath, params, headers=headers)
 
 
   def post(self, relpath=None, params=None, data=None, contenttype=None, headers=None):

+ 6 - 5
desktop/libs/hadoop/src/hadoop/yarn/mapreduce_api.py

@@ -35,14 +35,14 @@ _api_cache = None
 _api_cache_lock = threading.Lock()
 
 
-def get_mapreduce_api():
+def get_mapreduce_api(user):
   global _api_cache
   if _api_cache is None:
     _api_cache_lock.acquire()
     try:
       if _api_cache is None:
         yarn_cluster = cluster.get_cluster_conf_for_job_submission()
-        _api_cache = MapreduceApi(yarn_cluster.PROXY_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
+        _api_cache = MapreduceApi(user, yarn_cluster.PROXY_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
     finally:
       _api_cache_lock.release()
   return _api_cache
@@ -50,8 +50,9 @@ def get_mapreduce_api():
 
 class MapreduceApi(object):
 
-  def __init__(self, oozie_url, security_enabled=False, ssl_cert_ca_verify=False):
-    self._url = posixpath.join(oozie_url, 'proxy')
+  def __init__(self, user, mr_url, security_enabled=False, ssl_cert_ca_verify=False):
+    self._user = user
+    self._url = posixpath.join(mr_url, 'proxy')
     self._client = HttpClient(self._url, logger=LOG)
     self._root = Resource(self._client)
     self._security_enabled = security_enabled
@@ -114,4 +115,4 @@ class MapreduceApi(object):
 
   def kill(self, job_id):
     app_id = job_id.replace('job', 'application')
-    get_resource_manager().kill(app_id) # We need to call the RM
+    get_resource_manager(self._user).kill(app_id) # We need to call the RM

+ 28 - 1
desktop/libs/hadoop/src/hadoop/yarn/resource_manager_api.py

@@ -24,6 +24,7 @@ from django.utils.translation import ugettext as _
 
 from desktop.conf import DEFAULT_USER
 from desktop.lib.exceptions_renderable import PopupException
+from desktop.lib.i18n import smart_str
 from desktop.lib.rest.http_client import HttpClient
 from desktop.lib.rest.resource import Resource
 
@@ -108,8 +109,34 @@ class ResourceManagerApi(object):
     return self._execute(self._root.get, 'cluster/apps/%(app_id)s' % {'app_id': app_id}, params=params, headers={'Accept': _JSON_CONTENT_TYPE})
 
   def kill(self, app_id):
+    data = {'state': 'KILLED'}
+    token = None
+
+    # Tokens are managed within the kill method but should be moved out when not alpha anymore or we support submitting an app.
+    if self.security_enabled:
+      full_token = self.delegation_token()
+      if 'token' not in full_token:
+        raise PopupException(_('YARN did not return any token field.'), detail=smart_str(full_token))
+      data['X-Hadoop-Delegation-Token'] = token = full_token.pop('token')
+      LOG.debug('Received delegation token %s' % full_token)
+
+    try:
+      params = self._get_params()
+      return self._execute(self._root.put, 'cluster/apps/%(app_id)s/state' % {'app_id': app_id}, params=params, data=json.dumps(data), contenttype=_JSON_CONTENT_TYPE)
+    finally:
+      if token:
+        self.cancel_token(token)
+
+  def delegation_token(self):
+    params = self._get_params()
+    data = {'renewer': self._user}
+    return self._execute(self._root.post, 'cluster/delegation-token', params=params, data=json.dumps(data), contenttype=_JSON_CONTENT_TYPE)
+
+  def cancel_token(self, token):
     params = self._get_params()
-    return self._execute(self._root.put, 'cluster/apps/%(app_id)s/state' % {'app_id': app_id}, params=params, data=json.dumps({'state': 'KILLED'}), contenttype=_JSON_CONTENT_TYPE)
+    headers = {'Hadoop-YARN-RM-Delegation-Token': token}
+    LOG.debug('Canceling delegation token of ' % self._user)
+    return self._execute(self._root.delete, 'cluster/delegation-token', params=params, headers=headers)
 
   def _execute(self, function, *args, **kwargs):
     response = function(*args, **kwargs)