Browse Source

HUE-6659 [import] added hive as a new option in Rdbms data import.

Prachi Poddar 8 năm trước cách đây
mục cha
commit
1edcc00

+ 81 - 69
desktop/libs/indexer/src/indexer/api3.py

@@ -173,7 +173,22 @@ def guess_field_types(request):
 def get_databases(request):
   source = json.loads(request.POST.get('source', '{}'))
   user = User.objects.get(username=request.user)
-  query_server = rdbms.get_query_server_config(server=source['rdbmsType'])
+  if source['rdbmsMode'] == 'configRdbms':
+    query_server = rdbms.get_query_server_config(server=source['rdbmsType'])
+  else:
+    name = source['rdbmsType']
+    if name:
+      query_server = {
+        'server_name': str(name),
+        'server_host': str(source['rdbmsHostname']),
+        'server_port': int(source['rdbmsPort']),
+        'username': str(source['rdbmsUsername']),
+        'password': str(source['rdbmsPassword']),
+        'options': {},
+        'alias': name
+      }
+    LOG.debug("Query Server: %s" % query_server)
+
   db = rdbms.get(user, query_server=query_server)
   assist = Assist(db)
   data = assist.get_databases() #format of data ['abc','def','ghi',...,'xyz']
@@ -188,7 +203,9 @@ def get_databases(request):
     format_['data'] = list
     format_['status'] = 0
   else:
-    format_ = []
+    format_ = {}
+    format_['data'] = []
+    format_['status'] = 1
   print format_
   return JsonResponse(format_)
 
@@ -196,52 +213,22 @@ def get_databases(request):
 def get_tables(request):
   source = json.loads(request.POST.get('source', '{}'))
   user = User.objects.get(username=request.user)
-  query_server = rdbms.get_query_server_config(server=source['rdbmsType'])
-  db = rdbms.get(user, query_server=query_server)
-  assist = Assist(db)
-  data = assist.get_tables(source['rdbmsDatabaseName']) ##format of data ['abc','def','ghi',...,'xyz']
-  format_ = {}
-  if data:
-    list = []
-    for element in data:
-      dict = {}
-      dict['name'] = element
-      dict['value'] = element
-      list.append(dict)
-    format_['data'] = list
-    format_['status'] = 0
+  if source['rdbmsMode'] == 'configRdbms':
+    query_server = rdbms.get_query_server_config(server=source['rdbmsType'])
   else:
-    format_ = []
-  print format_
-  return JsonResponse(format_)
-
-
-def dbms_test_connection(request):
-  source = json.loads(request.POST.get('source', '{}'))
-  user = User.objects.get(username=request.user)
-  name = source['rdbmsType']
-  if name:
-    query_server = {
-      'server_name': name,
-      'server_host': source['rdbmsHostname'],
-      'server_port': int(source['rdbmsPort']),
-      'username': source['rdbmsUsername'],
-      'password': source['rdbmsPassword'],
-      'options': {},
-      'alias': name,
-      'name': name
-    }
-  LOG.debug("Query Server: %s" % query_server)
-  db = rdbms.get(user, query_server=query_server)
-  assist = Assist(db)
-  data = assist.get_databases()  # format of data ['abc','def','ghi',...,'xyz']
-  print "ABCABCABC"
-  print data
+    name = source['rdbmsType']
+    if name:
+      query_server = {
+        'server_name': str(name),
+        'server_host': str(source['rdbmsHostname']),
+        'server_port': int(source['rdbmsPort']),
+        'username': str(source['rdbmsUsername']),
+        'password': str(source['rdbmsPassword']),
+        'options': {},
+        'alias': name
+      }
+    LOG.debug("Query Server: %s" % query_server)
 
-  '''
-  source = json.loads(request.POST.get('source', '{}'))
-  user = User.objects.get(username=request.user)
-  query_server = rdbms.get_query_server_config(server=source['rdbmsType'])
   db = rdbms.get(user, query_server=query_server)
   assist = Assist(db)
   data = assist.get_tables(source['rdbmsDatabaseName']) ##format of data ['abc','def','ghi',...,'xyz']
@@ -257,11 +244,6 @@ def dbms_test_connection(request):
     format_['status'] = 0
   else:
     format_ = []
-  '''
-  format_ = {}
-  format_['data'] = 'true'
-  format_['status'] = 0
-  print "@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@@"
   print format_
   return JsonResponse(format_)
 
@@ -297,6 +279,10 @@ def importer_submit(request):
     job_handle = _create_database(request, source, destination, start_time)
   elif destination['outputFormat'] == 'file' and source['inputFormat'] == 'rdbms':
     job_handle = run_sqoop(request, source, destination, start_time)
+  elif destination['outputFormat'] == 'hive' and source['inputFormat'] == 'rdbms':
+    job_handle = run_sqoop(request, source, destination, start_time)
+  elif destination['outputFormat'] == 'hbase' and source['inputFormat'] == 'rdbms':
+    job_handle = run_sqoop(request, source, destination, start_time)
   else:
     job_handle = _create_table(request, source, destination, start_time)
   print JsonResponse(job_handle)
@@ -423,24 +409,50 @@ def _index(request, file_format, collection_name, query=None, start_time=None, l
   return indexer.run_morphline(request, collection_name, morphline, input_path, query, start_time=start_time, lib_path=lib_path)
 
 def run_sqoop(request, source, destination, start_time):
-  rdbmsName = source['rdbmsName']
-  rdbmsHost = DATABASES[rdbmsName].HOST.get()
-  rdbmsPort = DATABASES[rdbmsName].PORT.get()
-  rdbmsDatabaseName = source['rdbmsDatabaseName']
-  rdbmsTableName = source['rdbmsTableName']
-  rdbmsUserName = DATABASES[rdbmsName].USER.get()
-  rdbmsPassword = get_database_password(rdbmsName)
-  targetDir = conf.HDFS_CLUSTERS['default'].FS_DEFAULTFS.get()+destination['name']
-  print rdbmsName
-  print rdbmsHost
-  print rdbmsPort
-  print rdbmsDatabaseName
-  print rdbmsTableName
-  print rdbmsUserName
-  print rdbmsPassword
-  print targetDir
-  print 'import --connect jdbc:'+rdbmsName+'://'+'127.0.0.1'+':'+str(rdbmsPort)+'/'+rdbmsDatabaseName+' --username '+rdbmsUserName+' --password '+rdbmsPassword+' --query \'SELECT * FROM '+rdbmsTableName+' as a WHERE $CONDITIONS\' --target-dir '+targetDir+' --verbose --split-by a.empid'
+  rdbmsMode = str(source['rdbmsMode'])
+  rdbmsName = str(source['rdbmsType'])
+  rdbmsDatabaseName = str(source['rdbmsDatabaseName'])
+  allTablesSelected = str(source['allTablesSelected'])
+  destinationType = str(destination['outputFormat'])
+
+  if not allTablesSelected:
+    rdbmsTableName = str(source['rdbmsTableName'])
+
+  if rdbmsMode == 'configRdbms':
+    rdbmsHost = str(DATABASES[rdbmsName].HOST.get())
+    rdbmsPort = str(DATABASES[rdbmsName].PORT.get())
+    rdbmsUserName = str(DATABASES[rdbmsName].USER.get())
+    rdbmsPassword = str(get_database_password(rdbmsName))
+  else:
+    rdbmsHost = str(source['rdbmsHostname'])
+    rdbmsPort = str(source['rdbmsPort'])
+    rdbmsUserName = str(source['rdbmsUsername'])
+    rdbmsPassword = str(source['rdbmsPassword'])
+
+  if destinationType == 'file':
+    targetDir = conf.HDFS_CLUSTERS['default'].FS_DEFAULTFS.get()+str(destination['name'])+'/test'
+
+  #print rdbmsName
+  #print rdbmsHost
+  #print rdbmsPort
+  #print rdbmsDatabaseName
+  #print rdbmsUserName
+  #print rdbmsPassword
+  #print targetDir
+  #print 'import --connect jdbc:'+rdbmsName+'://'+'127.0.0.1'+':'+str(rdbmsPort)+'/'+rdbmsDatabaseName+' --username '+rdbmsUserName+' --password '+rdbmsPassword+' --query \'SELECT * FROM '+rdbmsTableName+' as a WHERE $CONDITIONS\' --target-dir '+targetDir+' --verbose --split-by a.empid'
+
+  if destinationType == 'file':
+    if allTablesSelected:
+      statement='import-all-tables --connect jdbc:'+rdbmsName+'://'+rdbmsHost+':'+rdbmsPort+'/'+rdbmsDatabaseName+' --username '+rdbmsUserName+' --password '+rdbmsPassword+' --warehouse-dir '+targetDir+' -m 1'
+    else:
+      statement = 'import --connect jdbc:'+rdbmsName+'://'+rdbmsHost+':'+rdbmsPort+'/'+rdbmsDatabaseName+' --username '+rdbmsUserName+' --password '+rdbmsPassword+' --table '+rdbmsTableName+' --target-dir '+ targetDir+' -m 1'
+  elif destinationType == 'hive':
+    if allTablesSelected:
+      statement = 'import-all-tables --connect jdbc:'+rdbmsName+'://'+rdbmsHost+':'+rdbmsPort+'/'+rdbmsDatabaseName+' --username '+rdbmsUserName+' --password '+rdbmsPassword+' --hive-import'
+    else:
+      statement = 'import --connect jdbc:'+rdbmsName+'://'+rdbmsHost+':'+rdbmsPort+'/'+rdbmsDatabaseName+' --username '+rdbmsUserName+' --password '+rdbmsPassword+' --table '+rdbmsTableName+' --hive-import'
 
+  print statement
   task = make_notebook(
       name=_('Indexer job for %(rdbmsDatabaseName)s.%(rdbmsDatabaseName)s to %(path)s') % {
           'rdbmsDatabaseName': source['rdbmsDatabaseName'],
@@ -448,7 +460,7 @@ def run_sqoop(request, source, destination, start_time):
           'path': destination['name']
         },
       editor_type='sqoop1',
-      statement='import --connect jdbc:mysql://172.31.114.131:3306/hue --username hue --password 12345678 --query SELECT * FROM employee WHERE $CONDITIONS --target-dir hdfs://nightly512-unsecure-1.gce.cloudera.com:8020/user/admin/test77 -m 1',
+      statement=statement,
       files = [{"path": "/user/admin/mysql-connector-java.jar", "type": "jar"}],
       status='ready',
       on_success_url='/filebrowser/view/%s(name)s' % destination,

+ 92 - 49
desktop/libs/indexer/src/indexer/templates/importer.mako

@@ -427,7 +427,7 @@ ${ assist.assistPanel() }
 
             <!-- /ko -->
 
-            <!-- ko if: createWizard.source.rdbmsMode() == 'configRdbms' || (createWizard.source.rdbmsMode() == 'customRdbms' && createWizard.source.dbmsIsValid() == 'true') -->
+            <!-- ko if: createWizard.source.rdbmsMode() == 'configRdbms' || (createWizard.source.rdbmsMode() == 'customRdbms' && createWizard.source.dbmsIsValid() == true) -->
               <!-- ko if: createWizard.source.rdbmsType -->
               <div class="control-group input-append">
                 <label for="rdbmsDatabaseName" class="control-label"><div>${ _('Database Name') }</div>
@@ -438,7 +438,7 @@ ${ assist.assistPanel() }
 
               <!-- ko if: createWizard.source.rdbmsDatabaseName -->
               <div class="control-group input-append">
-                <!-- ko if: createWizard.source.allTablesSelected() == 'false' -->
+                <!-- ko if: createWizard.source.allTablesSelected() == false -->
                 <label for="rdbmsTableName" class="control-label"><div>${ _('Table Name') }</div>
                   <select id="rdbmsTableName" data-bind="selectize: createWizard.source.rdbmsTableNames, value: createWizard.source.rdbmsTableName, optionsText: 'name', optionsValue: 'value'"></select>
                 </label>
@@ -488,6 +488,7 @@ ${ assist.assistPanel() }
       <!-- /ko -->
     </div>
 
+    <!-- ko if: createWizard.source.isAllTables() == false -->
     <div class="card step" style="min-height: 310px;">
       <!-- ko ifnot: createWizard.isGuessingFormat -->
       <!-- ko if: createWizard.isGuessingFieldTypes -->
@@ -517,6 +518,7 @@ ${ assist.assistPanel() }
       <!-- /ko -->
     </div>
     <!-- /ko -->
+    <!-- /ko -->
 
     <!-- /ko -->
 
@@ -533,8 +535,10 @@ ${ assist.assistPanel() }
             </label>
           </div>
           <div class="control-group">
+            <!-- ko if: outputFormat() != 'hive' && outputFormat() != 'hbase' && outputFormat() != '' -->
             <label for="collectionName" class="control-label "><div>${ _('Name') }</div></label>
-            <!-- ko if: outputFormat() != 'table' && outputFormat() != 'database' -->
+            <!-- /ko -->
+            <!-- ko if: outputFormat() == 'file' -->
               <input type="text" class="form-control name input-xlarge" id="collectionName" data-bind="value: name, filechooser: name, filechooserOptions: { linkMarkup: true, skipInitialPathIfEmpty: true, openOnFocus: true, selectFolder: true, uploadFile: false, uploadFolder: true}" placeholder="${ _('Name') }">
             <!-- /ko -->
 
@@ -544,7 +548,8 @@ ${ assist.assistPanel() }
             <span class="help-inline muted" data-bind="visible: !isTargetExisting() && isTargetChecking()">
               <i class="fa fa-spinner fa-spin"></i>
             </span>
-            <span class="help-inline muted" data-bind="visible: !$parent.createWizard.isValidDestination()">
+            <!-- ko if: outputFormat() != 'hive' && outputFormat() != 'hbase' && outputFormat.length > 0 -->
+            <span class="help-inline muted" data-bind="visible: ! $parent.createWizard.isValidDestination()">
               <i class="fa fa-warning" style="color: #c09853"></i> ${ _('Empty name or invalid characters') }
             </span>
             <span class="help-inline muted" data-bind="visible: isTargetExisting()">
@@ -556,6 +561,7 @@ ${ assist.assistPanel() }
               <!-- /ko -->
               <a href="javascript:void(0)" data-bind="hueLink: existingTargetUrl(), text: name" title="${ _('Open') }"></a>
             </span>
+            <!-- /ko -->
           </div>
         </div>
       </div>
@@ -1239,6 +1245,7 @@ ${ assist.assistPanel() }
         self.sample.removeAll();
         self.path('');
         resizeElements();
+        self.rdbmsMode('');
       });
       self.inputFormatsAll = ko.observableArray([
           {'value': 'file', 'name': 'File'},
@@ -1279,6 +1286,18 @@ ${ assist.assistPanel() }
       });
       // Rdbms
       self.rdbmsMode = ko.observable('');
+      self.rdbmsMode.subscribe(function (val) {
+        self.rdbmsType('');
+        self.rdbmsDatabaseName('');
+        self.rdbmsTableName('');
+        self.isAllTables(false);
+        self.rdbmsHostname('');
+        self.rdbmsPort('');
+        self.rdbmsUsername('');
+        self.rdbmsPassword('');
+        self.dbmsIsValid(false);
+        self.isConnection(false);
+      });
       self.rdbmsTypesAll = ko.observableArray([
           {'value': 'mysql', 'name': 'Mysql'},
           {'value': 'oracle', 'name': 'Oracle'},
@@ -1291,73 +1310,97 @@ ${ assist.assistPanel() }
       self.rdbmsType = ko.observable('');
       self.rdbmsType.subscribe(function (val) {
         self.path('');
+        self.isConnection(false);
         resizeElements();
-        $.post("${ url('indexer:get_databases') }", {
-          "source": ko.mapping.toJSON(self)
-        }, function (resp) {
-          if (resp.status == 0 && resp.data) {
-            self.rdbmsDatabaseNames(resp.data);
-          }
-        });
+        if(self.rdbmsMode() == 'configRdbms'){
+          $.post("${ url('indexer:get_databases') }", {
+            "source": ko.mapping.toJSON(self)
+          }, function (resp) {
+            if (resp.status == 0 && resp.data) {
+              self.rdbmsDatabaseNames(resp.data);
+            }
+          });
+        }
       });
       self.rdbmsDatabaseName = ko.observable('');
       self.rdbmsDatabaseName.subscribe(function (val) {
-        $.post("${ url('indexer:get_tables') }", {
-          "source": ko.mapping.toJSON(self)
-        }, function (resp) {
-          if (resp.status == 0 && resp.data) {
-            self.rdbmsTableNames(resp.data);
-          }
-        });
+        if(val != ''){
+          $.post("${ url('indexer:get_tables') }", {
+            "source": ko.mapping.toJSON(self)
+          }, function (resp) {
+            if (resp.status == 0 && resp.data) {
+              self.rdbmsTableNames(resp.data);
+            }
+          });
+        }
       });
       self.rdbmsDatabaseNames = ko.observableArray([]);
       self.rdbmsTableName = ko.observable('');
       self.rdbmsTableName.subscribe(function (val) {
-        wizard.guessFieldTypes();
+        if(val != ''){
+          wizard.guessFieldTypes();
+        }
       });
       self.rdbmsTableNames = ko.observableArray([]);
-      // Table
-      self.table = ko.observable('');
-      self.tableName = ko.computed(function() {
-        return self.table().indexOf('.') > 0 ? self.table().split('.', 2)[1] : self.table();
-      });
-      self.databaseName = ko.computed(function() {
-        return self.table().indexOf('.') > 0 ? self.table().split('.', 2)[0] : 'default';
-      });
-      self.table.subscribe(function(val) {
-        resizeElements();
-      });
-      self.apiHelperType = ko.observable('${ source_type }');
       self.rdbmsHostname = ko.observable('');
+      self.rdbmsHostname.subscribe(function (val) {
+        self.isConnection(false);
+      });
       self.rdbmsPort = ko.observable('');
+      self.rdbmsPort.subscribe(function (val) {
+        self.isConnection(false);
+      });
       self.rdbmsUsername = ko.observable('');
+      self.rdbmsUsername.subscribe(function (val) {
+        self.isConnection(false);
+      });
       self.rdbmsPassword = ko.observable('');
-      self.dbmsIsValid = ko.observable('');
-      self.allTablesSelected = ko.observable('false');
+      self.rdbmsPassword.subscribe(function (val) {
+        self.isConnection(false);
+      });
+      self.allTablesSelected = ko.observable(false);
+      self.isAllTables = ko.observable(false);
+      self.isAllTables.subscribe(function(newVal) {
+        self.rdbmsTableName('');
+        if(newVal){
+          self.allTablesSelected(true);
+        }else{
+          self.allTablesSelected(false);
+        }
+      });
+      self.dbmsIsValid = ko.observable(false);
       self.isConnection = ko.observable(false);
       self.isConnection.subscribe(function(newVal) {
         if(newVal){
-          $.post("${ url('indexer:dbms_test_connection') }", {
+          $.post("${ url('indexer:get_databases') }", {
             "source": ko.mapping.toJSON(self)
           }, function (resp) {
-            if (resp.status == 0 && resp.data) {
-              self.dbmsIsValid(resp.data);
+            if(resp.status == 0 && resp.data) {
+              console.log(resp.data)
+              self.dbmsIsValid(true);
+              self.rdbmsDatabaseNames(resp.data);
             }
+          }).fail(function (xhr, textStatus, errorThrown) {
+            $(document).trigger("error", "Connection Failed.");
+            self.dbmsIsValid(false);
+            self.isConnection(false);
           });
         }else{
-          self.dbmsIsValid('');
+          self.dbmsIsValid(false);
         }
       });
-      self.isAllTables = ko.observable(false);
-      self.isAllTables.subscribe(function(newVal) {
-        if(newVal){
-          self.allTablesSelected('true');
-        }else{
-          self.allTablesSelected('false');
-        }
+      // Table
+      self.table = ko.observable('');
+      self.tableName = ko.computed(function() {
+        return self.table().indexOf('.') > 0 ? self.table().split('.', 2)[1] : self.table();
       });
-
-
+      self.databaseName = ko.computed(function() {
+        return self.table().indexOf('.') > 0 ? self.table().split('.', 2)[0] : 'default';
+      });
+      self.table.subscribe(function(val) {
+        resizeElements();
+      });
+      self.apiHelperType = ko.observable('${ source_type }');
       // Queries
       self.query = ko.observable('');
       self.draggedQuery = ko.observable();
@@ -1387,7 +1430,7 @@ ${ assist.assistPanel() }
         } else if (self.inputFormat() == 'manual') {
           return true;
         } else if (self.inputFormat() == 'rdbms') {
-          return self.rdbmsDatabaseName().length > 0 && (self.rdbmsTableName().length > 0 || self.allTablesSelected() == 'true');
+          return self.rdbmsDatabaseName().length > 0 && (self.rdbmsTableName().length > 0 || self.allTablesSelected() == true);
         }
       });
       self.defaultName = ko.computed(function() {
@@ -1721,7 +1764,7 @@ ${ assist.assistPanel() }
          );
       });
       self.readyToIndex = ko.computed(function () {
-        var validFields = self.destination.columns().length || self.destination.outputFormat() == 'database' || self.destination.outputFormat() == 'file';
+        var validFields = self.destination.columns().length || self.destination.outputFormat() == 'database' || self.destination.outputFormat() == 'file' || self.destination.outputFormat() == 'hive' || self.destination.outputFormat() == 'hbase';
         var validTableColumns = self.destination.outputFormat() != 'table' || ($.grep(self.destination.columns(), function(column) {
             return column.name().length == 0;
           }).length == 0
@@ -1729,7 +1772,7 @@ ${ assist.assistPanel() }
             return column.name().length == 0 || (self.source.inputFormat() != 'manual' && column.partitionValue().length == 0);
           }).length == 0
         );
-        var isTargetAlreadyExisting = ! self.destination.isTargetExisting() || self.destination.outputFormat() == 'index' || self.destination.outputFormat() == 'file';
+        var isTargetAlreadyExisting = ! self.destination.isTargetExisting() || self.destination.outputFormat() == 'index' || self.destination.outputFormat() == 'file' || self.destination.outputFormat() == 'hive' || self.destination.outputFormat() == 'hbase';
         var isValidTable = self.destination.outputFormat() != 'table' || (
           self.destination.tableFormat() != 'kudu' || (self.destination.kuduPartitionColumns().length > 0 &&
               $.grep(self.destination.kuduPartitionColumns(), function(partition) { return partition.columns().length > 0 }).length == self.destination.kuduPartitionColumns().length && self.destination.primaryKeys().length > 0)

+ 0 - 1
desktop/libs/indexer/src/indexer/urls.py

@@ -57,7 +57,6 @@ urlpatterns += patterns('indexer.api3',
   url(r'^api/indexer/guess_field_types/$', 'guess_field_types', name='guess_field_types'),
   url(r'^api/indexer/get_databases/$', 'get_databases', name='get_databases'),
   url(r'^api/indexer/get_tables/$', 'get_tables', name='get_tables'),
-  url(r'^api/indexer/dbms_test_connection/$', 'dbms_test_connection', name='dbms_test_connection'),
 
   url(r'^api/importer/submit', 'importer_submit', name='importer_submit')
 )