瀏覽代碼

HUE-2935 [hadoop] Support YARN impersonation

Like for the other services Hue will need to be a proxy user, e.g. in core-site.xml:

  <property>
    <name>hadoop.proxyuser.hue.hosts</name>
    <value>*</value>
  </property>
  <property>
    <name>hadoop.proxyuser.hue.groups</name>
    <value>*</value>
  </property>
Romain Rigaux 10 年之前
父節點
當前提交
91a9336

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

@@ -213,7 +213,7 @@ class YarnApi(JobBrowserApi):
   """
   def __init__(self, user):
     self.user = user
-    self.resource_manager_api = resource_manager_api.get_resource_manager()
+    self.resource_manager_api = resource_manager_api.get_resource_manager(user)
     self.mapreduce_api = mapreduce_api.get_mapreduce_api()
     self.history_server_api = history_server_api.get_history_server_api()
 
@@ -302,7 +302,7 @@ class YarnApi(JobBrowserApi):
     return self.get_job(jobid).task(task_id)
 
   def get_tracker(self, node_manager_http_address, container_id):
-    api = node_manager_api.get_resource_manager_api('http://' + node_manager_http_address)
+    api = node_manager_api.get_node_manager_api('http://' + node_manager_http_address)
     return Container(api.container(container_id))
 
 

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

@@ -577,7 +577,7 @@ class MockResourceManagerApi:
     },
   }
 
-  def __init__(self, oozie_url=None): pass
+  def __init__(self, user, rm_url=None): pass
 
   def apps(self, **kwargs):
     return {

+ 14 - 14
apps/jobbrowser/src/jobbrowser/views.py

@@ -66,7 +66,7 @@ def check_job_permission(view_func):
       job = get_api(request.user, request.jt).get_job(jobid=jobid)
     except ApplicationNotRunning, e:
       if e.job.get('state', '').lower() == 'accepted' and 'kill' in request.path:
-        rm_api = resource_manager_api.get_resource_manager()
+        rm_api = resource_manager_api.get_resource_manager(request.user)
         job = Application(e.job, rm_api)
       else:
         # reverse() seems broken, using request.path but beware, it discards GET and POST info
@@ -108,21 +108,21 @@ def jobs(request):
   retired = request.GET.get('retired')
 
   if request.GET.get('format') == 'json':
-    try:
+#    try:
       # Limit number of jobs to be 10,000
       jobs = get_api(request.user, request.jt).get_jobs(user=request.user, username=user, state=state, text=text, retired=retired, limit=10000)
-    except Exception, ex:
-      ex_message = str(ex)
-      if 'Connection refused' in ex_message or 'standby RM' in ex_message:
-        raise PopupException(_('Resource Manager cannot be contacted or might be down.'))
-      elif 'Could not connect to' in ex_message:
-        raise PopupException(_('Job Tracker cannot be contacted or might be down.'))
-      else:
-        raise ex
-    json_jobs = {
-      'jobs': [massage_job_for_json(job, request) for job in jobs],
-    }
-    return JsonResponse(json_jobs, encoder=JSONEncoderForHTML)
+#    except Exception, ex:
+#      ex_message = str(ex)
+#      if 'Connection refused' in ex_message or 'standby RM' in ex_message:
+#        raise PopupException(_('Resource Manager cannot be contacted or might be down.'))
+#      elif 'Could not connect to' in ex_message:
+#        raise PopupException(_('Job Tracker cannot be contacted or might be down.'))
+#      else:
+#        raise ex
+      json_jobs = {
+        'jobs': [massage_job_for_json(job, request) for job in jobs],
+      }
+      return JsonResponse(json_jobs, encoder=JSONEncoderForHTML)
 
   return render('jobs.mako', request, {
     'request': request,

+ 3 - 3
desktop/libs/hadoop/src/hadoop/conf.py

@@ -191,7 +191,7 @@ def config_validator(user):
   for name in YARN_CLUSTERS.keys():
     cluster = YARN_CLUSTERS[name]
     if cluster.SUBMIT_TO.get():
-      res.extend(test_yarn_configurations())
+      res.extend(test_yarn_configurations(user))
       submit_to.append('yarn_clusters.' + name)
 
   if not submit_to:
@@ -201,7 +201,7 @@ def config_validator(user):
   return res
 
 
-def test_yarn_configurations():
+def test_yarn_configurations(user):
   # Single cluster for now
   from hadoop.yarn.resource_manager_api import get_resource_manager
 
@@ -209,7 +209,7 @@ def test_yarn_configurations():
 
   try:
     url = ''
-    api = get_resource_manager()
+    api = get_resource_manager(user)
     url = api._url
     api.apps()
   except Exception, e:

+ 3 - 3
desktop/libs/hadoop/src/hadoop/yarn/node_manager_api.py

@@ -32,12 +32,12 @@ _JSON_CONTENT_TYPE = 'application/json'
 
 
 
-def get_resource_manager_api(api_url):
+def get_node_manager_api(api_url):
   yarn_cluster = cluster.get_cluster_conf_for_job_submission()
-  return ResourceManagerApi(api_url, yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
+  return NodeManagerApi(api_url, yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
 
 
-class ResourceManagerApi(object):
+class NodeManagerApi(object):
   def __init__(self, oozie_url, security_enabled=False, ssl_cert_ca_verify=True):
     self._url = posixpath.join(oozie_url, 'ws', _API_VERSION)
     self._client = HttpClient(self._url, logger=LOG)

+ 27 - 9
desktop/libs/hadoop/src/hadoop/yarn/resource_manager_api.py

@@ -22,6 +22,7 @@ import threading
 
 from django.utils.translation import ugettext as _
 
+from desktop.conf import DEFAULT_USER
 from desktop.lib.exceptions_renderable import PopupException
 from desktop.lib.rest.http_client import HttpClient
 from desktop.lib.rest.resource import Resource
@@ -30,8 +31,8 @@ from hadoop import cluster
 
 
 LOG = logging.getLogger(__name__)
-DEFAULT_USER = 'hue'
 
+DEFAULT_USER = DEFAULT_USER.get()
 _API_VERSION = 'v1'
 _JSON_CONTENT_TYPE = 'application/json'
 
@@ -39,7 +40,7 @@ _api_cache = None
 _api_cache_lock = threading.Lock()
 
 
-def get_resource_manager():
+def get_resource_manager(user):
   global _api_cache
   if _api_cache is None:
     _api_cache_lock.acquire()
@@ -48,7 +49,7 @@ def get_resource_manager():
         yarn_cluster = cluster.get_cluster_conf_for_job_submission()
         if yarn_cluster is None:
           raise PopupException(_('No Resource Manager are available.'))
-        _api_cache = ResourceManagerApi(yarn_cluster.RESOURCE_MANAGER_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
+        _api_cache = ResourceManagerApi(user, yarn_cluster.RESOURCE_MANAGER_API_URL.get(), yarn_cluster.SECURITY_ENABLED.get(), yarn_cluster.SSL_CERT_CA_VERIFY.get())
     finally:
       _api_cache_lock.release()
   return _api_cache
@@ -60,8 +61,9 @@ class YarnFailoverOccurred(Exception):
 
 class ResourceManagerApi(object):
 
-  def __init__(self, oozie_url, security_enabled=False, ssl_cert_ca_verify=False):
-    self._url = posixpath.join(oozie_url, 'ws', _API_VERSION)
+  def __init__(self, user, rm_url, security_enabled=False, ssl_cert_ca_verify=False):
+    self._user = user
+    self._url = posixpath.join(rm_url, 'ws', _API_VERSION)
     self._client = HttpClient(self._url, logger=LOG)
     self._root = Resource(self._client)
     self._security_enabled = security_enabled
@@ -71,6 +73,16 @@ class ResourceManagerApi(object):
 
     self._client.set_verify(ssl_cert_ca_verify)
 
+  def _get_params(self):
+    params = {}
+
+    if DEFAULT_USER != self._user.username: # We impersonate if needed
+      params['doAs'] = self._user.username
+      if not self.security_enabled:
+        params['user.name'] = DEFAULT_USER
+
+    return params
+
   def __str__(self):
     return "ResourceManagerApi at %s" % (self._url,)
 
@@ -83,16 +95,22 @@ class ResourceManagerApi(object):
     return self._security_enabled
 
   def cluster(self, **kwargs):
-    return self._execute(self._root.get, 'cluster/info', params=kwargs, headers={'Accept': _JSON_CONTENT_TYPE})
+    params = self._get_params()
+    params.update(kwargs)
+    return self._execute(self._root.get, 'cluster/info', params=params, headers={'Accept': _JSON_CONTENT_TYPE})
 
   def apps(self, **kwargs):
-    return self._execute(self._root.get, 'cluster/apps', params=kwargs, headers={'Accept': _JSON_CONTENT_TYPE})
+    params = self._get_params()
+    params.update(kwargs)
+    return self._execute(self._root.get, 'cluster/apps', params=params, headers={'Accept': _JSON_CONTENT_TYPE})
 
   def app(self, app_id):
-    return self._execute(self._root.get, 'cluster/apps/%(app_id)s' % {'app_id': app_id}, headers={'Accept': _JSON_CONTENT_TYPE})
+    params = self._get_params()
+    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):
-    return self._execute(self._root.put, 'cluster/apps/%(app_id)s/state' % {'app_id': app_id}, data=json.dumps({'state': 'KILLED'}), contenttype=_JSON_CONTENT_TYPE)
+    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)
 
   def _execute(self, function, *args, **kwargs):
     response = function(*args, **kwargs)

+ 0 - 1
desktop/libs/liboozie/src/liboozie/oozie_api.py

@@ -16,7 +16,6 @@
 
 import logging
 import posixpath
-import threading
 
 from desktop.conf import TIME_ZONE
 from desktop.conf import DEFAULT_USER