6 ms·
If you're a Python user, I gave a talk at PyData called "Beating the GIL with Python" that covers all the tooling for parallel computing in Python, from the sim
by pixelmonkey 10y ago
If you're a Python user, I gave a talk at PyData called "Beating the GIL with Python" that covers all the tooling for parallel computing in Python, from the simplest (e.g. multiprocessing) to the more complex (Storm/Spark), with many tools in between, and what the trade-offs are. That may help you out.
https://www.youtube.com/watch?v=gVBLF0ohcrE https://www.youtube.com/watch?v=gVBLF0ohcrE
The rough summary is, if you need to do distributed multi-node computation with very low latency -- that is, end-to-end data processing time measured in millis or seconds -- then Storm will help you out. The "competitive" product in open source right now is pyspark + Spark Streaming.
The benefit of Storm is that it is an older (more mature) project with many production deployments. My team at Parse.ly wrote the open source Python + Storm integration libraries, pystorm and streamparse, which help Python programmers make use of Storm for large-scale stream processing: https://github.com/Parsely/streamparse https://github.com/Parsely/streamparse. We also wrote pykafka, which is a Kafka client for Python: https://github.com/Parsely/pykafka https://github.com/Parsely/pykafka -- these projects are related since Kafka is often the streaming data source used as the input to a Storm topology.
- RBerenguel 10y agoAfter the recent comparisons on real time processing between Flink and Spark, isn't Flink the actual "competitive" product? In the sense that it's real "real time", not the way Spark streaming (not dismissing it, it's what I actually use) works.
- pixelmonkey 10y agoThat's an interesting question. Unfortunately I haven't had a chance to evaluate Flink in-depth, it being relatively new.
- gshulegaard 10y agoThat's the gist. I am a Flink fan, but I get the impression it's still in the research and development phase...although rapidly approaching a potentially production-level quality. That said, Storm is Event-streaming processing which is usually the focal point of most comparisons between Flink and Spark Streaming.
- narsil 10y agoI prefer Celery for Python. Far more light-weight, and intended for distributed task processing. http://celery.readthedocs.org/ http://celery.readthedocs.org/
- brianwawok 10y agoAre those really competing though? Celery rocks for some background image processing. Not for multinode analytics of huge data sets.
- throwaway-anon 10y agoActually that's not true. Have a look at the airflow project by AirBnB. https://github.com/airbnb/airflow https://github.com/airbnb/airflow
- squeaky-clean 10y agoDoesn't seem like it can handle processing on real-time streams of data. They don't really seem comparable. From the readme: "Airflow is not a data streaming solution. Tasks do not move data from one to the other (though tasks can exchange metadata!). Airflow is not in the Spark Streaming or Storm space, it is more comparable to Oozie or Azkaban."
- brianwawok 10y agoThat doesn't really seem the same at all.
- throwaway-anon 10y agoI was replying to a comment not saying airflow is capable of stream processing. > Not for multinode analytics of huge data sets. Airflow does this just fine with celery.
- bcbrown 10y agoWhat's your thoughts on Dask? I saw a talk on it at Strata, and it seemed fairly promising.
- pixelmonkey 10y agoVery promising IMO. Lots of new stuff going on in their 'distributed' project: - https://github.com/dask/distributed https://github.com/dask/distributed - http://matthewrocklin.com/blog/work/2016/02/17/dask-distributed-part1 http://matthewrocklin.com/blog/work/2016/02/17/dask-distribu...
- gjulianm 10y agoAt what point is using Storm or any other alternative inevitable? I am always reminded of [1] in these cases. Have you found instances where an implementation of the core logic in C/C++ (maybe with thread parallelism) is fast enough? 1: http://aadrake.com/command-line-tools-can-be-235x-faster-than-your-hadoop-cluster.html http://aadrake.com/command-line-tools-can-be-235x-faster-tha...
- pixelmonkey 10y agoIt's always a valid question -- when trying to scale to handle large data volumes, it's always worth trying to scale "up" (bigger single machine) and "in" (faster code) before scaling "out" (more machines). In our case, we had pushed the limits of the largest EC2 box we could find. We have to keep up with tens of thousands of events per second, and do non-trivial amounts of CPU processing on each of them. Whether Python, Java, C++, or assembler, it wasn't going to happen on a single box. We proved this to ourself by getting the biggest box we could and rewriting the hotspots of our code in Cython. It still wasn't enough.
- e12e 10y ago> In our case, we had pushed the limits of the largest EC2 box we could find. We have to keep up with tens of thousands of events per second, and do non-trivial amounts of CPU processing on each of them. I understand there are many reasons for sticking with AWS, but looking at reserved instances[1], I find: r3.8xlarge which has 32 cores, 244GB of RAM and 2 x 320 GB SSD storge, for $2.66 per Hour -- or just under 2K a month. For just a little more (USD 2099/month, no minimum term, no setup), you can get something like: Dell R930, 4x Intel Xeon E7-4820v3 (that's 40 cores at base frequency of 1.9 GHz), 384GB DDR4, 2x480GB SSD HW RAID 1 (although I suspect 6x120GB, possibly in raid0 might be better) with 1Gbps Full-Duplex and 10 TB of egress bandwidth included from leasweb: https://www.leaseweb.com/dedicated-server/configure/22651 https://www.leaseweb.com/dedicated-server/configure/22651 Granted, one might want 10Gps - and managing your own server isn't free, even when the hosting company takes care of the physical hardware (not that managing an EC2 instance is free either). But I'm curious if you tested on dedicated hardware as well? 10.000 events at 4kb/event is "just" 300 mbps (call it 1 gbps including overhead). I'm certainly not claiming it's easy to any kind of processing at a sustained ~1gbps -- but I'm curious if the kind of workload you're discussing could be handled by a single (relatively cheap) dedicated server? [1] http://aws.amazon.com/ec2/pricing/#reserved-instances http://aws.amazon.com/ec2/pricing/#reserved-instances
- iagooar 10y ago> The rough summary is, if you need to do distributed multi-node computation with very low latency -- that is, end-to-end data processing time measured in millis or seconds -- then Storm will help you out. I still don't really know what kind of computation that is. Any specific, real-world examples? I'm genuinely curious.
- wlesieutre 10y agoCompression for live video streams maybe? You need to take a high quality input, continuously recompress it to several different qualities (for varying connection speeds and screen sizes) and then kick the data out to a CDN for viewers?
- pixelmonkey 10y agoWe use it for real-time web/content analytics. You can see our product tour @ http://parse.ly/tour http://parse.ly/tour. Other examples of "common" steaming data with real-time use cases: - Twitter firehose - Network packet analysis - Financial market data - Device and sensor data
- lkrubner 10y agoYou would use this where you need to take in a lot of data from multiple sources, and create some kind of information out of that, with soft real-time constraints. So, for instance, web analytics. Or Natural Language Processing, if the communication was between humans and machines, and the responses needed to be real-time. Think about getting a terabyte of data a day, or more. Or you can take a different approach, and look at a system that does not use Storm, and then imagine at what scale that system would break. Do consider the system that Matthias Nehlsen wrote about here: http://matthiasnehlsen.com/blog/2014/11/07/Building-Systems-in-Clojure-4/ http://matthiasnehlsen.com/blog/2014/11/07/Building-Systems-... Using Redis as a universal bus is a fine architecture for small amounts of data, where "small amounts of data" means a gigabyte a day, or maybe a little bit more. But what about 10 gigabytes of data per day? What about about 100 gigabytes of data per day? Universal bus architectures break down when each node sees too many messages that are not meant for that node. There is a limit on how far you can go with an architecture where every node gets every message. Assuming that each message is suppose to be read by certain nodes, there is an inefficiency to sending messages to nodes that don't want the message. But there is a wonderful flexibility to this system. And I've created similar systems. Nevertheless, at a certain scale, you need to give up that flexibility for the sake of scale. Again, Matthias Nehlsen says this well: http://matthiasnehlsen.com/blog/2014/10/30/Building-Systems-in-Clojure-3/ http://matthiasnehlsen.com/blog/2014/10/30/Building-Systems-... "What comes to mind immediately when regurgitating the requirements above is Storm and the Lambda Architecture. First I thought, great, such a search could be realized as a bolt in Storm. But then I realized, and please correct me if I’m wrong, that topologies are fixed once they are running. This limits the flexibility to add and tear down additional live searches. I am afraid that keeping a few stand-by bolts to assign to queries dynamically would not be flexible enough." I would (and have) done exactly what Matthias Nehlsen suggests: use something like Redis for as long as you can. But I would also do what Montalenti eventually did, and adopt Storm to deal with a certain level of scale -- when you reach a scale where a universal bus architecture no longer works, sacrifice flexibility for performance, and switch to Apache Storm. If I recall correctly, Montalenti and Parse.ly experimented with MongoDB and Cassandra and a bunch of other technologies before they found that Kafka/Storm is what they needed. Maybe he's written about the whole journey somewhere. Certainly, it would be interesting to read why all the other systems failed, and why moving to Kafka/Storm was the right choice for Parse.ly.