Apache Heron: A realtime, distributed, fault-tolerant stream processing engine
heron.incubator.apache.org
heron.incubator.apache.org
> https://www.youtube.com/watch?v=GFfLCGW_5-w
and another on event log architecture:
> https://www.youtube.com/watch?v=RlwO6CJbJjQ
After digging through all of this material and playing around with LMAX Disruptor & Raft, I have been able to develop a really good understanding of how to build these sorts of systems on my own. Fun constraints like "Only one thread actually mutates anything, and its the same one over and over" make for incredibly elegant implementation opportunities. Not having to constantly hunt down exotic thread-safe data structures means that you can focus on building actual value.
Latency is the biggest devil you will dance with in this arena, so almost everything you do will be oriented around mitigating that effect. Latency both at the network and inside the CPU/memory/storage. It applies at every level.
Sounds great in theory and usually results in fantastic average case, but you get a hot partition and suddenly you can't work share and things go south. It's not a perfect solution.
The trick is to minimize the raw amount of bytes that must enter into that synchronous context. Maybe the account disclosure PDFs can be referred to by some GUID token in AWS S3, but the actual decimal account balance/transaction facts should be included in the event data stream.
One other option is to talk to the business and see if you can break their 1 gigantic synchronous context into multiple smaller ones that can progress independently.
I think most popular options for high-volume self-hosted distributed stream processing solutions are still Spark, Flink, and Kafka Streams.
Kafka streams is simpler, as it is basically just a framework on top of Kafka itself, so if you already use Kafka for streaming data and don't have complex needs, it might be a good option.
Spark and Flink are similar. Both support both batch processing (on top of Hadoop, for example) and stream processing. Spark has better tooling, but Flink has more sophisticated support for streaming window functions. Spark also uses "micro-batches" instead of being truly real-time, so there will be a bit more latency when doing streaming with Spark, if that matters.
--
Another interesting project is Beam, which provides a unified way of writing jobs that can then be run on different engines that support it (both Flink and Spark do, as well as Dataflow on Google).
Apache hosts a lot of projects in this category. Most (like Storm) I would probably not pick up for a greenfield project today. Also, these things come with some significant operational overhead, so make sure you really need them. Stream processing at scale is hard. The compelling use case for these things is when you need to do window aggregations on a lot of streaming data and get results in real-time.
I have no relationship with Databricks.
At this point, all of these frameworks would probably benefit from a flowchart (or questionnaire tool) that can guide someone towards an informed decision. "Do you need redundancy?" - "Can you afford to lose some messages in situation XYZ?" - "How many events/sec do you want to process?" - "How much hardware can you throw at the problem?" etc.
For eg, storm or samza might check all the boxes, but the design of the system is poor enough that the performance will suck. For older versions of Storm, you should be able to write a multithreaded app on a single machine that outperforms a storm cluster.
I still don't get why they use VPN terms for event broker
Solace's first product was hardware appliances which are still used for high throughput and low latency usecases. Concept of VPN was used to set up isolated virtual brokers so different teams can have their own environments on a shared hardware appliance.
The concept was ported over to software as well and is extremely useful in an enterprise environment. It allows different teams to have their own virtual brokers but not have to pay for or manage multiple brokers.
That recontextualizes your previous post quite a bit...
On the other hand databases have a huge lock-in power (been trying to strangle an Oracle myself for over 5 years now). It is lucrative to be in the database business.
I'd say that every project with a claim of some improvement can, and should try to establish itself in the market; and that having it join the Apache foundation is a great way to get some brand recognition on the cheap.
----
Also, Heron is not that new. It has been developed at Twitter, for replacing Storm IIRC.
(I could frame it in a more glass half full way, but I find that the pessimistic way of looking at it helps a lot with trimming down options when you have far, far, far too many options.)
https://www.usenix.org/system/files/conference/hotos15/hotos...
I guess there are enough users to keep the communities alive?
In event processing there's a continuum of expected latencies from batch processing to realtime. Batch processing is typically running reports over large volume of events (good for throughput). Hadoop is a good example. On the other end, sub-second realtime report is possible with Heron/Storm. Spark is kind of in the middle with hybrid mini-batching. Reportedly Twitter has used Heron/Storm to track word counts in all the tweets to find trending topics, where the latency between a new tweet coming in to the word counts updated over the whole network is in 100s milliseconds.
...then some of the employees involved moved to Yahoo, built it again, and then open sourced it as Pulsar. Then moved onto form a Confluent v2 to sell their Kafka v2 (now with even more Zookeeper!)
Based on a quick Wikipedia skim, looks like the answer is "yes". That explains this bullet point:
> Heron is API compatible with Apache Storm and hence no code change is required for migration.
From the beginning, Heron was envisioned as a new kind of stream processing system, built to meet the most demanding of technological requirements, to handle even the most massive of workloads, and to meet the needs of organizations of all sizes and degrees of complexity. Amongst these requirements:
The ability to process billions of events per minute
Extremely low end-to-end latency
Predictable behavior regardless of scale and in the face of issue like extreme traffic spikes and pipeline congestion
Simple administration, including:
The ability to deploy on shared infrastructure
Powerful monitoring capabilities
Fine-grained configurability
Easy debuggability
I can't wait for my next startup interview. We have a requirement of 25 messages per hour, with 10KB per message, you think you can build the ingestion pipeline using Kafka and MongoDB on a 10 node M5d.24xlarge cluster?
They specifically want Kafka, there's no real reason other than they need a queue, which Kafka actively states that it's not. At that point it gets really tricky to reason with the developers about why they might be better served by something else. Generally speaking it's not much of an issue, because Kafka will deal with workloads just fine, it's just weird. I have seen one customer use Kafka as a database, that works less well.
We do see the same with Kubernetes. The developers pick Kubernetes and at that point it's to late. They specifically want Kubernetes even if you could more easily solve the problem with Nomad, Docker-Compose, plain old VMs or EC2, depending on the problem.
Confluent itself says that for a workload of 38Mbps/sec rabbitmq has 5 times less latency than Kafka and is 100x simpler to manage.
https://www.confluent.io/blog/kafka-fastest-messaging-system...
How do you teach someone to look at the problems first and then pick the tools?
By understanding what the person cares about. Everyone knows "pick the right tool for the problem". Not everyone uses such a simple calculus because life isn't that simple. People have their own agendas, backgrounds, experiences, career growth desires, personal lives, etc., that are all part of their personal objective function. If you want to convince someone that your tools are better, show that your tools have a higher payoff for their personal objective function. This is way more than a mere product question. In a team setting it's even harder, because you have to balance it across multiple people simultaneously.
"Due to CPU bottlenecks, we were not able to drive a throughput higher than 38K messages/s, and any attempt to measure latency at this rate showed significant degradation in performance clocking a p99 latency of almost two seconds."
I wonder if something similar is happening with Kafka?
As an engineering leader sometimes that even means knowing that people are making the wrong decision, and letting them do it anyway, and then helping them learn from it.
Where do the docs state that? Might be a useful link to keep handy.
They do go to great length to avoid calling Kafka a queue. No where does it directly state that Kafka is not a queue. The docs just never talks about Kafka as being a queue.
https://www.confluent.io/blog/kafka-fastest-messaging-system...
The specific thing I have experience with is in analytics/relational databases. Suddenly around 8-10 years ago it became imperative for every client I was dealing with to migrate their RDBMSes to Hadoop/Hive setups, even when their largest dataset was only about 120M rows denormalized.
They were trading three servers (primary, backup, DR) for sometimes 15 to 20. Queries that MSSQL was handling sub-second were suddenly taking 45s on Hive. It was utter madness and was as far as I can tell driven by good salespeople, FOMO, and the feeling of importance of being able to say your company is running Big Data(TM?).
I saw maybe one implementation (of dozens) that actually stayed in use for any length of time.
If you're using it, you're constrained by its engineering decisions, so you need to be sure it's a worthwhile tradeoff.
When scaling out on a standard message queue you generally have greedy consumers which means you can't assume stickyness, or have to create your own partitioning structures. It makes it great for realtime apps...
If there were cheaper alternatives I think people would use them but it does definitely have powerful benefits.
We actually have a use-case that exactly matches this. One service makes [stuff], the other service consumes [stuff], both services are "immutable infrastructure" with no local stable state storage, and [stuff] is individually too small and frequent to be affordable with IaaS managed-MQ per-message costs — but batching messages into reasonable chunks before send means potentially losing up-to-a-batch worth of messages if the producer dies.
This reminds me of something I experienced
Back in 2000 I worked at a company that hired me to take over their brand new fully redundant web infrastructure. that was architected and built by a firm on a $10,000,000+ contract.
They couldn't understand why the site only served 1 page every 2 seconds. To be very clear: _1_ page every _2_ seconds.
The short story was the development company didn't have anyone who understood databases. So they were using the oracle cluster as a key-value store, and then parsing/sorting XML files on the application servers on demand.
This was a $500,000/yr oracle license, on redundant dedicated $(million) Sun hardware, with huge high speed disk arrays. They were almost sitting idle at peak load... of 1 page every 2 second.
Bonus story: The developer's didn't understand why their dynamic uploading of files to their app servers only worked 1/x of the time where X was the number of app servers in the cluster. They were dropping the files locally on the app-servers and wondering why the other app-servers didn't know about them.
It took a dry erase board and about 30 minutes before they truly understood that the filesystems on those devices were not shared and why. (note, it was never part of the design specification their own team created)
I'm happy that I didn't have to explain ephemeral containers to them.
Sure, if you have the budget to run it, can do! Feel free to reach out, email in the profile ;)
Soon: https://thanos.io/
IMO, Apache Flink is the most complete project for those use cases. It is well maintained and the devs are very helpful when asked on the mailing lists.
The quick summary here is that this was a clean-house rewrite of Apache Storm done by an internal team at Twitter. As an open source project history refresher, Apache Storm was originally built by a startup called Backtype, and the project was led by Nathan Marz, the technical founder of Backtype. Then, Backtype was acquired by Twitter, and Storm became a major component for large-scale stream processing (of tweets, tweet analytics, and other things) at Twitter.
I wrote a summary of the "interesting bits" of Apache Storm here:
However, at a certain point, Nathan Marz left Twitter, and a different group of engineers tried to rethink Storm inside Twitter. There was also a lot of work going on around Apache Mesos at the time. Heron is kind of a merger of their "rethinking" of Storm while also making it possible to manage Storm-like Heron clusters using Mesos.
But, I don't think Heron really took off. Meanwhile, Storm got very, very stable in the 1.x series, and then had a clean-house rewrite from Clojure to Java in the 2.x series, mainly to improve performance even more. The last stable/major Storm release was in 2020.
Storm provides a stream processing programming API, a multi-lang wire protocol, and a cluster management approach. But certain cluster computing problems can probably be better solved at the infrastructure layer today. (For example, Storm was developed before the whole container + docker + k8s focus in cloud ops.) That said, it's still a very powerful system; on my team, we process 75K+ events per second across hundreds of vCPU cores and thousands of Python processes with sub-second latencies by combining Storm and Kafka with our open source Python project, streamparse.
https://github.com/Parsely/streamparse
The core problems Storm solves: modeling data processing as a computation graph; high-speed network communication between threads, processes, and nodes; message delivery guarantees and retry capabilities; tunable parallelism; built-in monitoring and logging; and much more.
(Also, I'd be remiss if I didn't mention -- if you're interested in stream processing and distributed computing, we are hiring Python Data Engineers to work on a stack involving Storm, Spark, Kafka, Cassandra, etc.) -- https://www.parse.ly/careers/python_data_engineer
I'm not even a clojure user, but my impression was that it was pretty performant. I remember a discussion that they didn't even really need the JVM invokedynamic because they were doing pretty well without it, so that made me think it was close to pure JVM speed.
Lately it's been starting to feel the same for distributed systems. How many streaming engines are there now under Apache? Four?
We've used Python at scale on petabyte-sized production data, multi-billion-request API tiers, and high-concurrency low-latency data processing all the while.
Java is the right tool for writing system-level code in some isolated contexts, but Python is the right tool for the job for a huge number of important use cases, with those use cases growing by the day.
I love Java (and appreciate Kotlin and Clojure), but with parallel compute power cheaper by the day, Python's focus on code simplicity, open source ecosystem, and programmer happiness continues to win the day.
I assume you mean "Kafka Streams" when you say this, as Kafka is just a event bus and Flink can be used to read messages from it.
The biggest advantage of Flink (IMO) is that you can write the Flink logic once, and reuse it for both batch and stream processing. So if I write a Flink job that consumes a Kafka stream and produces some aggregated outputs, that same job can be run against my data lake in S3/GCS/Azure Blob Storage, etc.
Kafka Streams does not support batch processing, or working on top of anything other than Kafka. Flink supports building on top of other message buses like Pulsar as well: https://flink.apache.org/news/2019/11/25/query-pulsar-stream...
* Flink work on other message queue tech other than Kafka like Amazon SQS, Pulsar, etc..
* Back when I last read about Kafka Stream, Flink has better support for stateful processing, the entire state of execution can be safely snapshot into storage and resume at any time.
* Kafka Stream shuffle data is slower, because it has to send data to a new topic in the broker, instead of sending it directly between compute node.
Advantage of Kafka Stream over Flink
* Kafka Stream deployment is simple, just start the jar like any other Java program, scaling up by running the jar multiple time. On the other hand, Flink need a mature orchestration framework like YARN, Meso or K8S, trying to manage a Flink deployment without them is very painful.
* Flink require a central, persistent storage like HDFS or S3 for its checkpoint mechanism; Kafka stream doesn't.
- Kafka is a distributed log.
- Storm, Samza & Flink are stream processing engines.
- Spark is a Map/Reduce framework that uses memory to cache computations to provide some performance increase over other disk-based frameworks. It can also do some streaming computations if you squint hard enough.
- Confluent is a company that sells an enterprise Kafka.
Not really sure the comparison you made is apt.
> Built-in Stream Processing > Process streams of events with joins, aggregations, filters, transformations, and more, using event-time and exactly-once processing.
So unlike Flink, Storm, Spark, Heron etc. it's only useful with Kafka.
> Non-idempotent stateful topologies are stateful topologies that do not apply processing logic along the model of "multiply by zero" and thus cannot provide effectively-once semantics. An example of a non-idempotent
https://heron.incubator.apache.org/docs/heron-delivery-seman...
(I'm puzzled by their idempotent vs non-idempotent stateful topology description, because if something is mutating an internal state upon receiving events, it will likely be non-idempotent by design... unless they just mean "idempotent stateful" here to refer to keeping track of source/output position state and such.)
(They also do say that can only support state storage in ZK or local FS, which feels like a likely non-starter compared to Flink for some of my use cases.)
1. Enriching event streams. Say you have a stream of log records with an IP address field. You want to enrich with a geo-location before sending the logs to Elasticsearch.
2. Windowed aggregation. Maybe you have an application that is emitting "login" events and you want to to detect login attempts from different IP addresses within X minutes of each other.
3. Joining multiple event streams. You have multiple different event streams and you want to join them together using some common join key (maybe session ID or something like that) to compute a metric that aggregates all of them.
There are plenty of more esoteric use cases as well.
The lowest latency model refresh time I've come across is ~5 minutes. Going lower is likely to be a lot of data transfers for model synchronization (sync the serving model with training model) for little value.
Is it apt? Reader exercise.