瀏覽代碼

HUE-8758 [connectors] Plugged-in HDFS connector when activated

Romain 6 年之前
父節點
當前提交
d5c19346d4

+ 11 - 5
desktop/core/src/desktop/lib/connectors/api.py

@@ -69,7 +69,7 @@ CONNECTOR_CLASSES += [
   {'nice_name': "Pig", 'dialect': 'pig', 'settings': [], 'category': 'editor', 'description': '', 'properties': {}},
   {'nice_name': "Java", 'dialect': 'java', 'settings': [], 'category': 'editor', 'description': '', 'properties': {}},
 
-  {'nice_name': "HDFS", 'dialect': 'hdfs', 'settings': [{'name': 'server_host', 'value': ''}], 'category': 'browsers', 'description': '', 'properties': {}},
+  {'nice_name': "HDFS", 'dialect': 'hdfs', 'interface': 'rest', 'settings': [{'name': 'server_url', 'value': 'http://localhost:9870/webhdfs/v1'}, {'name': 'default_fs', 'value': 'fs_defaultfs=hdfs://localhost:8020'}], 'category': 'browsers', 'description': '', 'properties': {}},
   {'nice_name': "YARN", 'dialect': 'yarn', 'settings': [], 'category': 'browsers', 'description': '', 'properties': {}},
   {'nice_name': "S3", 'dialect': 's3', 'settings': [], 'category': 'browsers', 'description': '', 'properties': {}},
   {'nice_name': "ADLS", 'dialect': 'adls-v1', 'settings': [], 'category': 'browsers', 'description': '', 'properties': {}},
@@ -195,7 +195,7 @@ def delete_connector(request):
     raise PopupException(_('No connector with the name %(name)s found.') % connector)
 
 
-def _get_installed_connectors(category=None):
+def _get_installed_connectors(category=None, dialect=None, interface=None):
   global CONNECTOR_INSTANCES
   global CONNECTOR_IDS
   config_connectors = CONNECTORS.get()
@@ -222,10 +222,16 @@ def _get_installed_connectors(category=None):
       CONNECTOR_INSTANCES.append(connector)
       CONNECTOR_IDS += 1
 
+  connectors = CONNECTOR_INSTANCES
+
   if category is not None:
-    return [connector for connector in CONNECTOR_INSTANCES if category == connector['category']]
-  else:
-    return CONNECTOR_INSTANCES
+    connectors = [connector for connector in connectors if category == connector['category']]
+  if dialect is not None:
+    connectors = [connector for connector in connectors if dialect == connector['dialect']]
+  if interface is not None:
+    connectors = [connector for connector in connectors if interface == connector['interface']]
+
+  return connectors
 
 
 def _get_connector_by_id(id):

+ 1 - 2
desktop/core/src/desktop/middleware.py

@@ -124,8 +124,7 @@ class ClusterMiddleware(object):
   """
   def process_view(self, request, view_func, view_args, view_kwargs):
     """
-    Sets request.fs and request.jt on every request to point to the
-    configured filesystem.
+    Sets request.fs and request.jt on every request to point to the configured filesystem.
     """
     request.fs_ref = request.GET.get('fs', view_kwargs.get('fs', 'default'))
     if "fs" in view_kwargs:

+ 13 - 4
desktop/libs/hadoop/src/hadoop/cluster.py

@@ -23,7 +23,8 @@ from django.utils.functional import wraps
 from hadoop import conf
 from hadoop.fs import webhdfs, LocalSubFileSystem
 
-from desktop.conf import DEFAULT_USER
+from desktop.conf import DEFAULT_USER, has_connectors
+from desktop.lib.connectors.api import _get_installed_connectors
 from desktop.lib.paths import get_build_dir
 
 
@@ -61,7 +62,7 @@ def rm_ha(funct):
 def get_hdfs(identifier="default", user=None):
   global FS_CACHE
   get_all_hdfs()
-  return FS_CACHE[identifier]
+  return FS_CACHE[FS_CACHE.keys()[0]] if has_connectors() else FS_CACHE[identifier]
 
 
 def get_defaultfs():
@@ -79,8 +80,16 @@ def get_all_hdfs():
     return FS_CACHE
 
   FS_CACHE = {}
-  for identifier in list(conf.HDFS_CLUSTERS.keys()):
-    FS_CACHE[identifier] = _make_filesystem(identifier)
+  if has_connectors():
+    for connector in _get_installed_connectors(category='browsers', dialect='hdfs', interface='rest'):
+      settings = {setting['name']: setting['value'] for setting in connector['settings']}
+      FS_CACHE[connector['name']] = webhdfs.WebHdfs(
+        url=settings['server_url'],
+        fs_defaultfs=settings['default_fs']
+      )
+  else:
+    for identifier in list(conf.HDFS_CLUSTERS.keys()):
+      FS_CACHE[identifier] = _make_filesystem(identifier)
   return FS_CACHE
 
 

+ 30 - 26
desktop/libs/hadoop/src/hadoop/fs/webhdfs.py

@@ -34,6 +34,10 @@ import urllib.request, urllib.error
 
 from django.utils.encoding import smart_str
 from django.utils.translation import ugettext as _
+
+import hadoop.conf
+import desktop.conf
+
 from desktop.lib.rest import http_client, resource
 from past.builtins import long
 from hadoop.fs import normpath as fs_normpath, SEEK_SET, SEEK_CUR, SEEK_END
@@ -42,10 +46,6 @@ from hadoop.fs.exceptions import WebHdfsException
 from hadoop.fs.webhdfs_types import WebHdfsStat, WebHdfsContentSummary
 from hadoop.hdfs_site import get_nn_sentry_prefixes, get_umask_mode, get_supergroup, get_webhdfs_ssl
 
-
-import hadoop.conf
-import desktop.conf
-
 if sys.version_info[0] > 2:
   from urllib.parse import unquote as urllib_quote
   from urllib.parse import urlparse
@@ -53,10 +53,11 @@ else:
   from urllib import unquote as urllib_quote
   from urlparse import urlparse
 
+
 DEFAULT_HDFS_SUPERUSER = desktop.conf.DEFAULT_HDFS_SUPERUSER.get()
 
 # The number of bytes to read if not specified
-DEFAULT_READ_SIZE = 1024*1024 # 1MB
+DEFAULT_READ_SIZE = 1024 * 1024 # 1MB
 
 LOG = logging.getLogger(__name__)
 
@@ -65,18 +66,20 @@ class WebHdfs(Hdfs):
   """
   WebHdfs implements the filesystem interface via the WebHDFS rest protocol.
   """
-  DEFAULT_USER = desktop.conf.DEFAULT_USER.get()        # This should be the user running Hue
+  DEFAULT_USER = desktop.conf.DEFAULT_USER.get()
   TRASH_CURRENT = 'Current'
 
-  def __init__(self, url,
-               fs_defaultfs,
-               logical_name=None,
-               hdfs_superuser=None,
-               security_enabled=False,
-               ssl_cert_ca_verify=True,
-               temp_dir="/tmp",
-               umask=0o1022,
-               hdfs_supergroup=None):
+  def __init__(
+      self,
+      url,
+      fs_defaultfs,
+      logical_name=None,
+      hdfs_superuser=None,
+      security_enabled=False,
+      ssl_cert_ca_verify=True,
+      temp_dir="/tmp",
+      umask=01022,
+      hdfs_supergroup=None):
     self._url = url
     self._superuser = hdfs_superuser
     self._security_enabled = security_enabled
@@ -103,14 +106,16 @@ class WebHdfs(Hdfs):
   def from_config(cls, hdfs_config):
     fs_defaultfs = hdfs_config.FS_DEFAULTFS.get()
 
-    return cls(url=_get_service_url(hdfs_config),
-               fs_defaultfs=fs_defaultfs,
-               logical_name=hdfs_config.LOGICAL_NAME.get(),
-               security_enabled=hdfs_config.SECURITY_ENABLED.get(),
-               ssl_cert_ca_verify=hdfs_config.SSL_CERT_CA_VERIFY.get(),
-               temp_dir=hdfs_config.TEMP_DIR.get(),
-               umask=get_umask_mode(),
-               hdfs_supergroup=get_supergroup())
+    return cls(
+        url=_get_service_url(hdfs_config),
+        fs_defaultfs=fs_defaultfs,
+        logical_name=hdfs_config.LOGICAL_NAME.get(),
+        security_enabled=hdfs_config.SECURITY_ENABLED.get(),
+        ssl_cert_ca_verify=hdfs_config.SSL_CERT_CA_VERIFY.get(),
+        temp_dir=hdfs_config.TEMP_DIR.get(),
+        umask=get_umask_mode(),
+        hdfs_supergroup=get_supergroup()
+    )
 
   def __str__(self):
     return "WebHdfs at %s" % self._url
@@ -244,7 +249,7 @@ class WebHdfs(Hdfs):
     path = self.normpath(path)
     if not self._is_remote:
       return path
-  
+
     split = urlparse(path)
     if not split.netloc:
       path = split._replace(netloc=self._netloc).geturl()
@@ -1029,8 +1034,7 @@ def _get_service_url(hdfs_config):
 
 def test_fs_configuration(fs_config):
   """
-  This is a config validation method. Returns a list of
-    [ (config_variable, error_message) ]
+  This is a config validation method. Returns a list of [(config_variable, error_message)].
   """
   fs = WebHdfs.from_config(fs_config)
   fs.setuser(fs.superuser)

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

@@ -19,7 +19,6 @@ from collections import OrderedDict
 
 from django.utils.translation import ugettext_lazy as _t
 
-
 from desktop import appmanager
 from desktop.conf import is_oozie_enabled, has_connectors
 from desktop.lib.conf import Config, UnspecifiedConfigSection, ConfigSection, coerce_json_dict, coerce_bool, coerce_csv