Forráskód Böngészése

HUE-4467 [indexer] Add field operation morphline generation tests for each operation

Aaron Peddle 9 éve
szülő
commit
de560a39f3

+ 17 - 1
desktop/libs/indexer/src/indexer/argument.py

@@ -17,6 +17,8 @@ from django.utils.translation import ugettext as _
 
 class Argument():
   _type = None
+  _default_value = None
+
   def __init__(self, name, description=None):
     self._name = name
     self._description = _(description if description else name)
@@ -29,14 +31,28 @@ class Argument():
   def type(self):
     return self._type
 
+  @property
+  def default_value(self):
+    return self._default_value
+
+
   def to_dict(self):
     return {"name": self._name, "type": self._type, "description": self._description}
 
+  def get_default_arg_pair(self):
+    return (self.name, self.default_value)
+
 class TextArgument(Argument):
   _type = "text"
+  _default_value = ""
 
 class CheckboxArgument(Argument):
   _type = "checkbox"
+  _default_value = False
 
 class MappingArgument(Argument):
-  _type = "mapping"
+  _type = "mapping"
+
+  @property
+  def default_value(self):
+    return []

+ 7 - 4
desktop/libs/indexer/src/indexer/fields.py

@@ -39,16 +39,16 @@ class FieldType():
     return pattern.match(field)
 
 class Field(object):
-  def __init__(self, name, field_type, operations=[]):
+  def __init__(self, name, field_type_name, operations=None):
     self.name = name
-    self.field_type = field_type
+    self.field_type_name = field_type_name
     self.keep = True
-    self.operations = operations
+    self.operations = operations if operations else []
     self.required = False
 
   def to_dict(self):
     return {'name': self.name,
-    'type': self.field_type,
+    'type': self.field_type_name,
     'keep': self.keep,
     'operations': [operation.to_dict() for operation in self.operations],
     'required': self.required}
@@ -61,6 +61,9 @@ FIELD_TYPES = [
   FieldType('date', "^([0-9]+-[0-9]+-[0-9]+T[0-9]+:[0-9]+:[0-9]+(\\.[0-9]*)?Z)?$")
 ]
 
+def get_field_type(type_name):
+  return [file_type for file_type in FIELD_TYPES if file_type.name == type_name][0]
+
 def guess_field_type_from_samples(samples):
   guesses = [_guess_field_type(sample) for sample in samples]
 

+ 12 - 2
desktop/libs/indexer/src/indexer/operations.py

@@ -30,6 +30,9 @@ class Operator():
   def args(self):
     return self._args
 
+  def _get_default_output_fields(self):
+    return []
+
   def to_dict(self):
     return {
       "name": self._name,
@@ -37,6 +40,13 @@ class Operator():
       "outputType": self._output_type
     }
 
+  def get_default_operation(self):
+    return {
+      "type": self._name,
+      "settings": dict([arg.get_default_arg_pair() for arg in self._args]),
+      "fields": self._get_default_output_fields()
+    }
+
 OPERATORS = [
   Operator(
     name="split",
@@ -106,11 +116,11 @@ OPERATORS = [
   ),
 ]
 
-def _get_operator(operation_name):
+def get_operator(operation_name):
   return [operation for operation in OPERATORS if operation.name == operation_name][0]
 
 def get_checked_args(operation):
-  operation_args = _get_operator(operation["type"]).args
+  operation_args = get_operator(operation["type"]).args
 
   kept_args = [arg for arg in operation_args if operation['settings'][arg.name]]
 

+ 4 - 4
desktop/libs/indexer/src/indexer/smart_indexer.py

@@ -23,7 +23,7 @@ from liboozie.oozie_api import get_oozie
 from oozie.models2 import Job
 from liboozie.submission2 import Submission
 
-from indexer.fields import Field, FIELD_TYPES
+from indexer.fields import Field, FIELD_TYPES, get_field_type
 from indexer.operations import get_checked_args
 from indexer.file_format import get_file_format_instance, get_file_format_class
 from indexer.conf import CONFIG_INDEXING_TEMPLATES_PATH
@@ -144,10 +144,10 @@ class Indexer(object):
     return base_name
 
   @staticmethod
-  def _get_regex_for_type(type_):
-    matches = filter(lambda field_type: field_type.name == type_, FIELD_TYPES)
+  def _get_regex_for_type(type_name):
+    field_type = get_field_type(type_name)
 
-    return matches[0].regex.replace('\\', '\\\\')
+    return field_type.regex.replace('\\', '\\\\')
 
   def generate_morphline_config(self, collection_name, data, uuid_name="__uuid"):
     """

+ 75 - 1
desktop/libs/indexer/src/indexer/tests_indexer.py

@@ -24,8 +24,9 @@ from hadoop.pseudo_hdfs4 import is_live_cluster
 
 from indexer.smart_indexer import Indexer
 from indexer.controller import CollectionManagerController
-
+from indexer.operations import get_operator
 from indexer.file_format import ApacheCombinedFormat, RubyLogFormat, HueLogFormat
+from indexer.fields import Field, get_field_type
 
 LOG = logging.getLogger(__name__)
 
@@ -40,6 +41,17 @@ def _test_fixed_type_format_generate_morphline(format_):
 
   assert_true(isinstance(morphline, basestring))
 
+def _test_generate_field_operation_morphline(operation_format):
+  fields = IndexerTest.simpleCSVFields[:]
+  fields[0]['operations'].append(operation_format)
+
+  indexer = Indexer("test", None)
+  morphline =indexer.generate_morphline_config("test_collection", {
+      "columns": fields,
+      "format": IndexerTest.simpleCSVFormat
+    })
+
+  assert_true(isinstance(morphline, basestring))
 
 class IndexerTest():
   simpleCSVString = """id,Rating,Location,Name,Time
@@ -156,6 +168,68 @@ class IndexerTest():
   def test_generate_hue_log_morphline(self):
     _test_fixed_type_format_generate_morphline(HueLogFormat)
 
+  def test_generate_split_operation_morphline(self):
+    split_dict = get_operator('split').get_default_operation()
+
+    split_dict['fields'] = [
+        Field("test_field_1", "string").to_dict(),
+        Field("test_field_2", "string").to_dict()
+      ]
+
+    _test_generate_field_operation_morphline(split_dict)
+
+  def test_generate_extract_uri_components_operation_morphline(self):
+    extract_uri_dict = get_operator('extract_uri_components').get_default_operation()
+
+    extract_uri_dict['fields'] = [
+        Field("test_field_1", "string").to_dict(),
+        Field("test_field_2", "string").to_dict()
+      ]
+
+    _test_generate_field_operation_morphline(extract_uri_dict)
+
+  def test_generate_grok_operation_morphline(self):
+    grok_dict = get_operator('grok').get_default_operation()
+
+    grok_dict['fields'] = [
+        Field("test_field_1", "string").to_dict(),
+        Field("test_field_2", "string").to_dict()
+      ]
+
+    _test_generate_field_operation_morphline(grok_dict)
+
+  def test_generate_convert_date_morphline(self):
+    convert_date_dict = get_operator('convert_date').get_default_operation()
+
+    _test_generate_field_operation_morphline(convert_date_dict)
+
+  def test_generate_geo_ip_morphline(self):
+    geo_ip_dict = get_operator('geo_ip').get_default_operation()
+
+    geo_ip_dict['fields'] = [
+        Field("test_field_1", "string").to_dict(),
+        Field("test_field_2", "string").to_dict()
+      ]
+
+    _test_generate_field_operation_morphline(geo_ip_dict)
+
+  def test_generate_translate_morphline(self):
+    translate_dict = get_operator('translate').get_default_operation()
+
+    translate_dict['fields'] = [
+      Field("test_field_1", "string").to_dict(),
+      Field("test_field_2", "string").to_dict()
+    ]
+
+    translate_dict['settings']['mapping'].append({"key":"key","value":"value"})
+
+    _test_generate_field_operation_morphline(translate_dict)
+
+  def test_generate_find_replace_morphline(self):
+    find_replace_dict = get_operator('find_replace').get_default_operation()
+
+    _test_generate_field_operation_morphline(find_replace_dict)
+
   def test_end_to_end(self):
     if not is_live_cluster():
       raise SkipTest()