Skip to content

skyer9/spark-scala-tutorial-ko

Folders and files

NameName
Last commit message
Last commit date

Latest commit

ย 

History

40 Commits
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 

Repository files navigation

Apache Spark Scala Tutorial For Korean

Spark

์ด ๋ฌธ์„œ๋Š” Spark ์—์„œ Scala ์–ธ์–ด๋ฅผ ์‚ฌ์šฉํ•˜๊ณ ์ž ํ•˜๋Š” ๊ฐœ๋ฐœ์ž๋ฅผ ์œ„ํ•œ ํŠœํ† ๋ฆฌ์–ผ์ด๋‹ค.

์ด ๋ฌธ์„œ๋ฅผ ์ฝ๊ธฐ ์œ„ํ•ด์„œ๋Š” Java ๋˜๋Š” Python ์— ๋Œ€ํ•œ ์ง€์‹์ด ์žˆ์–ด์•ผ ํ•˜๊ณ , Java ๋˜๋Š” Python ์„ ์ด์šฉํ•ด Spark ๋ฅผ ์ด์šฉํ•œ ๊ฒฝํ—˜์ด ์žˆ์–ด์•ผ ํ•œ๋‹ค.

๋ฌธ์„œ์˜ ์–‘์„ ์ค„์ด๊ธฐ ์œ„ํ•ด Spark ์— ๋Œ€ํ•œ ์„ค๋ช…์€ ์ƒ๋žต๋˜๋ฉฐ, ๋˜ํ•œ Scala ๋ฌธ๋ฒ•์ค‘ Spark ๊ฐœ๋ฐœ์— ๋ถˆํ•„์š”ํ•˜๋‹ค๊ณ  ํŒ๋‹จ๋˜๋Š” ๋ฌธ๋ฒ• ๋˜ํ•œ ์ƒ๋žตํ•œ๋‹ค.

์ฐธ๊ณ ์ž๋ฃŒ

์•„๋ž˜ ์ฃผ์†Œ์˜ ๋‚ด์šฉ์„ ์ฐธ๊ณ ํ•œ๋‹ค.

Scala ๊ฐœ๋ฐœํ™˜๊ฒฝ ๊ตฌ์„ฑํ•˜๊ธฐ

์ด ๋ฌธ์„œ์—์„œ๋Š” 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
$

Hello, World! ์ถœ๋ ฅํ•˜๊ธฐ

์•„๋ž˜์™€ ๊ฐ™์€ ๋ฐฉ๋ฒ•์œผ๋กœ ๊ฐ„๋‹จํ•œ ํ…Œ์ŠคํŠธ ํ”„๋กœ๊ทธ๋žจ์„ ์‹คํ–‰์‹œํ‚ฌ ์ˆ˜ ์žˆ๋‹ค.

$ mkdir sample
$ cd sample/
$ sbt console
scala> println("Hello, World!")
Hello, World!
scala> :q
$

Scala ๋ฌธ๋ฒ• ์„ค๋ช…ํ•˜๊ธฐ

์ดํ›„์— ์žˆ๋Š” ์ฝ”๋“œ๋“ค์€ sbt console ๋ช…๋ น์œผ๋กœ ์ฝ˜์†”์— ๋กœ๊ทธ์ธ๋˜์–ด ์žˆ๋Š” ๊ฒƒ์„ ์ „์ œ๋กœ ํ•œ๋‹ค.

๋ณ€์ˆ˜ ์ƒ์„ฑ(declare variable)

์•„๋ž˜์˜ ๋ฐฉ๋ฒ•์œผ๋กœ ๋ณ€์ˆ˜๋ฅผ ์ƒ์„ฑํ•  ์ˆ˜ ์žˆ๋‹ค.

val i: Int = 1
i = 2          // error

๋ณ€์ˆ˜๋ฅผ ์ƒ์„ฑํ•˜๋Š” ํ‚ค์›Œ๋“œ๋Š” val ๊ณผ var ๊ฐ€ ์žˆ๋‹ค. val ๋กœ ์ƒ์„ฑํ•œ ๋ณ€์ˆ˜๋Š” ๊ฐ’์˜ ๋ณ€๊ฒฝ์ด ๋ถˆ๊ฐ€๋Šฅํ•œ ๋ณ€์ˆ˜๊ฐ€ ๋œ๋‹ค. ๋ฐ˜๋ฉด์— var ๋กœ ์ƒ์„ฑํ•œ ๋ณ€์ˆ˜๋Š” ๊ฐ’์˜ ๋ณ€๊ฒฝ์ด ๊ฐ€๋Šฅํ•˜๋‹ค. ํ•˜์ง€๋งŒ, Scala ์—์„œ๋Š” var ๋ฅผ ์‚ฌ์šฉํ•˜์ง€ ์•Š์„ ๊ฒƒ์„ ๊ถŒ์žฅํ•˜๊ณ  ์žˆ๋‹ค.

val j = 2

Scala ์—์„œ๋Š” ๋ณ€์ˆ˜๊ฐ’์˜ ํƒ€์ž…์„ ์•Œ ์ˆ˜ ์žˆ๋Š” ๊ฒฝ์šฐ(์œ„์—์„œ 2 ๋Š” Int) ์œ„์™€๊ฐ™์ด ๋ณ€์ˆ˜ํƒ€์ž…(: Int)์„ ์ƒ๋žตํ•  ์ˆ˜ ์žˆ๋‹ค.

val k = 3
println("class: " + k.getClass)
// class: int

Scala ์—์„œ๋Š” ๋ชจ๋“  ๋ณ€์ˆ˜๋Š” ๊ฐ์ฒด๋‹ค. ์œ„์—์„œ k ๋Š” ๋‹จ์ˆœํ•œ ์ •์ˆ˜๊ฐ’์ด ์•„๋‹Œ ์ •์ˆ˜ํ˜• ๊ฐ์ฒด๊ฐ€ ๋œ๋‹ค.

ํ•จ์ˆ˜ ์ƒ์„ฑ(declare function)

์•„๋ž˜์™€ ๊ฐ™์ด ํ•จ์ˆ˜๋ฅผ ์ƒ์„ฑํ•  ์ˆ˜ ์žˆ๋‹ค.

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))

ํด๋ž˜์Šค ์ƒ์„ฑ(declare class)

์•„๋ž˜ ์ฝ”๋“œ๋ฅผ 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 ๋กœ ์ •์˜ํ•œ๋‹ค.

ํด๋ž˜์Šค ์ƒ์„ฑ์ž(class constructor)

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)       // ok

get,set ๋ฉ”์†Œ๋“œ๋Š” ์ž๋™์œผ๋กœ ์ถ”๊ฐ€๋˜๋ฏ€๋กœ ๋ณ„๋„๋กœ ์ž‘์—…ํ•  ํ•„์š”๊ฐ€ ์—†๋‹ค. ๋˜ํ•œ, ๋‹ค๋ฅธ ์–ธ์–ด์™€ ๋‹ค๋ฅด๊ฒŒ ๋งด๋ฒ„๋ณ€์ˆ˜ ๋ฐ ๋งด๋ฒ„ํ•จ์ˆ˜๊ฐ€ private ์„ ๋ณ„๋„๋กœ ์ง€์ •ํ•ด ์ฃผ์ง€ ์•Š๋Š” ํ•œ, public ์ด ๋””ํดํŠธ๋กœ ์ง€์ •๋œ๋‹ค.

ํŒจํ„ด ๋งค์นญ(switch case statment)

์ผ๋ฐ˜์ ์ธ 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)

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)

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)

-> ๋ฅผ ์ด์šฉํ•ด ํŠœํ”Œ์„ ์ƒ์„ฑํ•  ์ˆ˜ ์žˆ๋‹ค.

๋งต(Map)

key-value ํ˜•ํƒœ์˜ ๊ฐ’์˜ ๋ฌถ์Œ์ด Map ์ด๋‹ค.

val m = Map(1 -> "one", 2 -> "two")
println(m(2))

์œ„์—์„œ -> ๋Š” ํŠน๋ณ„ํ•œ ๋ฌธ๋ฒ•์ด ์•„๋‹ˆ๊ณ  ํŠœํ”Œ์˜ ์ƒ์„ฑ์— ๋ถˆ๊ณผํ•˜๋‹ค. ์œ„์—์„œ ์ƒ์„ฑ๋œ ๋งต์€ ์‹ค์ œ๋กœ๋Š” Map((1, "one"), (2, "two")) ์˜ ํ˜•ํƒœ๊ฐ€ ๋˜๊ณ , ๋งต์— ๋“ค์–ด์žˆ๋Š” ๋ฐ์ดํƒ€๋Š” ์ฒซ๋ฒˆ์งธ ๊ฐ’์ด key ๊ฐ€ ๋˜๊ณ , ๋‘๋ฒˆ์งธ ๊ฐ’์ด value ๊ฐ€ ๋œ๋‹ค.

ํ•จ์ˆ˜ ์กฐํ•ฉ(function combinator)

๋ฆฌ์ŠคํŠธ๋ฅผ ์ „๋‹ฌ๋ฐ›์•„ ์ผ์ •ํ•œ ์ฒ˜๋ฆฌ๋ฅผ ํ•˜๊ณ  ์ฒ˜๋ฆฌ๋œ ๊ฐ’์„ ์ „๋‹ฌํ•ด์ฃผ๋Š” ๊ฒƒ์„ ํ•ฉ์ˆ˜์กฐํ•ฉ์ด๋ผ๊ณ  ํ•œ๋‹ค.

map()

๋‹ค๋ฅธ ์–ธ์–ด์—์„œ๋Š” 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 ์—์„œ๋Š” ํŒŒ๋ผ๋ฏธํ„ฐ๊ฐ€ ํ•œ๊ฐœ์ผ ๊ฒฝ์šฐ ๊ด„ํ˜ธ๋ฅผ ์ƒ๋žตํ•  ์ˆ˜ ์žˆ๋‹ค.

foreach()

map() ์ด ์ž…๋ ฅ๋œ ๋ฐ์ดํƒ€๋ฅผ ๊ทธ๋Œ€๋กœ ๋‘๊ณ  ๋ณ€ํ˜•๋œ ๋ฐ์ดํƒ€๋ฅผ ๋ฐ˜ํ™˜ํ•˜๋Š”๊ฒƒ๊ณผ ๋‹ค๋ฅด๊ฒŒ, foreach() ๋Š” ์ž…๋ ฅ๋œ ๊ฐ’ ์ž์ฒด๋ฅผ ๋ณ€ํ™˜ํ•˜๊ณ  ๋ฆฌํ„ด๊ฐ’์ด ์—†๋‹ค.

val numbers = List(1, 2, 3, 4)
numbers.foreach(i => println(i))

foreach() ์—๊ฒŒ ๋ฆฌํ„ด๊ฐ’์„ ์š”์ฒญํ•˜๋ฉด Unit(๋‹ค๋ฅธ ์–ธ์–ด์—์„œ๋Š” void) ์ด ๋ฐ˜ํ™˜๋œ๋‹ค.

filter()

์ž…๋ ฅ๋œ ๊ฐ’์„ ํ•„ํ„ฐ๋งํ•ด์„œ ๊ฐ’์ด ์ฐธ์ธ ๊ฒƒ๋“ค๋กœ ์ด๋ฃจ์–ด์ง„ ๋ฆฌ์ŠคํŠธ๋ฅผ ๋ฐ˜ํ™˜ํ•œ๋‹ค.

val numbers = List(1, 2, 3, 4)
println(numbers.filter(i => i % 2 == 0))

zip()

๋‘๊ฐœ์˜ ๋ฆฌ์ŠคํŠธ๋ฅผ ๊ฐ๊ฐ์˜ ๋ฐ์ดํƒ€๋ฅผ ๋ฌถ์–ด ํŠœํ”Œ ๋ฆฌ์ŠคํŠธ๋กœ ๋งŒ๋“ ๋‹ค.

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()

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()

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()

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(), dropWhile()

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()

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()

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()

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()

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() ํ•œ๋‹ค.

ํ•จ์ˆ˜ ์กฐํ•ฉ(function combinator) ๊ณผ ํŒจํ„ด๋งค์นญ

ํ•จ์ˆ˜์กฐํ•ฉ๊ณผ ํŒจํ„ด๋งค์นญ์„ ํ•จ๊ป˜ ์‚ฌ์šฉํ•˜๋ฉด ์•„๋ž˜์™€ ๊ฐ™์ด ์ฝ”๋“œ๋ฅผ ์ž‘์„ฑํ•  ์ˆ˜ ์žˆ๋‹ค.

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)
}

์ •๊ทœํ‘œํ˜„์‹(Regular Expressions)

์ •๊ทœ ํ‘œํ˜„์‹์€ ์•„๋ž˜์™€ ๊ฐ™์ด ์‚ฌ์šฉํ•  ์ˆ˜ ์žˆ๋‹ค.

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 _ => ("", "")
}

์‹ฑ๊ธ€ํ†ค ๊ฐ์ฒด(Singleton Class, Static Object)

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 ํ”„๋กœ์ ํŠธ ์ƒ์„ฑํ•˜๊ธฐ

์ƒˆ ํ”„๋กœ์ ํŠธ ์ƒ์„ฑํ•˜๊ธฐ

์•„๋ž˜ ๋ช…๋ น์œผ๋กœ 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
> exit

Main.scala ์— ์•„๋ž˜ ๋‚ด์šฉ์„ ์ž…๋ ฅํ•œ๋‹ค.

$ mkdir -p src/main/scala
$ vi src/main/scala/Main.scala
# fix linter warning
import 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!
===================================
......

Scala Spark Example With Web Log

์›น๋กœ๊ทธ๋Š” ๊ณต๊ฐœํ•  ์ˆ˜๊ฐ€ ์—†์–ด ๊ฐœ์ธ์ ์œผ๋กœ ๊ตฌํ•˜๊ธฐ ๋ฐ”๋ž๋‹ˆ๋‹ค.

Spark Shell ์„ ์ด์šฉํ•˜๋Š” ๋ฐฉ๋ฒ•๊ณผ sbt ํˆด์„ ์ด์šฉํ•˜๋Š” ๋ฐฉ๋ฒ• ์ค‘ sbt ํˆด์„ ์ด์šฉํ•˜๋Š” ๋ฐฉ์‹์œผ๋กœ ์ง„ํ–‰ํ•œ๋‹ค.

์•„๋ž˜ ๋‚ด์šฉ์€ ์œ„์—์„œ ์ƒ์„ฑํ•œ ํ”„๋กœ์ ํŠธ ์ค‘ Main.scala ๋ฅผ ์ˆ˜์ •ํ•˜๋Š” ๋ฐฉ์‹์œผ๋กœ ์ง„ํ–‰ํ•œ๋‹ค.

์›น๋กœ๊ทธ์—์„œ 5๋ผ์ธ๋งŒ ์ถœ๋ ฅํ•˜๊ธฐ(with RDD)

$ 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

๋กœ๊ทธํŒŒ์ผ์€ ๋กœ๊ทธ์˜ ์ฒซ๋ถ€๋ถ„์— ๋กœ๊ทธํŒŒ์ผ์˜ ํฌ๋ฉง์ •๋ณด๊ฐ€ ์žˆ๋‹ค.

๊ฐ€์žฅ ๋งŽ์ด ์ ‘์†ํ•œ ํด๋ผ์ด์–ธํŠธ ์•„์ดํ”ผ ๊ตฌํ•˜๊ธฐ(with RDD)

๊ณต๋ฐฑ๋ฌธ์ž๋กœ ์ชผ๊ฐค ์ˆ˜ ์žˆ๊ฒŒ ๋˜์–ด ์žˆ๋‹ค. ๋กœ๊ทธํฌ๋ฉง์€ ์„œ๋ฒ„์„ค์ •์— ์˜ํ•ด ๊ฒฐ์ •๋˜๋Š”๋ฐ ์œ„์˜ ๊ฒฝ์šฐ ์ด 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)
......

์ด์ƒํ’ˆ์„ ์กฐํšŒํ•œ ๊ณ ๊ฐ์ด ์กฐํšŒํ•œ ๋‹ค๋ฅธ ์ƒํ’ˆ ๊ตฌํ•˜๊ธฐ(with RDD)

์•„๋งˆ์กด์— ๋ณด๋ฉด 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 ๋ฅผ ๊ตฌํ•  ์ˆ˜ ์žˆ๋‹ค.

์ด์ƒํ’ˆ์„ ์กฐํšŒํ•œ ๊ณ ๊ฐ์ด ์กฐํšŒํ•œ ๋‹ค๋ฅธ ์ƒํ’ˆ ๊ตฌํ•˜๊ธฐ(with Dataset)

์œ„ ์†Œ์Šค๋ฅผ 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 ์ด ๋” ๊น”๋”ํ•˜๊ฒŒ ๋™์ž‘ํ•œ๋‹ค.

๋”๋ณด๊ธฐ

RDD, DataFrame, DataSet ์˜ ์ฐจ์ด์ 

About

Tutorial for Scala on Spark only

Resources

License

Stars

12 stars

Watchers

5 watching

Forks

Releases

No releases published

Packages

 
 
 

Contributors