Skip to content
This repository was archived by the owner on Jan 20, 2022. It is now read-only.
johnynek edited this page Dec 11, 2014 · 29 revisions
Summingbird is a library that lets you write MapReduce programs that look like native Scala or Java collection transformations. So, while a word-counting aggregation in pure Scala might look like this: #+BEGIN_SRC scala def wordCount(source: Iterable[String], store: MutableMap[String, Long]) = source.flatMap { sentence => toWords(sentence).map(_ -> 1L) }.foreach { case (k, v) => store.update(k, store.get(k) + v) } #+END_SRC Counting words in Summingbird looks like this: #+BEGIN_SRC scala def wordCount[P <: Platform[P]] (source: Producer[P, String], store: P#Store[String, Long]) = source.flatMap { sentence => toWords(sentence).map(_ -> 1L) }.sumByKey(store) #+END_SRC The logic is exactly the same, and the code is almost the same. The main difference is that you can execute the Summingbird program in "batch mode" (using [[https://github.com/twitter/scalding][Scalding]]), in "realtime mode" (using [[https://github.com/nathanmarz/storm][Storm]]), or on both Scalding and Storm in a hybrid batch/realtime mode that offers your application very attractive fault-tolerance properties. Summingbird provides you with the primitives you need to build rock solid production systems. *** Around the Web - [[https://github.com/upio/summingbird-hybrid-example][Example hybrid summingbird github project]]. - [[https://blog.twitter.com/2013/streaming-mapreduce-with-summingbird][Open Sourcing Announcement, Twitter Eng Blog]] - [[http://www.youtube.com/watch?v=23scdoxHOLg&feature=youtu.be][Introduction to Summingbird]] (by [[https://twitter.com/posco][@posco]]) - [[https://github.com/sritchie/summingbird-workshop][LambdaJam 2013 Summingbird Workshop]] (by [[https://twitter.com/sritchie][@sritchie]]) - [[http://www.youtube.com/watch?v=Y3PETLJeP7o][Summingbird: StreamingMapReduce at Twitter]] (Sam's talk from the AK Data Science Summit) - [[https://speakerdeck.com/sritchie/summingbird-at-cufp][CUFP 2013: Realtime MapReduce at Twitter]] (by [[https://twitter.com/sritchie][@sritchie]]) - [[http://vimeo.com/75516079][Boston Storm Users: Summingbird, Scala & Storm]] ([[https://speakerdeck.com/sritchie/boston-storm-users-summingbird-scala-and-storm][slides]]) - [[http://www.youtube.com/watch?v=iuvauJZaMqA][PNW Scala 2013: Taking Hadoop Realtime with Summingbird]] - [[https://speakerdeck.com/sritchie/the-road-to-summingbird-stream-processing-at-every-scale][The Road to Summingbird: Stream Processing at (Every) Scale]] *** Future Plans We're very excited about Summingbird's development as we move beyond our initial release. Some of our future plans include: - Support for more execution platforms ([[https://github.com/mesos/spark][Spark]] and [[http://akka.io/][Akka]] seem like clear candidates) - Pluggable optimizations for the =Producer= graph layer - Projection and filter pushdown - Support for filter-aware data sources, like [[http://parquet.io/][Parquet]] - Libraries of higher-level mathematics and machine learning code on top of Summingbird's =Producer= primitives - More extensions to Summingbird's related projects (listed below) - More data structures with =Monoid= instances via [[https://github.com/twitter/algebird][Algebird]] - More key-value stores implementations via [[https://github.com/twitter/storehaus][Storehaus]] - More Storm data sources, via [[https://github.com/twitter/tormenta][Tormenta]] - More tutorials and examples with public data sources More information on these issues and many more can be found on the [[https://github.com/twitter/summingbird/issues][Summingbird issue tracker]]. *** Getting Involved Discussion about development and new features occurs primarily on the Summingbird mailing list. To join the mailing list, email <mailto:summingbird@librelist.com>. The same address, summingbird@librelist.com, is used for posting once you've joined. Issues should be reported on the [[https://github.com/twitter/summingbird/issues][GitHub issue tracker]]. Simpler issues appropriate for first-time contributors looking to help out are tagged "newbie". To see which issues are slated for a planned release, check out the [[https://github.com/twitter/summingbird/issues/milestones][milestones]] page. Follow [[https://twitter.com/summingbird][@summingbird]] on Twitter for updates. Do you use Summingbird? Please contact us so we can add you to the [[https://github.com/twitter/summingbird/wiki/Powered-by][Powered By]] wiki page. *** Related Projects The Summingbird projects spawned a number of related subprojects, notably: - [[https://github.com/twitter/algebird][Algebird]] Algebird is an abstract algebra library for Scala. Many of the data structures included in Algebird have Monoid implementations, making them ideal to use as values in Summingbird aggregations. - [[https://github.com/twitter/bijection][Bijection]] Summingbird uses the Bijection project's =Injection= typeclass to share serialization between different execution platforms and clients. - [[https://github.com/twitter/chill][Chill]] Summingbird's Storm and Scalding platforms both use the [[https://github.com/EsotericSoftware/kryo][Kryo]] library for serialization. Chill augments Kryo with a number of helpful configuration options, and provides modules for use with Storm, Scala, Hadoop. Chill is also used by the Berkeley Amp Lab's [[http://spark.incubator.apache.org/][Spark]] project. - [[https://github.com/twitter/tormenta][Tormenta]] Tormenta provides a type-safe layer over Storm's =Scheme= and =Spout= interfaces. - [[https://github.com/twitter/storehaus][Storehaus]] Summingbird's client is implemented using Storehaus's async key-value store traits. The Storm platform makes use of Storehaus's =MergeableStore= trait to perform real-time aggregations into a number of commonly used backing stores, including Memcached and Redis. *** Documentation TODO: - Testing - Subpackage Descriptions - Picking Monoids - Long-Term Serialization - The Storm Platform - The Scalding Platform

Clone this wiki locally