Do We Need Distributed Stream Processing?
lsds.doc.ic.ac.uk
lsds.doc.ic.ac.uk
A couple of great posts by Frank McSherry along the same lines a few years ago, doing page rank for a 128bil edge graph on a single laptop (in the second post below) -
The first post - "scalability! But at what cost" http://www.frankmcsherry.org/graph/scalability/cost/2015/01/...
The second post - "bigger data, same laptop" http://www.frankmcsherry.org/graph/scalability/cost/2015/02/...
Think ancient tools like grep versus Hadoop or Spark. Or our word2vec in Gensim vs word2vec in Spark or Tensorflow.
The occasional recalibration between "OMG-big-O-toward-infinity!-need-scale" versus "our-data-is-actually-finite-what's-a-good-operational-compromise" is a worthwhile exercise.
http://yellowstone.cs.ucla.edu/~yang/paper/sigmod2016-p958.p...
I wouldn't hold your breath on it getting faster.
Edit: and if you are interested in state-of-the-art stuff, here is a thread to start you out on the papers currently being presented at SIGMOD2018:
https://twitter.com/frankmcsherry/status/1002812204075364353
I work on IBM Streams, which is a distributed stream processing system [1,2]. When the cost of communicating between hosts is greater than the cost of your computation, distributed processing is not going to help. When your computation is basically a form of word count, then, yes, using multiple hosts may not help. But if your computation is, say, online speech to text for a large call center [3], then the application will eat just about as much compute resources as you can throw at it. That's going to mean using multiple machines, even if each machine is quite large.
As a plug for our system: we scale down as well as up. If you only have a single host, we will run the entire application in a single process, using threads as needed to achieve parallelism using either dedicate threads or a lightweight operator scheduler [4]. Users can configure which without changing application logic. If you have multiple machines, we will automatically partition your application across those machines for you [5]. By default, we do this all for you, without needing any input. But if you want to explicitly control what executes where, and using how many host machines, that's easy to configure, also without changing application logic.
[1] https://www.ibm.com/cloud/streaming-analytics
[2] http://ibmstreams.github.io/
[3] https://www.ibm.com/case-studies/verizon
[4] Low-Synchronization, Mostly Lock-Free, Elastic Scheduling for Streaming Runtimes. Scott Schneider and Kun-Lung Wu. PLDI 2017. http://www.scott-a-s.com/files/pldi2017_lf_elastic_schedulin...
[5] Automatic Fusion and Threading. https://developer.ibm.com/streamsdev/docs/automatic-fusion-a...
That's a good thing
I work on Hazelcast Jet [1], which is a Java based distributed stream processing engine. The core engine is fast enough that it can be used with very good throughput on a single node (several times faster compared to Flink or Spark) but usually several nodes are not only needed strictly for parallelization but also for tolerating node failures and being able to restart where you left off. As others have pointed out, not every computation can be parallelised efficiently. Jet also offers in memory storage, so adding more nodes also increases your storage capacity.
Since the core of Jet is small enough (~400kb JAR), we also considered making a non-distributed version that runs strictly in process. Mainly for lightweight usage or embedding but would also offer a path to distributed execution, if it was ever needed.
There are ways to achieve fault-tolerance, even for a centralised system (e.g, by maintaining an active/passive replica). In cases like this, where there is such a great gap between the performance of current systems, you can always "waste" another node(s) for fault-tolerance and still operate with less cost, if you want.
I agree that some types of computations (e.g., multiple/distributed sources) may not be benefited by a centralised approach and I definitely don't claim that this is a solution for everything. However, the point is to criticise the design choices that we make for a streaming system. Streaming support, even for popular systems today, is something like an extension (sometimes a hack) on top of the core of the system, hidden beneath multiple layers of abstraction. In addition, modern systems try to do many things at the same time (support AI, batch & stream processing, connectors to publish-subscribe systems, multiple wire protocols) and they end up doing most of them poorly.
Stream processing frameworks originally evolved to offer "big scale" through data partitioning compared to the traditional CEP systems. But CEP engines have been able to deal with windowing and similar concepts since many years ago - the main difference of the stream processing frameworks _is_ the distribution and scalability aspect.
My point was that the systems linked in the original article seem to match closely to the limitations of what distributed stream processing frameworks are able to do, but only run on a single node.
Product documentation for the operator: https://www.ibm.com/support/knowledgecenter/SSCRJU_4.2.0/com...
Academic paper: Partition and Compose: Parallel Complex Event Processing. Martin Hirzel. DEBS 2012. http://hirzels.com/martin/papers/debs12-cep.pdf
It's a bit of a chicken and egg problem, since when designing new systems, management asks "How fast can we run this?", to which engineers reply "How much data is there and what speed are we aiming at?".
At this point, a failure to specify the desired volume and latencies leads to a sad cycle where scales are overblown ("just in case") and a complex solution chosen. Distributed solutions necessarily come with a massive overhead, so it is later discovered that even larger systems are needed…
Many such anecdotes circulate the industry, and I can add our own: In 2012 we implemented SVD, a math algorithm at the core of many machine learning techniques like PCA, LSI etc. This was a fast streamed "local" SVD implementation in Gensim (and now even faster in ScaleText). Suddenly, use-cases that needed a cluster of 12 beefy Hadoop (Mahout) machines could be performed on a single laptop, faster, and after a single `pip install`.
Our perf testing of word2vec (yet another streamed ML algo) showed a similar ROI pattern of local-vs-generic: https://rare-technologies.com/machine-learning-hardware-benc... (note the Spark graph at the end).
In any case, there are two ways to distribute stream processing tasks:
* [~Heavyweight] Use a cluster with multiple runners like [Spark Streaming](https://spark.apache.org/streaming/) or [Flink](https://flink.apache.org/)
* [~Lightweight] Do stream processing in applications, at the edge, on gateways or devices like [Bistro Streams](https://github.com/asavinov/bistro/tree/master/server) or [Kafka Streams](https://kafka.apache.org/documentation/streams/)
Normally, distributed stream processing requires also partitioning the data, e.g., by user or session ids, so that these isolated streams can be processed independently at different nodes.
There is a clear value in making Stream processors work closer to hardware. That does not mean single node can handle does better under all use cases. When the complexity of queries increases and when queries can be portioned, there are use cases where still distributed setup can do better.
I work on WSO2 Stream Processor, https://wso2.com/analytics ( Opensource, Apache Licensed). We have chosen to support both worlds where SP has an HA mode that can two servers in a Hot-warm mode that do 100K event per the second and give you a two-node deployment. If you want more, SP runs on top of Kafka, scales, and support multi-data center deployments as well. If the user is in doubt, he can start small and later switch to Kafka without changing any code.
Of course these streams are more rigid than the kind the article discusses, but the tradeoffs are similar, and anyone architecting streams should know about the state of the art in database distribution.
I agree that the benchmark used here is not a good example of a real-world streaming application, as it's not even compute-intensive. However, it is still used as the main benchmark for the evaluation of streaming frameworks.
In this blogpost, we are trying to start a discussion about how modern streaming systems perform and how they are supposed to perform. Towards this direction, we should reconsider what we believe is regular. For example, in cases like this, if you can use a single node instead of 5 for your computations, why shouldn't you consider about it...
The Yahoo Streaming Benchmark isn't really much of a streaming benchmark, its a "reading things off of kakfa, deserializing JSON, sticking things in Redis benchmark". Really not a good benchmark for understanding the system at test.
> Databricks made a few modifications to the original benchmark, all of which are explained in their own post: > Removing Redis from step 5 > Generating data in memory for Flink and Spark; generating data via Spark and writing it to Kafka for Kafka Streams > Writing data out to Kafka instead of to Redis
I think the authors know that the benchmark is weak. The blog post isn't (to my reading) claiming that it is proof that SABER is great so much as that the distributed stream processors need better evidence that they should exist (vs the Databricks Structured Streaming SIGMOD 2018 paper, in which the YSB benchmark is (I think) the only quantitative evaluation they do).
[0]: https://www.usenix.org/system/files/conference/osdi12/osdi12...
streamr seems to be for low rate messages (0.001 to 1 per second) between multiple, disjoint users, using cloud services on the blockchain.
This post is about high rated messages (80,000,000 per second) within one organization using self-hosted hardware.