Created
June 17, 2019 13:56
-
-
Save rajanim/3f95c4de23d35232ff9c11b1678afb61 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| import org.apache.spark.sql.functions._ | |
| import org.apache.spark.sql.SaveMode | |
| import org.apache.hadoop.fs._ | |
| import org.apache.spark.scheduler.{SparkListener, SparkListenerTaskEnd} | |
| import java.io._ | |
| import spark.implicits._ | |
| import java.util.Calendar | |
| import java.text.SimpleDateFormat | |
| var recordsWrittenCount = 0L | |
| object SolrToParquet { | |
| /** | |
| * Export Solr documents to a Parquet file | |
| */ | |
| def write(collection:String, query:String, fields:String, filters:String, ds:String, repartitionSize:String, bucketName:String, iamRole:String, currentDate:String) { | |
| val srcMap = Map( | |
| "collection" -> collection, | |
| "query" -> query, | |
| "fields" -> fields, | |
| "filters" -> filters, | |
| //"solr.params" -> "start=10", // sample | |
| "flatten_multivalued" -> "false") | |
| val srcDF = spark.read.format("solr").options(srcMap).load | |
| // Upload a file as a new object with ContentType and title specified. | |
| srcDF.repartition(repartitionSize.toInt).write.mode(SaveMode.Overwrite).parquet("s3a://" + bucketName + File.separator + collection + File.separator + ds + File.separator + currentDate) | |
| } | |
| def handleException(e: Exception) { | |
| print("\n>>>>> " + e.getClass().getName() + ": " + e.getLocalizedMessage()) | |
| } | |
| } | |
| // avoid duplicate column errors for fields like "Content_Type_s" and "content_type_s" | |
| spark.sqlContext.sql("set spark.sql.caseSensitive=true") | |
| // add row count listener | |
| sc.addSparkListener(new SparkListener() { | |
| override def onTaskEnd(taskEnd: SparkListenerTaskEnd) { | |
| synchronized { | |
| recordsWrittenCount += taskEnd.taskMetrics.outputMetrics.recordsWritten | |
| } | |
| } | |
| }) | |
| // parse args | |
| val collection = sc.getConf.get("spark.parquet.publish.collection") | |
| val query = sc.getConf.get("spark.parquet.publish.query") | |
| val fields = sc.getConf.get("spark.parquet.publish.fields") | |
| val filters = sc.getConf.get("spark.parquet.publish.filters") | |
| val repartitionSize = sc.getConf.get("spark.parquet.publish.repartitionSize") | |
| val bucketName = sc.getConf.get("spark.parquet.publish.s3.bucket") | |
| val iamRole = sc.getConf.get("spark.parquet.publish.s3.iamRole") | |
| val ds = filters.split(":")(1) | |
| val now = Calendar.getInstance().getTime() | |
| val currentDate = new SimpleDateFormat("yyyyMMdd").format(now) | |
| sc.hadoopConfiguration.set("fs.s3a.credentialsType", "AssumeRole") | |
| sc.hadoopConfiguration.set("fs.s3a.stsAssumeRole.arn", iamRole) | |
| sc.hadoopConfiguration.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") | |
| // run it | |
| SolrToParquet.write(collection, query, fields, filters, ds, repartitionSize, bucketName, iamRole, currentDate) | |
| var df = Seq( | |
| (ds, recordsWrittenCount) | |
| ).toDF("ds", "count") | |
| df.write.mode(SaveMode.Append).json("s3a://" + bucketName + File.separator + collection + File.separator + ds + File.separator + currentDate) | |
| println(recordsWrittenCount) |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| { | |
| "type": "script", | |
| "id": "export_to_paerquet_iam", | |
| "script": "import org.apache.spark.sql.functions._\nimport org.apache.spark.sql.SaveMode\nimport org.apache.hadoop.fs._\nimport org.apache.spark.scheduler.{SparkListener, SparkListenerTaskEnd}\nimport java.io._\nimport spark.implicits._\nimport java.util.Calendar\nimport java.text.SimpleDateFormat\n\nvar recordsWrittenCount = 0L\n\nobject SolrToParquet {\n /**\n * Export Solr documents to a Parquet file\n */\n def write(collection:String, query:String, fields:String, filters:String, ds:String, repartitionSize:String, bucketName:String, iamRole:String, currentDate:String) {\n val srcMap = Map(\n \"collection\" -> collection,\n \"query\" -> query,\n \"fields\" -> fields,\n \"filters\" -> filters,\n //\"solr.params\" -> \"start=10\", // sample\n \"flatten_multivalued\" -> \"false\")\n val srcDF = spark.read.format(\"solr\").options(srcMap).load\n\n // Upload a file as a new object with ContentType and title specified.\n //srcDF.repartition(repartitionSize.toInt).write.mode(SaveMode.Overwrite).parquet(\"s3n://\" + bucketName + File.separator + collection + File.separator + ds + File.separator + currentDate)\n srcDF.repartition(repartitionSize.toInt).write.mode(SaveMode.Overwrite).parquet(\"s3a://\" + bucketName + File.separator + collection + File.separator + ds + File.separator + currentDate)\n }\n\n def handleException(e: Exception) {\n print(\"\\n>>>>> \" + e.getClass().getName() + \": \" + e.getLocalizedMessage())\n }\n}\n\n// avoid duplicate column errors for fields like \"Content_Type_s\" and \"content_type_s\"\nspark.sqlContext.sql(\"set spark.sql.caseSensitive=true\")\n\n// add row count listener\nsc.addSparkListener(new SparkListener() {\n override def onTaskEnd(taskEnd: SparkListenerTaskEnd) {\n synchronized {\n recordsWrittenCount += taskEnd.taskMetrics.outputMetrics.recordsWritten\n }\n }\n})\n\n// parse args\nval collection = sc.getConf.get(\"spark.parquet.publish.collection\")\nval query = sc.getConf.get(\"spark.parquet.publish.query\")\nval fields = sc.getConf.get(\"spark.parquet.publish.fields\")\nval filters = sc.getConf.get(\"spark.parquet.publish.filters\")\nval repartitionSize = sc.getConf.get(\"spark.parquet.publish.repartitionSize\")\nval bucketName = sc.getConf.get(\"spark.parquet.publish.s3.bucket\")\nval iamRole = sc.getConf.get(\"spark.parquet.publish.s3.iamRole\")\nval ds = filters.split(\":\")(1)\n\nval now = Calendar.getInstance().getTime()\nval currentDate = new SimpleDateFormat(\"yyyyMMdd\").format(now)\n\nsc.hadoopConfiguration.set(\"fs.s3a.credentialsType\", \"AssumeRole\")\nsc.hadoopConfiguration.set(\"fs.s3a.stsAssumeRole.arn\", iamRole)\nsc.hadoopConfiguration.set(\"fs.s3a.impl\", \"org.apache.hadoop.fs.s3a.S3AFileSystem\")\n\n// run it\nSolrToParquet.write(collection, query, fields, filters, ds, repartitionSize, bucketName, iamRole, currentDate)\n\nvar df = Seq(\n (ds, recordsWrittenCount)\n).toDF(\"ds\", \"count\")\n\ndf.write.mode(SaveMode.Append).json(\"s3a://\" + bucketName + File.separator + collection + File.separator + ds + File.separator + currentDate)\n\nprintln(recordsWrittenCount)\n", | |
| "sparkConfig": [ | |
| { | |
| "key": "spark.parquet.publish.collection", | |
| "value": "Amgen_Test" | |
| }, | |
| { | |
| "key": "spark.parquet.publish.query", | |
| "value": "*:*" | |
| }, | |
| { | |
| "key": "spark.parquet.publish.fields", | |
| "value": "id,body_t,_raw_content_" | |
| }, | |
| { | |
| "key": "spark.parquet.publish.filters", | |
| "value": "_lw_data_source_s:coretellis-clintrials" | |
| }, | |
| { | |
| "key": "spark.parquet.publish.repartitionSize", | |
| "value": "1" | |
| }, | |
| { | |
| "key": "spark.parquet.publish.s3.bucket", | |
| "value": "org.sbnc.roy" | |
| }, | |
| { | |
| "key": "spark.parquet.publish.s3.iamRole", | |
| "value": "arn:aws:iam::137571754341:role/TechnologyRole" | |
| } | |
| ], | |
| "shellOptions": [ | |
| { | |
| "key": "--jars", | |
| "value": "/Users/roykiesler/fusion-4.1.2/4.1.2/apps/libs/hadoop-aws-2.8.0.jar,/Users/roykiesler/fusion-4.1.2/4.1.2/apps/libs/aws-java-sdk-core-1.11.221.jar" | |
| } | |
| ] | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment