diff options
author | Bryan Newbold <bnewbold@archive.org> | 2018-08-20 18:41:51 -0700 |
---|---|---|
committer | Bryan Newbold <bnewbold@archive.org> | 2018-08-21 21:25:56 -0700 |
commit | 39bf4b57cd552e8042bfa25565b390cb2a456ab0 (patch) | |
tree | ec5ecb68a49dbc69727747f766a02ba91bd3c3b6 /scalding/src/main/scala/sandcrawler/HBaseStatusCodeCountJob.scala | |
parent | 6aeafb083d73be8cf3296707c3e558d825202bce (diff) | |
download | sandcrawler-39bf4b57cd552e8042bfa25565b390cb2a456ab0.tar.gz sandcrawler-39bf4b57cd552e8042bfa25565b390cb2a456ab0.zip |
distinction between status_code and status counting
Diffstat (limited to 'scalding/src/main/scala/sandcrawler/HBaseStatusCodeCountJob.scala')
-rw-r--r-- | scalding/src/main/scala/sandcrawler/HBaseStatusCodeCountJob.scala | 32 |
1 files changed, 32 insertions, 0 deletions
diff --git a/scalding/src/main/scala/sandcrawler/HBaseStatusCodeCountJob.scala b/scalding/src/main/scala/sandcrawler/HBaseStatusCodeCountJob.scala new file mode 100644 index 0000000..4d9880f --- /dev/null +++ b/scalding/src/main/scala/sandcrawler/HBaseStatusCodeCountJob.scala @@ -0,0 +1,32 @@ +package sandcrawler + +import java.util.Properties + +import cascading.property.AppProps +import cascading.tuple.Fields +import com.twitter.scalding._ +import com.twitter.scalding.typed.TDsl._ +import org.apache.hadoop.hbase.io.ImmutableBytesWritable +import org.apache.hadoop.hbase.util.Bytes +import parallelai.spyglass.base.JobBase +import parallelai.spyglass.hbase.HBaseConstants.SourceMode +import parallelai.spyglass.hbase.HBasePipeConversions +import parallelai.spyglass.hbase.HBaseSource + +class HBaseStatusCodeCountJob(args: Args) extends JobBase(args) with HBasePipeConversions { + + val source = HBaseCountJob.getHBaseSource( + args("hbase-table"), + args("zookeeper-hosts"), + "grobid0:status_code") + + val statusPipe : TypedPipe[Long] = source + .read + .toTypedPipe[(ImmutableBytesWritable,ImmutableBytesWritable)]('key, 'status_code) + .map { case (key, raw_code) => Bytes.toLong(raw_code.copyBytes()) } + + statusPipe.groupBy { identity } + .size + .debug + .write(TypedTsv[(Long,Long)](args("output"))) +} |