Skip to content

Instantly share code, notes, and snippets.

@rajanim
Created June 17, 2019 13:56
Show Gist options
  • Select an option

  • Save rajanim/3f95c4de23d35232ff9c11b1678afb61 to your computer and use it in GitHub Desktop.

Select an option

Save rajanim/3f95c4de23d35232ff9c11b1678afb61 to your computer and use it in GitHub Desktop.
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)
{
"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