Skip to main content

Spark SQL by gopal













SPARK SQL – READING AND WRITING THE DATA IN DIFF FORMATS
scala> val path = "/home/gopalkrishna/PRAC/SparkSQL/InputData.json"
path: String = /home/gopalkrishna/PRAC/SparkSQL/InputData.json
scala> val jsonFile = spark.read.json(path)
jsonFile: org.apache.spark.sql.DataFrame = [Address: string, Age: bigint .. 4 more fields]
scala> jsonFile.show
scala> jsonFile.printSchema
root
 |-- Address: string (nullable = true)
 |-- Age: long (nullable = true)
 |-- Desg: string (nullable = true)
 |-- State: string (nullable = true)
 |-- YrsOfExp: double (nullable = true)
 |-- name: string (nullable = true)
scala> jsonFile.select("name","Age","Address","State").show
scala> jsonFile.selectExpr("name","Desg").show
scala> jsonFile.select($"name",$"Desg").show
scala> jsonFile.select(jsonFile("name"),jsonFile("Desg")).show

scala> jsonFile.select($"name",$"Address",$"Age" , $"Age"+ 5).show
scala> jsonFile.selectExpr("name","Address","Age" , "Age+ 5").show
scala> jsonFile.select(jsonFile("name"),jsonFile("Age")+10).show

scala> jsonFile.filter($"Age" > 16).show
scala> jsonFile.filter($"Age" > 20 && $"Age" < 40).show
scala> jsonFile.filter($"Age" > 20 || $"Age" < 40).show
scala> jsonFile.filter("Age > 20").show
scala> jsonFile.where("Age > 20").show
scala> jsonFile.where(jsonFile("Address").contains("Hyderabad")).show
scala> jsonFile.limit(4).show
scala> jsonFile.head
scala> jsonFile.take(3)
scala> jsonFile.count
scala> jsonFile.groupBy($"Desg").count.show
+-------------+-----+                                                          
|         Desg  |count  |
+-------------+-----+
|      STA    |    7      |
|      null    |    9       |
|SeniorAnalyst|    1|
|  Sw Engineer|    9|
|           TA|           1|
+-------------+-----+


scala>
scala> jsonFile.sort(jsonFile("name"),jsonFile("YrsOfExp")).show
scala> jsonFile.orderBy(jsonFile("Address")).show
scala> jsonFile.groupBy(jsonFile("Address"),jsonFile("name")).count.show
scala> jsonFile.agg(max("Age"),min("Age") ).show
+--------+--------+
|max(Age)|min(Age)|
+--------+--------+
|      48|       6|
+--------+--------+

scala> jsonFile.describe("YrsOfExp","Age").show
+-------+------------------+------------------+
|summary|          YrsOfExp|               Age|
+-------+------------------+------------------+
|  count|                18|                18|
|   mean|6.6833333333333345| 25.88888888888889|
| stddev| 5.683852256224491|12.101963325835575|
|    min|               1.2|                 6|
|    max|              12.5|                48|
+-------+------------------+------------------+

scala> jsonFile.drop("name","Address","YrsOfExp","Age").show

scala> jsonFile.registerTempTable("empjsontab")
warning: there was one deprecation warning; re-run with -deprecation for details
scala> spark.sql("select name,COUNT(Age) from empjsontab GROUP BY name ORDER BY COUNT(Age) DESC").show
+-------+----------+                                                           
|   name|count(Age)|
+-------+----------+
|  Ramya|         9|
|Mounika|         9|
|  Gopal|         0|
+-------+----------+
scala> spark.sql("select name, Age , Age + 10 from empjsontab").show

CSV FILE:
scala> val path = "/home/gopalkrishna/PRAC/SparkSQL/empdata.csv"
path: String = /home/gopalkrishna/PRAC/SparkSQL/empdata.csv

scala> val csvFile = spark.read.csv(path)
csvFile: org.apache.spark.sql.DataFrame = [_c0: string, _c1: string ... 1 more field]

scala> csvFile.printSchema
root
 |-- _c0: string (nullable = true)
 |-- _c1: string (nullable = true)
 |-- _c2: string (nullable = true)


scala> csvFile.show
+----+-------+-----+
| _c0|    _c1|  _c2|
+----+-------+-----+
|1000|  Gopal|12000|
|1001|Krishna|14000|
|1002|   Ravi|16000|
|1002|  Ramya|24000|
|1003| Rakesh|34000|
|1004| Rajesh|22000|
|1005|Mounika|24000|
|1006| Sravya|26000|
|1007|   Siya|24000|
|1008| Bhavya|30000|
|1009|Trinath|34000|
+----+-------+-----+
scala> csvFile.orderBy("_c2").show
+----+-------+-----+
| _c0|    _c1|  _c2|
+----+-------+-----+
|1000|  Gopal|12000|
|1001|Krishna|14000|
|1002|   Ravi|16000|
|1004| Rajesh|22000|
|1005|Mounika|24000|
|1002|  Ramya|24000|
|1007|   Siya|24000|
|1006| Sravya|26000|
|1008| Bhavya|30000|
|1003| Rakesh|34000|
|1009|Trinath|34000|
+----+-------+-----+


scala>


scala> csvFile.groupBy("_c2").count.show
+-----+-----+                                                                  
|  _c2|count|
+-----+-----+
|14000|    1|
|26000|    1|
|12000|    1|
|24000|    3|
|22000|    1|
|30000|    1|
|16000|    1|
|34000|    2|
+-----+-----+
==============================================
scala>
PARQUETFILE
scala> val path = "/home/gopalkrishna/PRAC/SparkSQL/users.parquet"
path: String = /home/gopalkrishna/PRAC/SparkSQL/users.parquet

scala> val parFile = spark.read.parquet(path)
SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
SLF4J: Defaulting to no-operation (NOP) logger implementation
SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.
parFile: org.apache.spark.sql.DataFrame = [name: string, favorite_color: string ... 1 more field]

scala> parFile.printSchema
root
 |-- name: string (nullable = true)
 |-- favorite_color: string (nullable = true)
 |-- favorite_numbers: array (nullable = true)
 |    |-- element: integer (containsNull = true)
scala> parFile.show
17/06/08 04:20:03 WARN ParquetRecordReader: Can not initialize counter due to context is not a instance of TaskInputOutputContext, but is org.apache.hadoop.mapreduce.task.TaskAttemptContextImpl
+------+--------------+----------------+
|  name|favorite_color|favorite_numbers|
+------+--------------+----------------+
|Alyssa|          null|  [3, 9, 15, 20]|
|   Ben|           red|              []|
+------+--------------+----------------+
scala> parFile.select($"name").show
+------+
|  name|
+------+
|Alyssa|
|   Ben|
+------+

scala>

scala>
++++++++++++++++++++++++++++++++++
textFile
scala> val path = "/home/gopalkrishna/PRAC/SparkSQL/kv1.txt"
path: String = /home/gopalkrishna/PRAC/SparkSQL/kv1.txt

scala> val textFile1 = spark.read.textFile(path)
textFile1: org.apache.spark.sql.Dataset[String] = [value: string]

scala> textFile1.printSchema
root
 |-- value: string (nullable = true)


scala> textFile1.show
+-----------+
|      value|
+-----------+
|238val_238|
|  86val_86|
|311val_311|
|  27val_27|
|165val_165|
|409val_409|
|255val_255|
|278val_278|
|  98val_98|
|484val_484|
|265val_265|
|193val_193|
|401val_401|
|150val_150|
|273val_273|
|224val_224|
|369val_369|
|  66val_66|
|128val_128|
|213val_213|
+-----------+
only showing top 20 rows


scala>



----------------------

RDD OPERATIONS RELATED

scala> val data = Array(1,2,3,4,5,6,6,7,8)
data: Array[Int] = Array(1, 2, 3, 4, 5, 6, 6, 7, 8)

scala> val distriData = sc.parallelize(data)
distriData: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:26

scala> distriData.map(_ + 2).collect().mkString("\n")
res0: String =
3
4
5
6
7
8
8
9
10

scala> val sumData = distriData.map(_ + 2)    // Transformed RDD
sumData: org.apache.spark.rdd.RDD[Int] = MapPartitionsRDD[2] at map at <console>:28

scala> sumData.reduce(_+_)       // ActionRDD
res1: Int = 60

scala>

DATA FRAME ( schemaRDD)

DataFrame is an abstraction which gives a schema view of data. Which means it gives us a view of data as columns with column name and types info, We can think data in data frame like a table in the database.

DATA FRAME using CASE CLASS

scala> case class Person(name : String , age:Int , address:String)
defined class Person

scala> val df = List(Person("Raja",21,"HYD"),Person("Ramya",34,"BAN"),Person("Rani",30,"MUM")).toDF
df: org.apache.spark.sql.DataFrame = [name: string, age: int ... 1 more field]

scala> df.collect().mkString("\n")
res2: String =
[Raja,21,HYD]
[Ramya,34,BAN]
[Rani,30,MUM]

scala> df.show
+-----+---+-------+
| name|age|address|
+-----+---+-------+
| Raja| 21|    HYD|
|Ramya| 34|    BAN|
| Rani| 30|    MUM|
+-----+---+-------+


scala> df.filter("age > 25").show
+-----+---+-------+
| name|age|address|
+-----+---+-------+
|Ramya| 34|    BAN|
| Rani| 30|    MUM|
+-----+---+-------+

scala> df.filter("salary > 25").show
org.apache.spark.sql.AnalysisException: cannot resolve '`salary`' given input columns: [name, age, address]; line 1 pos 0

DATA SET [DS]

Data Set is an extension to Dataframe API, the latest abstraction which tries to provide best of both RDD and Dataframe.

CONVERT “DATA FRAME(DF)”  TO  “DATASET(DS)”
NOTE: we can always convert a data frame at any point of time into a dataset by calling ‘as’ method on Dataframe. Example:  df.as[MyClass]

i.e by providing the case class only we can convert a DATA FRAME into DATA SET.

scala> val ds = df.as[Person]
ds: org.apache.spark.sql.Dataset[Person] = [name: string, age: int ... 1 more field]


scala> ds.show
+-----+---+-------+
| name|age|address|
+-----+---+-------+
| Raja| 21|    HYD|
|Ramya| 34|    BAN|
| Rani| 30|    MUM|
+-----+---+-------+


scala> ds: org.apache.spark.sql.Dataset[Person] = [name: string, age: int ... 1 more field]

scala> ds.show
+-----+---+-------+
| name|age|address|
+-----+---+-------+
| Raja| 21|    HYD|
|Ramya| 34|    BAN|
| Rani| 30|    MUM|
+-----+---+-------+

scala> ds.filter(_.age > 21).show()
+-----+---+-------+
| name|age|address|
+-----+---+-------+
|Ramya| 34|    BAN|
| Rani| 30|    MUM|
+-----+---+-------+


scala> ds.filter(_.salary > 21).show()
<console>:28: error: value salary is not a member of Person
       ds.filter(_.salary > 21).show()

OBSERVATION : Unlike data frame , which is giving a runtime exception saying “can not resolve salary” , analysis exception , dataset is showing the COMPILE TIME ERROR only.

So Datasets API provides compile time safety which was not available in Data frames

CONVERTING “DATASET[DS]” to “DATA FRAME[DF]”

We can directly use toDF method to convert Data Set back to Data Frame , No need of using any Case Class over here

scala> val newdf = ds.toDF
newdf: org.apache.spark.sql.DataFrame = [name: string, age: int ... 1 more field]

scala> newdf.show
+-----+---+-------+
| name|age|address|
+-----+---+-------+
| Raja| 21|    HYD|
|Ramya| 34|    BAN|
| Rani| 30|    MUM|
+-----+---+-------+


scala>

READING “JSON” DATA using “DATA FRAME” & CONVERTING INTO “DATA SET”

scala> case class Emp(name:String,Desg:String,YrsOfExp:Double,Address:String,State:String)
defined class Emp

scala> val df = spark.read.json("file:///home/gopalkrishna/PRAC/SparkSQL/InputData.json")
df: org.apache.spark.sql.DataFrame = [Address: string, Age: bigint ... 4 more fields]

scala> df.show(4)
+---------+----+-----------+---------+--------+-------+
|  Address| Age|       Desg|    State|YrsOfExp|   name|
+---------+----+-----------+---------+--------+-------+
|Hyderabad|null|        STA|Telangana|    12.5|  Gopal|
|Bangalore|   6|       null|Karnataka|    null|Mounika|
|  Chennai|  22|Sw Engineer|TamilNadu|     1.2|  Ramya|
|Hyderabad|null|         TA|Telangana|    12.5| Sekhar|
+---------+----+-----------+---------+--------+-------+
only showing top 4 rows


scala> val ds = spark.read.json("file:///home/gopalkrishna/PRAC/SparkSQL/InputData.json").as[Emp]
ds: org.apache.spark.sql.Dataset[Emp] = [Address: string, Age: bigint ... 4 more fields]

scala>

TO CONVERT THE ABOVE “DATA FRAME(df)”  to “DATA SET(ds)” – as[case class]


Comments

Popular posts from this blog

Interview questions by company wise

Interview questions: Amazon link TCS Interview: 1) How you handle duplicates in hive table? 2) What is map side join and what is SMB join? 3) Why we use splitby? 4) How to import/view only schema using sqoop? 5) Diff b/n reduceByKey and groupByKey? 6) Broadcast variables and accumulators? 7) How to add new column to DF? 8) Diff b/n lineage and DAG? 9) Spark job submission command? Round-2: 1) Singleton scala object? 2) Higher order function in scala? 3) How to load data from file having multiple delimiters? 4) Compare orc and parquet 5) DAG and lineage difference? 6) Difference b/n rdd, Df and ds 7) How DS is created from DF and how rdd is converted to DF? Accenture interview: 1) Explain Ur project. 2) How u do incremental import for updating tables in Ur project? 3) DAG and lineage difference. 4) reduceByKey functioning? 5) Spark Architecture. 6) How u export jar in Ur project and what spark libraries u use for creatin...

Interview written