From acb5d550f497dcf344a888a689994b357c5bb113 Mon Sep 17 00:00:00 2001 From: Cloud User Date: Sat, 20 Aug 2016 01:21:45 -0400 Subject: [PATCH 1/3] Checking in the changes to the code --- resource/dava.properties | 7 ++- resource/dava.sh | 8 +-- resource/dava.spark.properties | 26 +++++++++ resource/electr_prod.json | 1 - .../org/chombo/spark/etl/DataValidator.scala | 54 +++++++++++-------- .../chombo/validator/ValidatorFactory.java | 8 ++- 6 files changed, 72 insertions(+), 32 deletions(-) create mode 100644 resource/dava.spark.properties diff --git a/resource/dava.properties b/resource/dava.properties index cc617fb6..b03cfc84 100644 --- a/resource/dava.properties +++ b/resource/dava.properties @@ -7,13 +7,12 @@ mapreduce.reduce.maxattempts=2 #ValidationChecker vac.filter.invalid.records=false -vac.invalid.data.file.path=/user/pranab/output/dava/invalid.txt -vac.schema.file.path=/user/pranab/meta/dava/electr_prod.json +vac.invalid.data.file.path=/usr/avi/chombo/invalid.txt +vac.validation.schema.file.path=/usr/avi/chombo/meta/dava/electr_prod.json #vac.cleanser.schema.file.path= +vac.validator.0=notMissing vac.validator.1=membership,notMissing vac.validator.2=membership,notMissing vac.validator.3=exactLength,notMissing vac.validator.4=min,max,notMissing vac.validator.5=min,max,notMissing - - diff --git a/resource/dava.sh b/resource/dava.sh index 3b3c5176..9feed306 100755 --- a/resource/dava.sh +++ b/resource/dava.sh @@ -2,12 +2,12 @@ JAR_NAME=/home/pranab/Projects/chombo/target/chombo-1.0.jar CLASS_NAME=org.chombo.mr.ValidationChecker echo "running mr" -IN_PATH=/user/pranab/dava/input -OUT_PATH=/user/pranab/dava/output +IN_PATH=/usr/avi/chombo/input/dava +OUT_PATH=/usr/avi/chombo/output/dava echo "input $IN_PATH output $OUT_PATH" hadoop fs -rmr $OUT_PATH echo "removed output dir" -hadoop fs -rm /user/pranab/output/dava/* +hadoop fs -rm /usr/avi/chombo/output/dava/* echo "removed invalid data file" -hadoop jar $JAR_NAME $CLASS_NAME -Dconf.path=/home/pranab/Projects/bin/chombo/dava.properties $IN_PATH $OUT_PATH +hadoop jar $JAR_NAME $CLASS_NAME -Dconf.path=/home/ec2-user/pranab/chombo/resource/dava.properties $IN_PATH $OUT_PATH diff --git a/resource/dava.spark.properties b/resource/dava.spark.properties new file mode 100644 index 00000000..fda640e4 --- /dev/null +++ b/resource/dava.spark.properties @@ -0,0 +1,26 @@ +system.master=local[*] +#system.master=spark://172.31.8.69:7077 +app.filter.invalid.records=true +app.output.invalid.records=true +app.invalid.records.output.file=hdfs://localhost:9000/usr/avi/chombo/invalid.txt +app.field.delim.in=, +app.field.delim.out=, +app.val.tag.separator=, +app.schema.file.path=/usr/avi/chombo/meta/dava/electr_prod.json +app.config.file.path=/usr/avi/chombo/meta/dava/electr_prod.conf + +field.delim.regex=, +debug.on=true +num.reducer=1 +mapreduce.map.maxattempts=2 +mapreduce.reduce.maxattempts=2 + +#ValidationChecker +#vac.cleanser.schema.file.path= +app.validator.0=notMissing +app.validator.1=membership,notMissing +app.validator.2=membership,notMissing +app.validator.3=exactLength,notMissing +app.validator.4=min,max,notMissing +app.validator.5=min,max,notMissing + diff --git a/resource/electr_prod.json b/resource/electr_prod.json index 1c1f078e..6067268c 100644 --- a/resource/electr_prod.json +++ b/resource/electr_prod.json @@ -1,5 +1,4 @@ { - "name" : "electronicProduct", "attributes" : [ { diff --git a/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala b/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala index 187254da..14d59630 100644 --- a/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala +++ b/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala @@ -23,6 +23,7 @@ import org.apache.spark.SparkContext import org.chombo.util.Utility import org.chombo.validator.ValidatorFactory import com.typesafe.config.Config +import com.typesafe.config.ConfigValue import org.chombo.validator.Validator import org.chombo.util.ProcessorAttributeSchema import org.chombo.util.NumericalAttrStatsManager @@ -47,16 +48,14 @@ object DataValidator extends JobConfiguration { * @return */ def main(args: Array[String]) { - val Array(master: String, inputPath: String, outputPath: String, configFile: String) = getCommandLineArgs(args, 3) + val Array( inputPath: String, outputPath: String, configFile: String) = getCommandLineArgs(args, 3) val config = createConfig(configFile) val sparkConf = createSparkConf("app.data validation", config, false) val sparkCntxt = new SparkContext(sparkConf) if (config.hasPath("app.invalid.records.output.file")) config.getString("app.invalid.records.output.file") - else - "" - + val fieldDelimIn = config.getString("app.field.delim.in") val fieldDelimOut = config.getString("app.field.delim.out") val valTagSeparator = config.getString("app.val.tag.separator") @@ -66,26 +65,36 @@ object DataValidator extends JobConfiguration { if (config.hasPath("app.invalid.records.output.file")) config.getString("app.invalid.records.output.file") else - "" - val validationSchema = Utility.getProcessingSchema( config.getString("app.schema.file.path")) + "" + val validationSchema = Utility.getProcessingSchema( config.getString("app.schema.file.path")) + val validatorConfig = config.atPath("app") - ValidatorFactory.initialize( config.getString( "app,custom.valid.factory.class"), validatorConfig ) + + val configClass = + if (config.hasPath("app.custom.valid.factory.class")) + config.getString("app.custom.valid.factory.class") + else + null + ValidatorFactory.initialize(configClass, validatorConfig ) val ordinals = validationSchema.getAttributeOrdinals() - val tagSep = config.getString( "app,vaidator.tag.separator") + val tagSep = config.getString( "app.val.tag.separator") //initialize stats manager - getAttributeStats(config.getString("app.stats.file.path")) - getAttributeMeds(config.getString("app.med.stats.file.path"), config.getString("app.mad.stats.file.path"), - Utility.intArrayFromString(config.getString("app.id.ordinals"), ",") ) + if(config.hasPath("app.stats.file.path")) + getAttributeStats(config.getString("app.stats.file.path")) + if(config.hasPath("app.med.stats.file.path")) + getAttributeMeds(config.getString("app.med.stats.file.path"), config.getString("app.mad.stats.file.path"), + Utility.intArrayFromString(config.getString("app.id.ordinals"), ",") ) //simple validators var foundSimpleValidators = false + ordinals.foreach(ord => { val key = "app.validator." + ord if (config.hasPath(key)) { - val validatorTag = config.getString(key) - val valTags = validatorTag.split(tagSep); + val validatorTag : String = config.getString(key) + val valTags :Array[String] = validatorTag.split(tagSep); createValidators(config, valTags, ord, validationSchema, mutValidators) foundSimpleValidators = true } @@ -110,12 +119,12 @@ object DataValidator extends JobConfiguration { //apply all validators for the field val taggedItems = itemsZipped.map(z => { - val valList = validators.get(z._2).get + val valList : Array[Validator] = validators.get(z._2).get val valStatuses = valList.map(validator => { val status = validator.isValid(z._1) (validator.getTag(), status) }) - + //only failed validators val failedValidators = valStatuses.filter(s => { !s._2 @@ -124,17 +133,18 @@ object DataValidator extends JobConfiguration { val field = if (failedValidators.isEmpty) z._1 else - z._1 + valTagSeparator + failedValidators.mkString(fieldDelimOut) + z._1 + ":" + failedValidators.mkString(fieldDelimOut) field }) - taggedItems.mkString(fieldDelimOut) }) taggedData.cache //filter valid data - val validData = taggedData.filter(line => !line.contains(valTagSeparator)) + val validData = taggedData.filter(line => !line.contains(":")) + val list = validData.collect() + list.foreach(println) validData.saveAsTextFile(outputPath) //filter invalid data @@ -154,9 +164,9 @@ object DataValidator extends JobConfiguration { private def createValidators( config : Config , valTags : Array[String], ord : Int, validationSchema : ProcessorAttributeSchema, mutValidators : scala.collection.mutable.HashMap[Int, Array[Validator]]) { val validatorList = List[Validator]() - val prAttr = validationSchema.findAttributeByOrdinal(ord) - val validatorConfig = config.atPath("app") - val validators = valTags.map(tag => { + val prAttr = validationSchema.findAttributeByOrdinal(ord) + val validatorConfig = config + val validators = valTags.filter(_.length > 0 ).map(tag => { val validator = tag match { case "zscoreBasedRange" => { getAttributeStats(config.getString("app.stats.file.path")) @@ -206,4 +216,4 @@ private def createValidators( config : Config , valTags : Array[String], ord validationContext.clear() validationContext.put("stats", medStatManager.get) } -} \ No newline at end of file +} diff --git a/src/main/java/org/chombo/validator/ValidatorFactory.java b/src/main/java/org/chombo/validator/ValidatorFactory.java index 69d56fcc..506cb3a9 100644 --- a/src/main/java/org/chombo/validator/ValidatorFactory.java +++ b/src/main/java/org/chombo/validator/ValidatorFactory.java @@ -171,7 +171,7 @@ public static Validator create(String validatorType, ProcessorAttribute prAttr, validator = new NumericalValidator.StatsBasedRangeValidator(validatorType, prAttr, validatorContext); } else if (validatorType.equals( ROBUST_ZCORE_BASED_RANGE_VALIDATOR)) { validator = new NumericalValidator.RobustZscoreBasedRangeValidator(validatorType, prAttr, validatorContext); - } else { + } else if (null != valConfig){ //custom validator with configured validator class names validator = createCustomValidator(validatorType, prAttr, valConfig); @@ -221,6 +221,12 @@ private static Validator createCustomValidator(String validatorType, ProcessorA * @return */ public static Config getValidatorConfig(Config transformerConfig ,String validatorTag, ProcessorAttribute prAttr) { + if(null == transformerConfig) + return null; + + if(!transformerConfig.hasPath("validators." + validatorTag)) + return null; + Config valConfig = transformerConfig.getConfig("validators." + validatorTag); Config config = null; try { From 3fe652fa420f2708595549b65424c5d07c041cf8 Mon Sep 17 00:00:00 2001 From: Cloud User Date: Sat, 20 Aug 2016 01:47:43 -0400 Subject: [PATCH 2/3] Changes ... --- pom.xml | 16 ++++++++ resource/dava.spark.properties | 28 ++++++------- .../org/chombo/spark/etl/DataValidator.scala | 39 ++++++++++--------- 3 files changed, 50 insertions(+), 33 deletions(-) diff --git a/pom.xml b/pom.xml index 52fcf8a4..b7bc6874 100644 --- a/pom.xml +++ b/pom.xml @@ -54,6 +54,22 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/ma + + org.apache.maven.plugins + maven-shade-plugin + + + package + + shade + + + + + uber-${artifactId}-${version} + + + diff --git a/resource/dava.spark.properties b/resource/dava.spark.properties index fda640e4..4ca18dfc 100644 --- a/resource/dava.spark.properties +++ b/resource/dava.spark.properties @@ -1,13 +1,13 @@ system.master=local[*] #system.master=spark://172.31.8.69:7077 -app.filter.invalid.records=true -app.output.invalid.records=true -app.invalid.records.output.file=hdfs://localhost:9000/usr/avi/chombo/invalid.txt -app.field.delim.in=, -app.field.delim.out=, -app.val.tag.separator=, -app.schema.file.path=/usr/avi/chombo/meta/dava/electr_prod.json -app.config.file.path=/usr/avi/chombo/meta/dava/electr_prod.conf +filter.invalid.records=true +output.invalid.records=true +invalid.records.output.file=hdfs://localhost:9000/usr/avi/chombo/invalid.txt +field.delim.in=, +field.delim.out=, +val.tag.separator=, +schema.file.path=/usr/avi/chombo/meta/dava/electr_prod.json +config.file.path=/usr/avi/chombo/meta/dava/electr_prod.conf field.delim.regex=, debug.on=true @@ -17,10 +17,10 @@ mapreduce.reduce.maxattempts=2 #ValidationChecker #vac.cleanser.schema.file.path= -app.validator.0=notMissing -app.validator.1=membership,notMissing -app.validator.2=membership,notMissing -app.validator.3=exactLength,notMissing -app.validator.4=min,max,notMissing -app.validator.5=min,max,notMissing +validator.0=notMissing +validator.1=membership,notMissing +validator.2=membership,notMissing +validator.3=exactLength,notMissing +validator.4=min,max,notMissing +validator.5=min,max,notMissing diff --git a/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala b/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala index 14d59630..874739cf 100644 --- a/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala +++ b/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala @@ -50,41 +50,42 @@ object DataValidator extends JobConfiguration { def main(args: Array[String]) { val Array( inputPath: String, outputPath: String, configFile: String) = getCommandLineArgs(args, 3) val config = createConfig(configFile) + val localConfig = config.atPath("app") val sparkConf = createSparkConf("app.data validation", config, false) val sparkCntxt = new SparkContext(sparkConf) - if (config.hasPath("app.invalid.records.output.file")) - config.getString("app.invalid.records.output.file") + if (localConfig.hasPath("app.invalid.records.output.file")) + localConfig.getString("app.invalid.records.output.file") - val fieldDelimIn = config.getString("app.field.delim.in") - val fieldDelimOut = config.getString("app.field.delim.out") - val valTagSeparator = config.getString("app.val.tag.separator") - val filterInvalidRecords = config.getBoolean("app.filter.invalid.records") - val outputInvalidRecords = config.getBoolean("app.output.invalid.records") + val fieldDelimIn = localConfig.getString("app.field.delim.in") + val fieldDelimOut = localConfig.getString("app.field.delim.out") + val valTagSeparator = localConfig.getString("app.val.tag.separator") + val filterInvalidRecords = localConfig.getBoolean("app.filter.invalid.records") + val outputInvalidRecords = localConfig.getBoolean("app.output.invalid.records") val invalidRecordsOutputFile = - if (config.hasPath("app.invalid.records.output.file")) - config.getString("app.invalid.records.output.file") + if (localConfig.hasPath("app.invalid.records.output.file")) + localConfig.getString("app.invalid.records.output.file") else "" - val validationSchema = Utility.getProcessingSchema( config.getString("app.schema.file.path")) + val validationSchema = Utility.getProcessingSchema( localConfig.getString("app.schema.file.path")) val validatorConfig = config.atPath("app") val configClass = - if (config.hasPath("app.custom.valid.factory.class")) - config.getString("app.custom.valid.factory.class") + if (localConfig.hasPath("app.custom.valid.factory.class")) + localConfig.getString("app.custom.valid.factory.class") else null ValidatorFactory.initialize(configClass, validatorConfig ) val ordinals = validationSchema.getAttributeOrdinals() - val tagSep = config.getString( "app.val.tag.separator") + val tagSep = localConfig.getString( "app.val.tag.separator") //initialize stats manager - if(config.hasPath("app.stats.file.path")) - getAttributeStats(config.getString("app.stats.file.path")) - if(config.hasPath("app.med.stats.file.path")) - getAttributeMeds(config.getString("app.med.stats.file.path"), config.getString("app.mad.stats.file.path"), - Utility.intArrayFromString(config.getString("app.id.ordinals"), ",") ) + if(localConfig.hasPath("app.stats.file.path")) + getAttributeStats(localConfig.getString("app.stats.file.path")) + if(localConfig.hasPath("app.med.stats.file.path")) + getAttributeMeds(localConfig.getString("app.med.stats.file.path"), localConfig.getString("app.mad.stats.file.path"), + Utility.intArrayFromString(localConfig.getString("app.id.ordinals"), ",") ) //simple validators @@ -93,7 +94,7 @@ object DataValidator extends JobConfiguration { ordinals.foreach(ord => { val key = "app.validator." + ord if (config.hasPath(key)) { - val validatorTag : String = config.getString(key) + val validatorTag : String = localConfig.getString(key) val valTags :Array[String] = validatorTag.split(tagSep); createValidators(config, valTags, ord, validationSchema, mutValidators) foundSimpleValidators = true From 4d9bb4454d8095906819b4669fb6747632a3e53e Mon Sep 17 00:00:00 2001 From: Cloud User Date: Sat, 20 Aug 2016 03:10:41 -0400 Subject: [PATCH 3/3] FInal changes --- resource/dava.spark.properties | 4 ++-- .../src/main/scala/org/chombo/spark/etl/DataValidator.scala | 5 +++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/resource/dava.spark.properties b/resource/dava.spark.properties index 4ca18dfc..254f7cdd 100644 --- a/resource/dava.spark.properties +++ b/resource/dava.spark.properties @@ -1,5 +1,5 @@ -system.master=local[*] -#system.master=spark://172.31.8.69:7077 +#system.master=local[*] +system.master=spark://172.31.8.69:7077 filter.invalid.records=true output.invalid.records=true invalid.records.output.file=hdfs://localhost:9000/usr/avi/chombo/invalid.txt diff --git a/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala b/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala index 874739cf..10cd61e9 100644 --- a/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala +++ b/spark/src/main/scala/org/chombo/spark/etl/DataValidator.scala @@ -93,7 +93,7 @@ object DataValidator extends JobConfiguration { ordinals.foreach(ord => { val key = "app.validator." + ord - if (config.hasPath(key)) { + if (localConfig.hasPath(key)) { val validatorTag : String = localConfig.getString(key) val valTags :Array[String] = validatorTag.split(tagSep); createValidators(config, valTags, ord, validationSchema, mutValidators) @@ -120,6 +120,7 @@ object DataValidator extends JobConfiguration { //apply all validators for the field val taggedItems = itemsZipped.map(z => { + println("The value of z is " + z) val valList : Array[Validator] = validators.get(z._2).get val valStatuses = valList.map(validator => { val status = validator.isValid(z._1) @@ -134,7 +135,7 @@ object DataValidator extends JobConfiguration { val field = if (failedValidators.isEmpty) z._1 else - z._1 + ":" + failedValidators.mkString(fieldDelimOut) + z._1 + valTagSeparator + failedValidators.mkString(fieldDelimOut) field })