This repository was archived by the owner on Jan 20, 2022. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 256
Home
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