์ด ๋ฌธ์๋ Spark ์์ Scala ์ธ์ด๋ฅผ ์ฌ์ฉํ๊ณ ์ ํ๋ ๊ฐ๋ฐ์๋ฅผ ์ํ ํํ ๋ฆฌ์ผ์ด๋ค.
์ด ๋ฌธ์๋ฅผ ์ฝ๊ธฐ ์ํด์๋ Java ๋๋ Python ์ ๋ํ ์ง์์ด ์์ด์ผ ํ๊ณ , Java ๋๋ Python ์ ์ด์ฉํด Spark ๋ฅผ ์ด์ฉํ ๊ฒฝํ์ด ์์ด์ผ ํ๋ค.
๋ฌธ์์ ์์ ์ค์ด๊ธฐ ์ํด Spark ์ ๋ํ ์ค๋ช ์ ์๋ต๋๋ฉฐ, ๋ํ Scala ๋ฌธ๋ฒ์ค Spark ๊ฐ๋ฐ์ ๋ถํ์ํ๋ค๊ณ ํ๋จ๋๋ ๋ฌธ๋ฒ ๋ํ ์๋ตํ๋ค.
์๋ ์ฃผ์์ ๋ด์ฉ์ ์ฐธ๊ณ ํ๋ค.
์ด ๋ฌธ์์์๋ SBT ๋ฅผ ์ด์ฉํด ์ํ์ฝ๋๋ฅผ ๋น๋ํ๊ณ ์คํํ๋ค. ํด์ ๋ค์ด๋ก๋ ๋ฐ ์ค์น๋ฐฉ๋ฒ์ ์ฌ๊ธฐ๋ฅผ ์ฐธ๊ณ ํ๊ธฐ ๋ฐ๋๋ค.
$ curl https://bintray.com/sbt/rpm/rpm > bintray-sbt-rpm.repo
$ sudo mv bintray-sbt-rpm.repo /etc/yum.repos.d/
$ sudo yum install sbt
$ sbt
sbt:ec2-user> exit
$์๋์ ๊ฐ์ ๋ฐฉ๋ฒ์ผ๋ก ๊ฐ๋จํ ํ ์คํธ ํ๋ก๊ทธ๋จ์ ์คํ์ํฌ ์ ์๋ค.
$ mkdir sample
$ cd sample/
$ sbt console
scala> println("Hello, World!")
Hello, World!
scala> :q
$์ดํ์ ์๋ ์ฝ๋๋ค์ sbt console ๋ช
๋ น์ผ๋ก ์ฝ์์ ๋ก๊ทธ์ธ๋์ด ์๋ ๊ฒ์ ์ ์ ๋ก ํ๋ค.
์๋์ ๋ฐฉ๋ฒ์ผ๋ก ๋ณ์๋ฅผ ์์ฑํ ์ ์๋ค.
val i: Int = 1
i = 2 // error๋ณ์๋ฅผ ์์ฑํ๋ ํค์๋๋ val ๊ณผ var ๊ฐ ์๋ค. val ๋ก ์์ฑํ ๋ณ์๋ ๊ฐ์ ๋ณ๊ฒฝ์ด ๋ถ๊ฐ๋ฅํ ๋ณ์๊ฐ ๋๋ค. ๋ฐ๋ฉด์ var ๋ก ์์ฑํ ๋ณ์๋ ๊ฐ์ ๋ณ๊ฒฝ์ด ๊ฐ๋ฅํ๋ค. ํ์ง๋ง, Scala ์์๋ var ๋ฅผ ์ฌ์ฉํ์ง ์์ ๊ฒ์ ๊ถ์ฅํ๊ณ ์๋ค.
val j = 2Scala ์์๋ ๋ณ์๊ฐ์ ํ์
์ ์ ์ ์๋ ๊ฒฝ์ฐ(์์์ 2 ๋ Int) ์์๊ฐ์ด ๋ณ์ํ์
(: Int)์ ์๋ตํ ์ ์๋ค.
val k = 3
println("class: " + k.getClass)
// class: intScala ์์๋ ๋ชจ๋ ๋ณ์๋ ๊ฐ์ฒด๋ค. ์์์ k ๋ ๋จ์ํ ์ ์๊ฐ์ด ์๋ ์ ์ํ ๊ฐ์ฒด๊ฐ ๋๋ค.
์๋์ ๊ฐ์ด ํจ์๋ฅผ ์์ฑํ ์ ์๋ค.
def addOne(m: Int): Int = m + 1
val three = addOne(2)
println(three)
// three: Int = 3์์์ = ๋ค ๋ถ๋ถ์ด ํจ์์ ๋ด์ฉ์ด์ง๋ง return ํค์๋๊ฐ ์๋ค. Scala ์์๋ return ํค์๋์ ์๋ต์ ๊ถ์ฅํ๋ค.
def three() = 1 + 2
three()
three์์์ ํจ์ ์ ์์์ 1 + 2 ์ ๊ฐ์ด ์ ์์ด๋ฏ๋ก : Int ๊ฐ ์๋ต๋์ด๋ ์ ์ ์๋ํ๋ค.
ํ๋ผ๋ฏธํฐ๊ฐ ์๋ ๊ฒฝ์ฐ () ๋ฅผ ์๋ตํ ์ ์๋ค.
ํจ์๊ฐ ๋ผ์ธ์๊ฐ ๋ง์ ๊ฒฝ์ฐ ์๋์ ๊ฐ์ด ๊ดํธ๋ฅผ ์ถ๊ฐํ๋ค.
def addOne(m: Int): Int = {
m + 1
}
print(addOne(2))์๋ ์ฝ๋๋ฅผ Scala ์ฝ์์ ์ ๋ ฅํด๋ณด์
class Calculator {
val brand: String = "HP"
def add(m: Int, n: Int): Int = m + n
}
val calc = new Calculator
println(calc.add(1, 2))
println(calc.brand)ํ๋(๋ฉค๋ฒ ๋ณ์)๋ val ๋ก, ๋ฉ์๋(๋ฉค๋ฒ ํจ์)๋ def ๋ก ์ ์ํ๋ค.
Scala ์์ ์์ฑ์๋ ๊ดํธ์ ์์ฒด์ด๋ค.
class Calculator(brand: String) {
println("start constructor")
val color: String = if (brand == "TI") {
"blue"
} else if (brand == "HP") {
"black"
} else {
"white"
}
def add(m: Int, n: Int): Int = m + n
println("end constructor")
}
val calc = new Calculator("HP")
println(calc.color)์ ์ฝ๋์์ println ์ด ๋๋ฒ ์คํ๋ ๊ฒ์ ๋ณผ ์ ์๋ค.
๋ํ, if ๋ฌธ์ฅ์ด ๋ฆฌํด๊ฐ์ ๋ฐํํด์ ๋ณ์์ ์
๋ ฅ๋๊ณ ์๋ ๊ฒ์ ๋ณผ ์ ์๋ค. Scala ์์๋ ๋๋ถ๋ถ์ ํํ์์ด ๋ฆฌํด๊ฐ์ ๊ฐ์ง๋ฉฐ return ํค์๋ ์์ด๋ ํจ์์ ๋ฆฌํด๊ฐ์ผ๋ก ๋ฐํ๋๋ค.
์์ฑ์์ ์ ๋ฌ๋ ํ๋ผ๋ฏธํฐ๋ ์์ฑ์๊ฐ ์คํ๋ ํ์๋ ์ฌ๋ผ์ง๋ค.
class Person(name: String, age: Int)
val person = new Person("mong", 9)
println(person.age) // error์ ๋ฌ๋ ํ๋ผ๋ฏธํฐ๋ฅผ ํด๋์ค์ ๋งด๋ฒํ๋๋ก ๋ง๋ค๋ ค๋ฉด ์๋์ ๊ฐ์ด val ์ ๋ถ์ฌ์ฃผ์ด์ผ ํ๋ค.
class Person(val name: String, val age: Int)
val person = new Person("mong", 9)
println(person.age) // okget,set ๋ฉ์๋๋ ์๋์ผ๋ก ์ถ๊ฐ๋๋ฏ๋ก ๋ณ๋๋ก ์์
ํ ํ์๊ฐ ์๋ค. ๋ํ, ๋ค๋ฅธ ์ธ์ด์ ๋ค๋ฅด๊ฒ ๋งด๋ฒ๋ณ์ ๋ฐ ๋งด๋ฒํจ์๊ฐ private ์ ๋ณ๋๋ก ์ง์ ํด ์ฃผ์ง ์๋ ํ, public ์ด ๋ํดํธ๋ก ์ง์ ๋๋ค.
์ผ๋ฐ์ ์ธ switch case ๋ฌธ๋ณด๋ค ๋ ๋ง์ ๊ธฐ๋ฅ์ ์ ๊ณตํ๋ค.
val times = 3
times match {
case 1 => "one"
case 2 => "two"
case i if i == 3 => "three"
case i if i == 4 => "four"
case _ => "some other number"
}๋จ์ํ ์ ์๋งค์นญ์ด๋ ๋ฌธ์์ด๋งค์นญ ๋ฟ๋ง ์๋๋ผ ์กฐ๊ฑด๋ฌธ์ ์ด์ฉํด ๋งค์นญํ ์ ์๋ค.
๋ง์ง๋ง์ ๋ณด์ด๋ _ ์ ์์ผ๋์นด๋๋ก ์ฌ์ฉ๋๋ค. ์ฌ๊ธฐ์๋ case else ๋ก ์ฌ์ฉ๋๊ณ import org.apache.spark.SparkContext._ ์ ๊ฐ์ ๊ฒฝ์ฐ์๋ ํ์์ ์๋ ๋ชจ๋ ๊ฒ์ ์ํฌํธํ๋ค. ์์์ case _ ๊ฐ ์๋ค๋ฉด ๋งค์นญ๋๋ ๊ฐ์ด ์์ ๋ ์๋ฌ๊ฐ ๋ฐ์ํ๋ค.
๊ฐ์ ๋ํ ๋งค์นญ ๋ฟ๋ง ์๋๋ผ ํ์ ์ ๋ํด์๋ ํจํด ๋งค์นญ์ด ๊ฐ๋ฅํ๋ค.
def bigger(o: Any): Any = {
o match {
case i: Int if i < 0 => i - 1
case i: Int => i + 1
case d: Double if d < 0.0 => d - 0.1
case d: Double => d + 0.1
case text: String => text + "s"
case _ => "what is it?"
}
}
println(bigger("cat"))์์์ ์ ์ ์ค์ ๋ฟ๋ง ์๋๋ผ ๋ฌธ์์ด๊ณผ๋ ๋งค์นญํจ์ ๋ณผ ์ ์๋ค.
ํด๋์ค์ ๋ํด์๋ ๋์ผํ ๋ฐฉ์์ผ๋ก ํจํด ๋งค์นญ์ด ๊ฐ๋ฅํ๋ค.
class Person(val name: String, val age: Int)
def isYoungPerson(person: Person) = person match {
case p if p.age < 20 => "Yes"
case _ => "No"
}case class ๋ฅผ ์ด์ฉํ๋ฉด new ๋ฅผ ์ฌ์ฉํ์ง ์์๋ ํด๋์ค๋ฅผ ์์ฑํ ์ ์๋ค.
class Person(name: String, age: Int)
val a = Person("Lee", 21) // error
val a = new Person("Lee", 21) // ok
println(a.age) // error
case class Person(name: String, age: Int)
val a = Person("Lee", 21) // ok
println(a.age) // ok์ผ์ด์ค ํด๋์ค๋ ํจํด ๋งค์นญ์ ์ฌ์ฉํ๊ธฐ ์ํด ์ค๊ณ๋์๋ค.
case class Person(name: String, age: Int)
def isYoungPerson(person: Person) = person match {
case Person("Lee", 12) => "Yes"
case Person(_, 12) => "Yes"
case _ => "No"
}
val p = Person("Lee", 12)
println(isYoungPerson(p))
val p2 = Person("Moon", 12)
println(isYoungPerson(p2))์์์ new ํค์๋ ์์ด ํด๋์ค๊ฐ ์์ฑ๋จ์ ๋ณผ ์ ์๋ค. ๋ํ ์์ผ๋์นด๋ ๋ฌธ์์ธ _ ๊ฐ ํด๋์ค์์ฑ์๋ ์ฌ์ฉ๋์์์ ๋ณผ ์ ์๋ค.
List ์๋ ๋์ผ ํ์
์ ๋ฐ์ดํ๋ง ์
๋ ฅํ ์ ์๊ณ ์ค๋ณต๋ ๋ฐ์ดํ๋ ์
๋ ฅ ๊ฐ๋ฅํ๋ค. Set ์๋ ์ค๋ณต๋๋ ๋ฐ์ดํ๋ฅผ ์
๋ ฅํ ์ ์๋ค. Tuple ์๋ ์๋ก ๋ค๋ฅธ ํ์
์ ๋ฐ์ดํ๋ฅผ ๋ฌถ์ ์ ์๋ค. Tuple ์ ์ฒซ๋ฒ์งธ ๋ฐ์ดํ ํธ์ถ์ ._0 ์ด ์๋ ._1 ์ ์ฌ์ฉํ๊ณ ์๋ค.
val numbers = List(1, 2, 3, 4)
println(numbers(2))
val animals = Set("Cat", "Dog", "Tiger")
println(animals("Cat"))
val hostPort = ("localhost", 80)
println(hostPort._1)
val a = 1 -> 2
println(a)-> ๋ฅผ ์ด์ฉํด ํํ์ ์์ฑํ ์ ์๋ค.
key-value ํํ์ ๊ฐ์ ๋ฌถ์์ด Map ์ด๋ค.
val m = Map(1 -> "one", 2 -> "two")
println(m(2))์์์ -> ๋ ํน๋ณํ ๋ฌธ๋ฒ์ด ์๋๊ณ ํํ์ ์์ฑ์ ๋ถ๊ณผํ๋ค. ์์์ ์์ฑ๋ ๋งต์ ์ค์ ๋ก๋ Map((1, "one"), (2, "two")) ์ ํํ๊ฐ ๋๊ณ , ๋งต์ ๋ค์ด์๋ ๋ฐ์ดํ๋ ์ฒซ๋ฒ์งธ ๊ฐ์ด key ๊ฐ ๋๊ณ , ๋๋ฒ์งธ ๊ฐ์ด value ๊ฐ ๋๋ค.
๋ฆฌ์คํธ๋ฅผ ์ ๋ฌ๋ฐ์ ์ผ์ ํ ์ฒ๋ฆฌ๋ฅผ ํ๊ณ ์ฒ๋ฆฌ๋ ๊ฐ์ ์ ๋ฌํด์ฃผ๋ ๊ฒ์ ํฉ์์กฐํฉ์ด๋ผ๊ณ ํ๋ค.
๋ค๋ฅธ ์ธ์ด์์๋ for (int i = 0; i < 10; i++) { ... } ์คํ์ผ๋ก ์ฝ๋ฉํ๋ ๊ฒฝ์ฐ๊ฐ ๋ง์ง๋ง, Scala ์์๋ ๋ณ์ ์์ฑ์ ์ง์ํ๋ค.
val numbers = List(1, 2, 3, 4)
println(numbers.map((i: Int) => i * 2))
println(numbers.map(i => i * i))์์ ๊ฐ์ด for ๋ฌธ ๋์ ์ map() ์ ์ด์ฉํด ์
๋ ฅ๋ ๋ฐ์ดํ๋ฅผ ๊ฐ ๊ฐ๋ค์ ์ฐ์ฐํ ์ ์๋ค. ์
๋ ฅ๋๋ ๋ฐ์ดํ๊ฐ ์ ์ํ์ด ํ์คํ๋ฏ๋ก : Int ๋ ์๋ตํ ์ ์๋ค.
map() ๊ณผ ๋ณ๋์ ํจ์๋ฅผ ์กฐํฉํ ์๋ ์๋ค.
val numbers = List(1, 2, 3, 4)
def square(i: Int) = i * i
println(numbers.map(square _))์์์ square _ ์ square(_) ์ ๋์ผํ ๋ด์ฉ์ด๋ค. Scala ์์๋ ํ๋ผ๋ฏธํฐ๊ฐ ํ๊ฐ์ผ ๊ฒฝ์ฐ ๊ดํธ๋ฅผ ์๋ตํ ์ ์๋ค.
map() ์ด ์
๋ ฅ๋ ๋ฐ์ดํ๋ฅผ ๊ทธ๋๋ก ๋๊ณ ๋ณํ๋ ๋ฐ์ดํ๋ฅผ ๋ฐํํ๋๊ฒ๊ณผ ๋ค๋ฅด๊ฒ, foreach() ๋ ์
๋ ฅ๋ ๊ฐ ์์ฒด๋ฅผ ๋ณํํ๊ณ ๋ฆฌํด๊ฐ์ด ์๋ค.
val numbers = List(1, 2, 3, 4)
numbers.foreach(i => println(i))foreach() ์๊ฒ ๋ฆฌํด๊ฐ์ ์์ฒญํ๋ฉด Unit(๋ค๋ฅธ ์ธ์ด์์๋ void) ์ด ๋ฐํ๋๋ค.
์ ๋ ฅ๋ ๊ฐ์ ํํฐ๋งํด์ ๊ฐ์ด ์ฐธ์ธ ๊ฒ๋ค๋ก ์ด๋ฃจ์ด์ง ๋ฆฌ์คํธ๋ฅผ ๋ฐํํ๋ค.
val numbers = List(1, 2, 3, 4)
println(numbers.filter(i => i % 2 == 0))๋๊ฐ์ ๋ฆฌ์คํธ๋ฅผ ๊ฐ๊ฐ์ ๋ฐ์ดํ๋ฅผ ๋ฌถ์ด ํํ ๋ฆฌ์คํธ๋ก ๋ง๋ ๋ค.
val numbers = List(1, 2, 3, 4)
val animals = List("dog", "cat", "lion", "tiger")
println(numbers.zip(animals))
// List((1,dog), (2,cat), (3,lion), (4,tiger))
val numbers = List(1, 2, 3, 4)
val animals = List("dog", "cat", "lion")
println(numbers.zip(animals))
// List((1,dog), (2,cat), (3,lion))
val numbers = List(1, 2, 3)
val animals = List("dog", "cat", "lion", "tiger")
println(numbers.zip(animals))
// List((1,dog), (2,cat), (3,lion))๋ฐ์ดํ์ ๊ฐฏ์๊ฐ ๋ง์ง ์์ผ๋ฉด ๋ง๋ ๋งํผ๋ง ๋ฌถ์ด์ ๋ฆฌ์คํธ๋ฅผ ๋ฐํํ๋ค.
partition() ๋ ์
๋ ฅ๋ ๋ฆฌ์คํธ๋ฅผ ๋๋ก ์ชผ๊ฐ์ด ๋๊ฐ์ ๋ฆฌ์คํธ๋ฅผ ๋ฐํํ๋ค.
val numbers = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
val two = numbers.partition(_ % 2 == 0)
println(two._1)๋ฐํ๋๋ ๊ฐ์ ํํ๋ก ๋ฌถ์ฌ ์๋ค. ํ๊ฐ์ ๋ฆฌ์คํธ๋ฅผ ์ด์ฉํ๋ ค๋ฉด ํํ์ ์ ๊ทผ๋ฒ๊ณผ ๋์ผํ๊ฒ ._1 ๋๋ ._2 ๋ฅผ ์ด์ฉํ๋ฉด ๋๋ค.
find() ๋ ์กฐ๊ฑด์ ๋ง์กฑํ๋ ์ฒซ๋ฒ์งธ ๊ฐ์ ๋ฐํํ๋ค.
val numbers = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
println(numbers.find(i => i > 5))
val tup = List((1,"dog"), (2,"cat"), (3,"lion"))
println(tup.find(t => t._1 > 1 && t._2 == "lion"))ํํ์ ์ ๋ ฅ๊ฐ์ผ๋ก ๋ฐ์ ์ ์๋ค.
Option() ์ ์ด๋ค ๊ฐ์ด ์์ ์๋ ์๊ณ ์์ ์๋ ์์ ๋ ์ฌ์ฉ๋๋ค. find() ์์ ๋ฆฌํด๋๋ ๊ฐ์ด Option() ์ด๋ค.
val numbers = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
val res1 = numbers.find(i => i > 5)
val res2 = numbers.find(i => i > 10)
val result = if (res1.isDefined) { res1.get * 2 } else { 0 }
println(result)
val result = res2.getOrElse(0) * 2
println(result)find() ๋ ๋ฆฌํด๊ฐ์ด ์์ ๋ None ์ ๋ฐํํ๋ค. ๋ฐ๋ผ์ res1.isDefined ๋ฅผ ์ด์ฉํด ๊ฐ์ด ์๋์ง ์ฒดํฌํ๋ ๋ฐฉ๋ฒ์ด ์๋ค. ๋๋, res1.getOrElse(0) ๋ฅผ ์ด์ฉํด ๋ํดํธ๊ฐ์ ์ง์ ํด ์ค ์๋ ์๋ค.
drop() ์ ์
๋ ฅ๋๋ ๋ฆฌ์คํธ์์ ์์์ n ๊ฐ์ ๊ฐ์ ์์ค ๋๋จธ์ง ๋ฆฌ์คํธ๋ฅผ ๋ฐํํ๋ค.
val numbers = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
println(numbers.drop(5))
println(numbers.drop(20))
println(numbers.dropWhile(_ % 2 != 0))dropWhile() ์ ์กฐ๊ฑด์ ๋ง์กฑํ์ง ์๋ ๊ฐ์ด ์์ ๋๊น์ง์ ๊ฐ์ ์์ค ๋๋จธ์ง ๋ฆฌ์คํธ๋ฅผ ๋ฐํํ๋ค. ์์์ 2 ์์ ์กฐ๊ฑด์ ๋ง์กฑํ์ง ์์ drop ์ ์ค๋จํ๊ณ ๋๋จธ์ง ๋ฆฌ์คํธ๋ฅผ ๋ฐํํ๊ฒ ๋๋ค.
foldLeft() ๋ ์
๋ ฅ๋๋ ๋ฆฌ์คํธ์ ๊ฐ ๊ฐ๋ค์ ์ฐ์ฐํ ๊ฐ์ ๋์ ๊ฐ์ ๋ฐํํ๋ค.
val numbers = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
println(numbers.foldLeft(0) { (acc, i) => println("acc: " + acc + " i: " + i); acc + i })
println(numbers.foldLeft(1000) { (acc, i) => println("acc: " + acc + " i: " + i); acc + i })
// acc: 1000 i: 1
// acc: 1001 i: 2
// acc: 1003 i: 3
// acc: 1006 i: 4
// acc: 1010 i: 5
// acc: 1015 i: 6
// acc: 1021 i: 7
// acc: 1028 i: 8
// acc: 1036 i: 9
// acc: 1045 i: 10
// 1055์์์ 0, 1000 ์ ์์๊ฐ์ด ๋๊ณ , acc ์ ๋์ ๊ฐ์ด ์ ์ฅ๋๋ฉฐ, i ๊ฐ ์
๋ ฅ๋ ๋ฆฌ์คํธ์ ๋ฐ์ดํ์ด๋ค.
foldRight() ๋ foldLeft() ์ ๋์ผํ ๊ธฐ๋ฅ์ ํ๋๋ฐ ๋ฐฉํฅ๋ง ๊ฑฐ๊พธ๋ก์ด๋ค.
val numbers = List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
println(numbers.foldRight(1000) { (acc, i) => println("acc: " + acc + " i: " + i); acc + i })flatten() ์ ์
๋ ฅ๋ ๋ฐ์ดํ์ ์ค์ฒฉ(nested)๋จ๊ณ๋ฅผ ํ๋จ๊ณ ํ์ด์ค๋ค.
val nestedNumbers = List(List(1, 2), List(3, 4), List(5, 6))
println(nestedNumbers.flatten)
// List(1, 2, 3, 4, 5, 6)flatMap() ์ flatten() ๊ณผ map() ์ ํฉ์น๊ฒ์ด๋ค.
val nestedNumbers = List(List(1, 2), List(3, 4), List(5, 6))
println(nestedNumbers.flatMap(x => x.map(_ * 2)))
// List(2, 4, 6, 8, 10, 12)๋ฆฌ์คํธ์ ๊ฐ ๋ฐ์ดํ์ ๋ํด map() ์ ์ ์ฉํ๊ณ ๋ฆฌํด๋ ๊ฐ๋ค์ flatten() ํ๋ค.
ํจ์์กฐํฉ๊ณผ ํจํด๋งค์นญ์ ํจ๊ป ์ฌ์ฉํ๋ฉด ์๋์ ๊ฐ์ด ์ฝ๋๋ฅผ ์์ฑํ ์ ์๋ค.
val log = Array(("2018-04-11", "11:22:33", "itemid=112233"), ("2018-04-12", "11:12:32", "itemid=443322"))
val parsed = log.map(i => i match {
case (yyyymmdd, hhmmss, params) =>
println(yyyymmdd)
})ํ์ง๋ง ์ต๋ช
ํจ์(anonymous function) ๋ฅผ ์ฌ์ฉํด match ํค์๋์์ด ๊ฐ๊ฒฐํ๊ฒ ์ฝ๋๋ฅผ ์์ฑํ ์ ์๋ค. ์ ์ด๋ฐ ์ฝ๋๊ฐ ์๋ํ๋์ง ํ์ธํ๋ ค๋ฉด PartialFunction ์ ์์์ผ ํ๋๋ฐ, ๊ทธ๋ฅ ์์๋ณด์ง ์์ ๊ฒ์ ๊ถ์ฅํ๋ค. (-.-)
val log = Array(("2018-04-11", "11:22:33", "itemid=112233"), ("2018-04-12", "11:12:32", "itemid=443322"))
val parsed = log.map({
case (yyyymmdd, hhmmss, params) =>
println(yyyymmdd)
})Scala ์์๋ ํจ์์ ํ๋ผ๋ฏธํฐ๊ฐ ํ๊ฐ์ผ ๊ฒฝ์ฐ () ๋ฅผ ์๋ตํ ์ ์๋ค. ๋ฐ๋ผ์ ์์ ์์ค์ฝ๋๋ ์๋์ ๊ฐ์ด ์ธ ์ ์๋ค.
val log = Array(("2018-04-11", "11:22:33", "itemid=112233"), ("2018-04-12", "11:12:32", "itemid=443322"))
val parsed = log.map{
case (yyyymmdd, hhmmss, params) =>
println(yyyymmdd)
}์ ๊ท ํํ์์ ์๋์ ๊ฐ์ด ์ฌ์ฉํ ์ ์๋ค.
val params = "itemid=1234&uid=abcd1234"
val regexItemid = "itemid=[0-9]+".r
val matchOne = regexItemid.findFirstIn(params).getOrElse("").replace("itemid=", "")
println(matchOne)
val regex = "([0-9]+)".r
regex.findAllIn(params).matchData.foreach(item => println(item.group(0)))๋๋ ํจํด๋งค์นญ์ ์ด์ฉํด ๊ฐ๋จํ ์ ๊ทํํ์์ ์ฌ์ฉํ ์ ์๋ค.
val pattern_itemid_uid = "itemid=([0-9]+).*uid=([a-zA-Z0-9]+)".r
val pattern_itemid = "itemid=([0-9]+)".r
val pattern_uid = "uid=([a-z0-9]+)".r
val params = "itemid=1234&uid=abcd1234"
val (itemid, uid) = params match {
case pattern_itemid_uid(itemid, uid) => (itemid, uid)
case pattern_itemid(itemid) => (itemid, "")
case pattern_uid(uid) => ("", uid)
case _ => ("", "")
}Scala ์์๋ ์ฑ๊ธํค ๊ฐ์ฒด๋ฅผ ์ํด object ํค์๋๋ฅผ ์ฌ์ฉํ๋ค.
object Timer {
var count = 0
def currentCount(): Long = {
count += 1
count
}
}
println(Timer.currentCount())
// 1
println(Timer.currentCount())
// 2
val e = new Timer() // error์ฑ๊ธํค ๊ฐ์ฒด๋ new ๋ฅผ ์ด์ฉํด ์ธ์คํด์ค๋ก ๋ง๋ค ์ ์๋ค.
์๋ ๋ช ๋ น์ผ๋ก Scala ๋ฒ์ ์ ํ์ธํ๋ค.
$ spark-shell
......
Using Scala version 2.11.8 (OpenJDK 64-Bit Server VM, Java 1.8.0_161)
......
scala> :q์ ํ๋ก์ ํธ๋ฅผ ์์ฑํ๋ค.
$ mkdir my_project
$ cd my_project/
$ sbt
> set name := "MyProject"
> set version := "0.1"
> set scalaVersion := "2.11.8"
> session save
> exitMain.scala ์ ์๋ ๋ด์ฉ์ ์
๋ ฅํ๋ค.
$ mkdir -p src/main/scala
$ vi src/main/scala/Main.scala
# fix linter warningimport org.apache.spark.SparkContext
import org.apache.spark.SparkContext._
import org.apache.spark.SparkConf
object Main {
def main(args: Array[String]) {
val conf = new SparkConf().setAppName("HelloWorld")
val sc = new SparkContext(conf)
println("===================================")
println("Hello, world!")
println("===================================")
sc.stop()
}
}$ vi project/plugins.sbt
---------------------------------------------------------------------
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "0.14.5")
---------------------------------------------------------------------
$ vi build.sbt
......
val sparkVersion = "2.3.0"
libraryDependencies ++= Seq(
"org.apache.spark" %% "spark-core" % sparkVersion % "provided",
"org.apache.spark" %% "spark-mllib" % sparkVersion % "provided"
)
......์ปดํ์ผํ๊ณ ์คํํ๋ค.
$ sbt assembly
$ spark-submit --class Main --master local target/scala-2.11/MyProject-assembly-0.1.jar
......
===================================
Hello, world!
===================================
......์น๋ก๊ทธ๋ ๊ณต๊ฐํ ์๊ฐ ์์ด ๊ฐ์ธ์ ์ผ๋ก ๊ตฌํ๊ธฐ ๋ฐ๋๋๋ค.
Spark Shell ์ ์ด์ฉํ๋ ๋ฐฉ๋ฒ๊ณผ sbt ํด์ ์ด์ฉํ๋ ๋ฐฉ๋ฒ ์ค sbt ํด์ ์ด์ฉํ๋ ๋ฐฉ์์ผ๋ก ์งํํ๋ค.
์๋ ๋ด์ฉ์ ์์์ ์์ฑํ ํ๋ก์ ํธ ์ค Main.scala ๋ฅผ ์์ ํ๋ ๋ฐฉ์์ผ๋ก ์งํํ๋ค.
$ vi src/main/scala/Main.scala
......
val conf = new SparkConf().setAppName("HelloWorld")
val sc = new SparkContext(conf)
val log_RDD = sc.textFile("/home/ec2-user/dev/www2-www-18XXXX17.gz")
log_RDD.take(5).map(line => println(line))
sc.stop()
......์ปดํ์ผํ๊ณ ์คํํ๋ค.
$ sbt assembly
$ spark-submit --class Main --master local target/scala-2.11/MyProject-assembly-0.1.jar
#Software: Microsoft Internet Information Services 7.5
#Version: 1.0
#Date: 2018-04-19 08:00:00
#Fields: date time s-ip cs-method cs-uri-stem cs-uri-query s-port cs-username c-ip cs(User-Agent) cs(Referer) sc-status sc-substatus sc-win32-status time-taken
2018-04-19 08:00:00 110.93.XXX.83 GET /login/loginpage.asp vType=G 80 - 106.XXX.166.106 Mozilla/5.0+(Windows+NT+6.1;+WOW64;+Trident/7.0;+rv:11.0)+like+Gecko http://www.test.co.kr/ 302 0 0 0๋ก๊ทธํ์ผ์ ๋ก๊ทธ์ ์ฒซ๋ถ๋ถ์ ๋ก๊ทธํ์ผ์ ํฌ๋ฉง์ ๋ณด๊ฐ ์๋ค.
๊ณต๋ฐฑ๋ฌธ์๋ก ์ชผ๊ฐค ์ ์๊ฒ ๋์ด ์๋ค. ๋ก๊ทธํฌ๋ฉง์ ์๋ฒ์ค์ ์ ์ํด ๊ฒฐ์ ๋๋๋ฐ ์์ ๊ฒฝ์ฐ ์ด 15๊ฐ์ ํ๋๊ฐ ์๊ณ 9๋ฒ์งธ์ ํด๋ผ์ด์ธํธ ์์ดํผ๊ฐ ์๋ค. ์ด๋ฐ ์ ๋ณด๋ฅผ ๋ฐํ์ผ๋ก ํด๋ผ์ด์ธํธ ์์ดํผ๋ณ ์กฐํ๊ฑด์๋ฅผ ๊ตฌํด๋ณด์.
$ vi src/main/scala/Main.scala
......
val conf = new SparkConf().setAppName("HelloWorld")
val sc = new SparkContext(conf)
val log_RDD = sc.textFile("/home/ec2-user/dev/www2-www-18XXXX17.gz")
val filtered_log_RDD = log_RDD.map(line => line.split(" "))
.filter(line => line.size == 15)
.map(arr => (arr(8), 1))
.reduceByKey(_ + _)
.sortBy(_._2)
filtered_log_RDD.zipWithIndex()
.sortBy(_._2, ascending = false)
.collect
.foreach(row => println(row._1))
sc.stop()
......$ sbt assembly
$ spark-submit --class Main --master local target/scala-2.11/MyProject-assembly-0.1.jar
(106.XXX.166.106,1545)
(211.XXX.239.243,927)
(118.XXX.84.92,280)
(112.XXX.95.50,262)
(211.XXX.98.3,256)
(1.XXX.66.232,244)
(185.XXX.151.187,218)
(210.XXX.101.198,200)
......์๋ง์กด์ ๋ณด๋ฉด Customers who viewed this item also viewed ๋ผ๋ ์๋น์ค๊ฐ ์๋๋ฐ ๋์ผํ ๊ธฐ๋ฅ์ ๊ตฌํํด ๋ณธ๋ค.
import org.apache.spark.SparkContext
import org.apache.spark.SparkContext._
import org.apache.spark.SparkConf
import org.apache.spark.mllib.rdd.RDDFunctions._
object Main {
def main(args: Array[String]) {
val conf = new SparkConf().setAppName("HelloWorld")
val sc = new SparkContext(conf)
val log_RDD = sc.textFile("/home/ec2-user/www2-www-18XXXX15.gz")
val filtered_01_log_RDD = log_RDD.map(line => line.split(" "))
.filter(line => line.size == 15)
println("====================================================")
println("get : log")
println("====================================================")
filtered_01_log_RDD.filter(row => row(4) == "/shopping/Product.asp")
.take(1).
foreach(row => println(row.mkString(" ")))
val pattern_itemid_uid = "itemid=([0-9]+).*uid=([a-zA-Z0-9]+)".r
val pattern_itemid = "itemid=([0-9]+)".r
val pattern_uid = "uid=([a-z0-9]+)".r
val filtered_02_log_RDD = filtered_01_log_RDD.filter(row => row(4) == "/shopping/Product.asp")
.map(row => (row(0), row(1), row(5)))
.map{
case (yyyymmdd, hhmmss, params) =>
val format = new java.text.SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
val datetime = format.parse(yyyymmdd + " " + hhmmss)
val (itemid, uid) = params match {
case pattern_itemid_uid(itemid, uid) => (itemid, uid)
case pattern_itemid(itemid) => (itemid, "")
case pattern_uid(uid) => ("", uid)
case _ => ("", "")
}
(datetime.getTime(), itemid, uid)
}
println("====================================================")
println("get : time / itemid / userkey")
println("====================================================")
filtered_02_log_RDD.take(5).foreach(row => println(row.productIterator.mkString(" ")))
val filtered_03_log_RDD = filtered_02_log_RDD.filter(row => row._3 != "")
.sortBy(row => (row._3, row._1))
.sliding(2)
println("====================================================")
println("get : ((prev time / prev itemid / prev userkey), (curr time / curr itemid / curr userkey))")
println("====================================================")
filtered_03_log_RDD.take(5).foreach(row =>
println("((" + row(0).productIterator.mkString(",") + "), (" + row(1).productIterator.mkString(",") + "))")
)
val filtered_04_log_RDD = filtered_03_log_RDD.map{
case Array(x, y) => if ((x._3 == y._3) && ((y._1 - x._1) < 8 * 60 * 1000)) {
Array(x._2, y._2)
} else { Array("", "") }
}
.filter(x => x(0) != "")
println("====================================================")
println("get : prev itemid / curr itemid")
println("====================================================")
filtered_04_log_RDD.take(5).foreach(row =>
println(row(0) + " " + row(1))
)
sc.stop()
}
}
/*
====================================================
get : log
====================================================
2018-04-27 06:00:04 110.93.XXX.83 GET /shopping/Product.asp itemid=171335 80 - 175.XXX.22.75 Mozilla/5.0+(Linux;+Android+6.0.1;+LG-F700K+Build/MMB29M)+AppleWebKit/537.36+(KHTML,+like+Gecko)+Chrome/66.0.3359.126+Mobile+Safari/537.36 https://msearch.shopping.naver.com/ 302 0 0 15
====================================================
get : time / itemid / userkey
====================================================
1524808804000 171335
1524808805000 193274 F378BA44654B92108664BEDCAB4
1524808805000 148742 4831D53499EA3522AEEB4003212
1524808818000 126473
1524808821000 182321 B78862B4FEBA7805F5BDAB0EBD6
====================================================
get : ((prev time / prev itemid / prev userkey), (curr time / curr itemid / curr userkey))
====================================================
((1524810683000,136280,B24B0AF44FB9E25FD84BB82194E), (1524810927000,79182,D7ECC004F7385033A97194A60E8))
((1524810927000,79182,D7ECC004F7385033A97194A60E8), (1524812204000,141565,B134FC34326B4149296DC5695D8))
((1524812204000,141565,B134FC34326B4149296DC5695D8), (1524812077000,150778,5F9EF4C4921A750B782FE271D11))
((1524812077000,150778,5F9EF4C4921A750B782FE271D11), (1524810578000,129051,60C82C0484E80F61E84AA1F334F))
((1524810578000,129051,60C82C0484E80F61E84AA1F334F), (1524810678000,137584,8FD3448447EAD6D88A76AA554F8))
====================================================
get : prev itemid / curr itemid
====================================================
195509 195103
195103 195340
195340 195414
195414 195078
195078 195899
*/์น URL ์ userKey ๋ฅผ ๋ฃ์ด๋์๋ค๋ฉด ๊ทธ๊ฑธ ์ด์ฉํ๋ฉด ๋๊ณ , ์๋ค๋ฉด clientip ๋ฅผ ์ด์ฉํด๋ ๊ฝค ์ ์ฌํ ๊ฒฐ๊ณผ๋ฅผ ์ป์ ์ ์๋ค.
sliding() ์ ์ด์ฉํ๋ฉด ์ด์ ๋ก๊ทธ๋ผ์ธ๊ณผ ํ์ฌ ๋ก๊ทธ๋ผ์ธ์ ํ์ค๋ก ๋ง๋ค ์ ์๋ค. ๊ทธ์ค์ userKey ๊ฐ ๋์ผํ๊ณ , ์ํํ์ด์ง ์กฐํ์๊ฐ ๊ฐ๊ฒฉ์ด 8๋ถ ๋ฏธ๋ง์ธ ๋ด์ญ๋ง ๋ฝ์ผ๋ฉด ์์ ๊ฐ์ด ์ด์ ์ ์กฐํํ ์ํ์ฝ๋์ ํ์ฌ ์กฐํํ๊ณ ์๋ ์ํ์ฝ๋๋ฅผ ๊ตฌํ ์ ์๋ค.
์ฌ๊ธฐ์ ๋ค์ ์ด์ ์ํ์ฝ๋๋ก groupBy() ํ๋ฉด ํน์ ์ํ์ ์กฐํํ ๊ณ ๊ฐ์ด ๋ค์์ ์กฐํํ ์ํ, ์ฆ Customers who viewed this item also viewed ๋ฅผ ๊ตฌํ ์ ์๋ค.
์ ์์ค๋ฅผ Dataset ์ ์ด์ฉํด ๋ค์ ๊ตฌํํด ๋ณด์. Dataset ์ RDD ์ ๋นํด 2๋ฐฐ ์๋๊ฐ ๋น ๋ฅด๊ณ , ๋ฉ๋ชจ๋ฆฌ ์๋ชจ๋์ด 1/4 ์ด๋ผ๊ณ ํ๋ ์ด์ฉ ์ ์๋ ์ํฉ์ด ์๋๋ฉด Dataset ์ ์ฐ๋๋ก ํ์.
$ vi build.sbt
......
val sparkVersion = "2.3.0"
libraryDependencies ++= Seq(
"org.apache.spark" %% "spark-core" % sparkVersion % "provided"
, "org.apache.spark" %% "spark-sql" % sparkVersion % "provided"
// , "org.apache.spark" %% "spark-mllib" % sparkVersion % "provided"
, "org.apache.hadoop" % "hadoop-aws" % "2.7.6" % "provided"
)
......// -*- coding: utf-8 -*-
import org.apache.spark.SparkContext
import org.apache.spark.SparkContext._
import org.apache.spark.SparkConf
import org.apache.spark.sql._
object Main {
case class RawLog(yyyymmdd: String, hhmmss: String, params: String)
case class FilteredLog(unixtime: Long, itemid: String, uid: String)
def main(args: Array[String]) {
val conf = new SparkConf().setAppName("DS Project")
val sc = new SparkContext(conf)
val spark = SparkSession.builder()
.appName("DS Project")
.getOrCreate()
val sqlContext= new SQLContext(sc)
import sqlContext.implicits._
// spark-shell ์์ ํ
์คํธํ๋ ค๋ฉด ์๋ ๋ด์ฉ์ ์
๋ ฅํด ์ฃผ์ด์ผ ํ๋ค.
// $ vi spark/conf/spark-defaults.conf
// spark.jars.packages org.apache.hadoop:hadoop-aws:2.7.6
// ==========================================================
// S3 ์ ์์ ์ํ ์ค์ ํ๊ธฐ
val region = "ap-northeast-2"
System.setProperty("com.amazonaws.services.s3.enableV4", "true")
sc.hadoopConfiguration.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
sc.hadoopConfiguration.set("com.amazonaws.services.s3.enableV4", "true")
sc.hadoopConfiguration.set("fs.s3a.endpoint", "s3." + region + ".amazonaws.com")
// ==========================================================
// RDD ํฌ๋ฉง์ผ๋ก ๋ก๊ทธํ์ผ ์ด๊ธฐ
val rdd = sc.textFile("s3a://๋ฒํท์ด๋ฆ/ํ์ผ์ด๋ฆ")
.map(line => line.split(" "))
.filter(line => line.size == 15 && line(4) == "/shopping/Product.asp")
.map(row => RawLog(row(0), row(1), row(5)))
// ==========================================================
// DataFrame ์ผ๋ก ๋ณ๊ฒฝ
val df = rdd.toDF()
// df.show()
val pattern_itemid_uid = "itemid=([0-9]+).*uid=([a-zA-Z0-9]+)".r
val pattern_itemid = ".*itemid=([0-9]+).*".r
val pattern_uid = ".*uid=([a-z0-9]+).*".r
val format = new java.text.SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
// ==========================================================
// Dataset ์ผ๋ก ๋ณ๊ฒฝ
val ds = df.map(row => {
val yyyymmdd = row.getAs[String]("yyyymmdd")
val hhmmss = row.getAs[String]("hhmmss")
val params = row.getAs[String]("params")
val datetime = format.parse(yyyymmdd + " " + hhmmss)
val (itemid, uid) = params match {
case pattern_itemid_uid(itemid, uid) => (itemid, uid)
case pattern_itemid(itemid) => (itemid, "")
case pattern_uid(uid) => ("", uid)
case _ => ("", "")
}
val unixtime = datetime.getTime() / 1000
FilteredLog(unixtime, itemid, uid)
})
.filter(row => row.uid != "")
//ds.show()
ds.createOrReplaceTempView("tv_row_log")
val resultDS = spark.sql("""
SELECT T.prev_itemid, T.itemid, COUNT(*) as cnt
FROM
(
SELECT
unixtime, itemid
, lag(itemid) OVER (PARTITION BY uid ORDER BY unixtime) AS prev_itemid
, lag(unixtime) OVER (PARTITION BY uid ORDER BY unixtime) AS prev_unixtime
FROM tv_row_log
) T
WHERE
1 = 1
AND T.prev_itemid is not NULL
AND (unixtime - prev_unixtime) <= 8 * 60
AND itemid <> prev_itemid
GROUP BY
T.prev_itemid, T.itemid
ORDER BY
cnt desc, T.prev_itemid, T.itemid
""")
resultDS.show()
sc.stop()
}
}๊ธฐ์กด์ sliding ํจ์์ ๋นํด lag ์ด ๋ ๊น๋ํ๊ฒ ๋์ํ๋ค.
