diff options
Diffstat (limited to 'scald-mvp/src/main')
| -rw-r--r-- | scald-mvp/src/main/scala/example/WordCount.scala | 12 | ||||
| -rw-r--r-- | scald-mvp/src/main/scala/example/WordCountJob.scala | 12 | 
2 files changed, 12 insertions, 12 deletions
| diff --git a/scald-mvp/src/main/scala/example/WordCount.scala b/scald-mvp/src/main/scala/example/WordCount.scala deleted file mode 100644 index 0de6ae0..0000000 --- a/scald-mvp/src/main/scala/example/WordCount.scala +++ /dev/null @@ -1,12 +0,0 @@ - -package example - -import com.twitter.scalding._ - -class WordCount(args : Args) extends Job(args) { -  TypedPipe.from(TextLine(args("input"))) -    .flatMap { line => line.split("""\s+""") } -    .groupBy { word => word } -    .size -    .write(TypedTsv(args("output"))) -} diff --git a/scald-mvp/src/main/scala/example/WordCountJob.scala b/scald-mvp/src/main/scala/example/WordCountJob.scala new file mode 100644 index 0000000..83a8dd0 --- /dev/null +++ b/scald-mvp/src/main/scala/example/WordCountJob.scala @@ -0,0 +1,12 @@ +package com.twitter.scalding.examples + +import com.twitter.scalding._ + +class WordCountJob(args: Args) extends Job(args) { +  TypedPipe.from(TextLine(args("input"))) +    .flatMap { line => line.split("\\s+") } +    .map { word => (word, 1L) } +    .sumByKey +    // The compiler will enforce the type coming out of the sumByKey is the same as the type we have for our sink +    .write(TypedTsv[(String, Long)](args("output"))) +} | 
