4 ms·
Answer the question about how to process a few TB on a multicore server, and you might still find yourself using Spark, or something like it. If you start from
by zten 6y ago
Answer the question about how to process a few TB on a multicore server, and you might still find yourself using Spark, or something like it.
If you start from the assumption that you've been ingesting data and storing it in a compressed columnar format like Parquet or ORC, then you're already locked into a solution that exists in the Hadoop ecosystem. This turns out to be an effective way to deal with terabytes of data because depending on the query you want to execute, the file format (and a smart partitioning structure) helps turn your problem of reading terabytes into one of reading gigabytes. And, everything in Hadoop land is generally a query, so you're going to need a query engine like Hive, Impala, Spark, etc. AND, you probably need something like Hive metadata so that you're not crawling some directory structure for the schema and input files every time you start up a new process to run a query.
You _could_ write something on your own that just forks out a bunch of threads in a single process to rip through the data, but why? Think about what you're effectively implementing - a bespoke query engine that runs one query plan. Spark has already written an API and a query engine (multi-process distributed, unlike whatever you're likely to hand-roll) and lots of input/output code. You can spark-submit --master local[*] and use up all of the cores and ram everything into one monster JVM if you really wanted to.
Finally, you could stuff all of this in memory, but at what cost, and what are you going to use to do it? (Worse, what happens if the VM goes down - those terabytes are going to take their sweet time reloading over a network link.)
- throwaway_pdp09 6y agoHow; plenty of disks to give you the IO. It's usually about the IO. If you want memory, go to a server site and configure to max out a server with lots of DRAM slots. You might be surprised. But really getting enough mem is less important than IO, so stick with disks, SSDs I suppose. You can get single chip with dozens of cores for not too much. I guess that's how I'd do it. > This turns out to be an effective way to deal with terabytes of data A few terabytes don't need cluster, typically. > You _could_ write something on your own that just forks out a bunch of threads in a single process to rip through the data, but why? Because it's simple and easy. I wrote one in a few days. Not much code. > You can spark-submit --master local[] and use up all of the cores and ram everything into one monster JVM if you really wanted to. Point is, do you need to? > those terabytes are going to take their sweet time reloading over a network link You do not run big IO over a network like that if you can avoid it. With a single server with plenty of SSDs, you trivially don't. My take anyway.
- zten 6y agoAh, so the thing I was attacking was the "if it fits in RAM, it isn't big data" meme. AWS is happy to sell you on the idea that you can use S3 as the persistent storage and then bring up the compute whenever. If you bring up the compute on-demand, then the network link matters, whether it's EBS spoon-feeding you data or you using multiple VMs and the right object storage strategy in S3 to suck the data out as fast as possible. You pay a premium to have this stuff loaded up on faster hardware. If you look at how Redshift works on their compute-optimized nodes, everything's sitting in local NVMe SSDs (a slightly different strategy is used on ra3 nodes). They handle the fact that this is ephemeral with automatic backups. Cluster restores aren't exactly fast. I actually agree with you; if you've got some monster SSDs attached, which are comparatively cheap, why focus on the RAM... I believe there are reasons to do that sometimes, but not everything demands quite that level of performance. For this point: >> You _could_ write something on your own that just forks out a bunch of threads in a single process to rip through the data, but why? > Because it's simple and easy. I wrote one in a few days. Not much code. I think it depends on what formats you're using and how it's laid out on disk. A lot of people reshape their data into a table-like structured or semi-structured format, and that makes it a candidate for putting it into a database like Postgres or Redshift, and other times it makes more sense to bring the database to the data (the Hadoop ecosystem of stuff.) For example my company still has stuff that writes out row-like objects into S3, and it's not too hard to write a single process job that spawns threads, reads the input files, and does some computation on those. They're on S3, so throughput kinda sucks, and copying it to local SSD only makes sense if you want to make multiple passes. But the Parquet format alternative of this data is tremendously faster to work with and genuinely easier to use, and the only barrier to entry is that you run Spark on your local machine and commit to using Spark with Elastic MapReduce to do processing for this data. Sure, you might only end up filtering through tens or hundreds of gigabytes; terabytes is usually rare. But part of that is because you only read 10-20% of every input file to do the work. It also integrates extremely well with other stuff we're using - it's even easier to just use Snowflake against the same data set, for example. Sorry, I think I'm rambling at this point and not presenting a really coherent argument.
- throwaway_pdp09 6y ago