We've become so accustomed to extremely inefficient software systems that we've lost all perspective on what is possible.
We've become so accustomed to extremely inefficient software systems that we've lost all perspective on what is possible.
Are you saying this can be optimised to fit inside a single 10 core server in terms of compute loads?
You use a cluster when your data and compute requirements are large and parallel enough that the tax paid on network latency trumps the 10-20X speedup you get on SSD and 1000X speedup you get from just keeping data in RAM.
250 Gigs is tiny enough that you could probably much get better performance running on high memory instance in AWS or GCP. You'll generally have to write your own multiprocessing code though which is fairly simple - your existing library may also be able to support it.
I once actually ran this kind of workload on just my laptop using a compiled language that performed better than pyspark on a cluster.
However, even if I could load it all into memory at once, and assuming it takes 200 gb, I'm still using a master's student access to a cluster. So I get preempted like it's nobody's business. Hence why I prefer a smaller memory footprint even if I take up cpus at variable rates through a single execution.
I did try to write my own multiprocessing code for this, but the operations are sometimes too complicated (like groupby) for me to rewrite everything from the ground up. If I'm not reliant on serial data communication between processes (like you'd need to sort a column), I can get it done pretty easily. In fact, I wrote my data cleaning code with this and cleaned up the entire file in half an hour because single chunks didn't rely on others.
However, if you have some idea of how to run these computational loads in parallel in python or any other language on single compute instances (like the size of a laptop's memory of 16 gb), I'd really love to see it. Thanks.
I do have fast SSD storage because it's on the scratch drive of a cluster and from what I've seen it can do ~300-400 MB/s easily. I haven't had a chance to test more than that since I'm mostly memory constrained in much of my testing.
My current attempt is to push this data into a pure database handling system like SQL so that I can query it. But like I said, I work with a less-than-stellar set of tools and I have to literally set up a postgres server from ground up to write to it. Which shouldn't be a big deal except when it's on a non-root user and I have to keep remapping dependencies (took 5-6 hours to set it up on the instance I have access to).
My other option was to write the entire 250 GB to a sqllite database using the sqlalchemy library in Python, but that seems to fail whether I do it with parallel writes, or serial writes. In both cases, it fails after I create ~64-70 tables.
> Generally speaking, Dask.dataframe groupby-aggregations are roughly same performance as Pandas groupby-aggregations, just more scalable.
The dask.distributed scheduler can also run on one high-RAM instance (with threads or processes) https://docs.dask.org/en/latest/setup.html
Pandas docs > Ecosystem > Out-of-core: https://pandas.pydata.org/pandas-docs/stable/ecosystem.html#...
Reading from Parquet into Apache Arrow is much faster than CSV because the data can just be directly loaded into RAM. https://ursalabs.org/blog/2019-10-columnar-perf/
If you have GPU instances, cuDF has a Pandas-like API on top of Apache Arrow. https://github.com/rapidsai/cudf
> Built based on the Apache Arrow columnar memory format, cuDF is a GPU DataFrame library for loading, joining, aggregating, filtering, and otherwise manipulating data.
> cuDF provides a pandas-like API that will be familiar to data engineers & data scientists, so they can use it to easily accelerate their workflows without going into the details of CUDA programming.
Dask-ML makes scalable scikit-learn, XGBoost, TensorFlow really easy. https://dask-ml.readthedocs.io/en/latest/
... re: the OT: While it's possible to write C++ code that's really fast, it's generally inflexible, expensive to develop, and dangerous for devs with experience in their respective domains of experience to write. Much saner to put a Python API on top and optimize that during compilation.
There are a few C++ frameworks in the top quartile of the TechEmpower framework benchmarks. https://www.techempower.com/benchmarks/
Hardware/hosting is relatively cheap. Developers and memory vulnerabilities aren't.
"Dask on HPC, what works and what doesn't" https://github.com/dask/dask-blog/issues/5
Maybe you should spend some time developing a job visualization system for end users from scratch, for end users with lots of C, JS, and HTML experience https://jobqueue.dask.org/en/latest/interactive.html
You can create memory mapped ndarrays, these act like normal numpy arrays but don't need to fit into RAM. Numpy maps the array to a binary file on disk. The array otherwise acts like an ndarray so you can build a DataFrame with it. Whenever you access an array index Numpy in the background (essentially) seeks that many values into the file to grab the value of that index.
Since you're on a fast SSD and Numpy is fairly smart you'll be able to access your arrays close to your drive's speed. It's slower than if the whole database was in RAM but far faster than distributing the data over a network to a bunch of worker nodes. Memory mapped files let you have array-like access to data on disk as if it lived in RAM. When building a pandas DataFrame from a memmapped ndarray I believe you just need to set copy=False in the constructor for it to Just Work.
I don't know what your data looks like but I doubt loading it into SQLite is going to improve your performance.
Your second paragraph is essentially what I want. I'm willing to wait a day for code that may run in 1 hour from memory, so time isn't entirely an issue unless it's starting to bleed into weeks. The read_csv function in pandas has a parameter called memory_map, but when I tried using it on a smaller 7GB dataset, it read the whole thing into memory (32GB instance) even when I set it to True.
SQLite is definitely not my best option here. It was the only server-less implementation I could find, so I tried to use it and it didn't work. However, a database like implementation will be helpful because each operation I need to do will require data that satisfies certain timestamp and arithmetic conditions. I figured it'd be best to load the whole thing into a db and query it for every operation to train my model.
Nowadays there's another reason to use clusters: to autoscale your expenditure wrt workload. A little inefficiency might be acceptable if you don't have to pay for a huge beefed up server idling at, say, 30%.
Biased: you might want to look into DOE libraries. For IO I suggest ADIOS2 [0]. There's python bindings too.
One of the biggest things you can do is use a different storage type like BP (ADIOS) or hdf5. These are readable but binary. But to really determine how to speed up your problem you have to know where the bottleneck is. Is it IO or compute? With 100 workers (threads or nodes?) you aren't highly parallelized. I mean that could be a single node if it's threads.
I'm curious as to why this seems like a common mistake (I've seen it a few times already in this comment section).
If you have some good indexes and do some push-down work (give the database aggregation tasks to do instead of your python code), you should probably be more than fine.
For a 250Gb file.. should be ok.. maybe add some partitioning too.
https://pandas.pydata.org/pandas-docs/version/0.22/generated...
https://pandas.pydata.org/pandas-docs/version/0.22/io.html#w...
Obligatory read: https://aadrake.com/command-line-tools-can-be-235x-faster-th...
I'm currently employed to write software which does DNA analysis. DNA is known for being Big Data.
Some applications are very compute intensive and others are very data-relation intensive. The very compute intensive applications process about 1GB of data in about 30 minutes on a 32 core Xeon 6xxx with 32GB of RAM assigned to it. The data-relation intensive application processes about 350GB of data in 15 minutes on a single core of the same CPU but with 700GB of RAM assigned to it.
Both are heavily optimized but in different ways. So without knowing more about what you're doing with that data, it's hard to say.
The compute-heavy workload calculates edit distance [0] of short paired-end sequenced DNA [1] vs the human genome [2]. There is open source software to manipulate the FASTA/FASTQ [3] and SAM files [4] and run the calculations [5]. The aligned file is processed in a couple of minutes to genetic variation report [6] which is used for some of the analysis products that were purchased. One popular product will give you a haplogroup [7] which basically tells you where you are in a genetic tree.
The relationship estimator uses a different sequencing technology and basically consumes a CSV file from the sequencer's manufacturer. It uses a proprietary algorithm to calculate centimorgans [8]. That then gives relationship estimates between you and other people who've purchased the product.
[0]: https://en.wikipedia.org/wiki/Needleman%E2%80%93Wunsch_algor...
[1]: https://en.wikipedia.org/wiki/Next-generation_sequencing
[2]: https://en.wikipedia.org/wiki/Reference_genome
[3]: https://en.wikipedia.org/wiki/FASTQ_format
[4]: https://github.com/samtools/hts-specs
[5]: https://en.wikipedia.org/wiki/List_of_sequence_alignment_sof...
[6]: https://en.wikipedia.org/wiki/Variant_Call_Format
While I hope that each 250GB chunk has everything in order and I can separate it cleanly, I don't trust it. So I broke up all the files by stock in an embarrassingly parallelised code I wrote and got 500 separate files. However, the problem is, I'd want to process these 500 files parallely so each stock gets only so much memory (and hence, my original constraint remains). And for each stock, I want the data to be queried on timestamps, so I need a way to quickly say "I have data at 9am on Thursday, I want data from 4pm Wednesday to 8am Thursday to create a model".
I figured the best way to do this efficiently was to have the code create a massive database for all the stocks and query it efficiently in SQL. But I'm stuck there due to a lack of tools.
I have no doubt I could solve your problems. I honestly don't care to do so here though.
I imagine there's a software development/engineering team at your work or school you could ask for guidance though.
That's fair. I did not think it'd be such a difficult problem when I first set out to do it myself. But every single turn leading to a dead end kinda bummed me out. I'm going to finally resort to the database method of storing all the data in one query-able file and work off of that.
Segmenting your raw data and using memory mapped files will let you work with large data sets without needing huge amounts of RAM. From there it's a question of your single system's processing speed and IO capacity. This is only necessary if your processing needs random access to the entire dataset.
If your CSV data is more like a streaming data source, you're processing each record as it's read in, you can just stream it in through `stdin`. At 1GB/s you're looking at five minutes or so to process your 250GB of raw data. A SATA SSD might take twenty minutes to stream that raw data.
It's important to note that often your disks aren't directly attached to your compute. That's frequently the case in (particularly cheap) cloud instances.
It's also a domain where you can buy an off-the-shelf desktop for a few hundred dollars to do the work. That's the thrust of this whole thread, because scalable "cloud" systems exist and look cheap people obsess about throwing more instances at problems.
Modern commodity systems are ridiculously powerful and far more capable than people tend to assume. Even "the cloud" gets underestimated because people look at the low end cheap instances and assume they need to spin up hundreds of those when one beefy image for a short duration could do the same work.
Nit: 1GB/s is ok, not even solid let alone fast.
A fast SSD pretty much saturates 4x3.0 links (which explains why they universally tend to cap out at 3.5GB/s). In fact there are now a few PCIe4 SSDs (e.g. Corsair's MP600) which close in on 5GB/s.
If your latency requirements are slack, then you can get away with one machine, because you can reboot or reprovision it and carry on processing without meet your requirements.
If your latency requirements are tight, you don't have time to fail over anyway, so you might as well run one machine and make sure you can deal with failing to meet your requirements.
Go back to two when maintenance is done.