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
Post a Comment