Sfoglia il codice sorgente

[importer_direct_upload_hive] Supporting hive dialect for the direct upload

ayush.goyal 4 anni fa
parent
commit
4aa4b05796

+ 18 - 14
desktop/libs/indexer/src/indexer/indexers/sql.py

@@ -295,18 +295,18 @@ class SQLIndexer(object):
 
     columns = destination['columns']
 
-    sql = '''CREATE TABLE IF NOT EXISTS %(table_name)s (
-%(columns)s
-);
-''' % {
-          'database': database,
-          'table_name': table_name,
-          'columns': ',\n'.join(['  `%(name)s` %(type)s' % col for col in columns]),
-          'primary_keys': ', '.join(destination.get('indexerPrimaryKey'))
-      }
+    if editor_type == 'hive':
+      sql = '''CREATE TABLE IF NOT EXISTS %(database)s.%(table_name)s (
+%(columns)s);
+      ''' % {
+              'database': database,
+              'table_name': table_name,
+              'columns': ',\n'.join(['  `%(name)s` %(type)s' % col for col in columns]),
+            }
 
     path = urllib_unquote(source['path'])
-    if path:
+
+    if path:                                                     # data insertion
       with open(BASE_DIR + path, 'r') as local_file:
         reader = csv.reader(local_file)
         list_of_tuples = list(map(tuple, reader))
@@ -315,10 +315,14 @@ class SQLIndexer(object):
           list_of_tuples = list_of_tuples[1:]
 
         csv_rows = str(list_of_tuples)[1:-1]
-        sql += '''INSERT INTO %(table_name)s VALUES %(csv_rows)s;'''% {
-            'table_name': table_name,
-            'csv_rows': csv_rows
-          }
+
+        if editor_type == 'hive':
+          sql += '''\nINSERT INTO %(database)s.%(table_name)s VALUES %(csv_rows)s;
+          '''% {
+                  'database': database,
+                  'table_name': table_name,
+                  'csv_rows': csv_rows
+                }
 
     on_success_url = reverse('metastore:describe_table', kwargs={'database': database, 'table': final_table_name}) + \
         '?source_type=' + source_type

+ 2 - 3
desktop/libs/indexer/src/indexer/indexers/sql_tests.py

@@ -817,7 +817,7 @@ def test_create_table_from_local():
 
   statement = '''USE default;
 
-CREATE TABLE IF NOT EXISTS test1 (
+CREATE TABLE IF NOT EXISTS default.test1 (
   `date` timestamp,
   `hour` bigint,
   `minute` bigint,
@@ -831,7 +831,6 @@ CREATE TABLE IF NOT EXISTS test1 (
   `plane` string,
   `cancelled` boolean,
   `time` bigint,
-  `dist` bigint
-);'''
+  `dist` bigint);'''
 
   assert_equal(statement, sql)