3 ms·
Single process "streaming" or basically CEP engines have been around for a very long time and used to be the norm before distributed stream processing engines c
by cangencer 8y ago
Single process "streaming" or basically CEP engines have been around for a very long time and used to be the norm before distributed stream processing engines came around (such as Storm). CEP engines have much richer functionality than distributed stream processing engines because they don't need to deal with partitioning or data distribution. I'm not quite sure I see the appeal of making a non-distributed engine but with the same limitations of a distributed engine.
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.
[1] https://github.com/hazelcast/hazelcast-jet https://github.com/hazelcast/hazelcast-jet
- jonaf 8y agoI couldn't find this on your product page, but does Hazelcast Jet support "global windows"? I've found 90% of stream processing systems I've come across only work within a window, but if I need to perform a computation "for all time," I'm SOL.
- capkutay 8y agoStream processing engines are meant to work with ephemeral data, fast data. Couldn't you just use a database to query 'all time'?
- cangencer 8y agoWe are adding something called a "rolling aggregation", where you receive a record, accumulate it and then emit the current accumulated value. I'm not sure if this matches what you want.
- grt17 8y agoWhich are the limitations shared by both a single node and a distributed engine? I am confused... 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.
- cangencer 8y agoI meant that products like Esper, StreamBase, InfoSphere have been around for a long time, which have a _very_ rich set of features [1], and are mostly designed around single process usage. Lot of the type of queries they support are not possible to implement in a performant way in a distributed system. Though nowadays Esper claim to have horizontal scalability - it was originally designed as a single threaded system. They do also have a passive/active type solution as you mentioned. 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. [1] http://www.espertech.com/esper/ http://www.espertech.com/esper/
- scott_s 8y agoI assume that by "InfoSphere" you mean what used to be called IBM InfoSphere Streams, but what is now just IBM Streams. I do research and development on IBM Streams (see a sibling comment), and I can say with certainty that it was designed as a distributed, parallel stream processing system from the start. It is a more general computing platform than CEP engines - in fact, we actually implemented a CEP engine as a Streams operator. Product documentation for the operator: https://www.ibm.com/support/knowledgecenter/SSCRJU_4.2.0/com.ibm.streams.toolkits.doc/spldoc/dita/tk$com.ibm.streams.cep/op$com.ibm.streams.cep$MatchRegex.html 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 http://hirzels.com/martin/papers/debs12-cep.pdf