Command-line Tools can be 235x Faster than your Hadoop Cluster (2014)
adamdrake.com
adamdrake.com
Just because your throw-away 40 line script worked from cron for five years without issue doesn't mean that a seven node hadoop cluster didn't come with benefits. You got to write in a language called "pig"! so fun.
That is why Rust is so awesome. It still allows me to get stuff in my resume, but still make an executable that runs on my laptop with high performance.
Here I sit, running a query on a fancy cloud-based tool we pay nontrivial amounts of money for, which takes ~15 minutes.
If I download the data set to a Linux box I can do the query in 3 seconds with grep and awk.
Oh but that is not The Way. So here I sit waiting ~15 minutes every time I need to fine tune and test the query.
Also, of course the query now is written in the vendor's homegrown weird query language which is lacking a lot of functionality, so whenever I need to do some different transformation or pull apart data a bit differently, I get to file a feature request and wait a few month for it to be implemented. On the linux box I could just change my awk parameters a little bit (or throw perl in the pipeline for heavier lifting) and be done in a minute. But hey at least I can put the ticket in blocked state for a few months while waiting for the vendor.
Why are we doing this?
someone got promoted
The result? 1. I don't have access to basic logs for debugging because apparently the infra guy would have to give me access to the whole cluster. 2. Production ends up dying from time to time because apparently they don't know how to set it up. 3. The boss likes him more because he's using big boy tools.
People who weren’t developing around this time can’t appreciate how game changing SSDs were then spinning rust.
I/O was no longer the bottleneck post SSD’s.
Even today, people way underestimate the power of NVME.
> The issues of how to parallelize the computation, distribute the data, and
> handle failures conspire to obscure the original simple computation with large
> amounts of complex code to deal with these issues. As a reaction to this
> complexity,we designed anew abstraction that allows us to express the simple
> computations we were trying to perform but hides the messy details of
> parallelization, fault-tolerance, data distribution and load balancing in
> a library .
If you are not meeting this complexity (and today with 16 TB of RAM and 192 cores, many jobs don't) then Map-Reduce / Hadoop is not for you...Makes sense, we are told that vertical has limits in university and we should prioritise horizontal; but I feel a little like the "mid-wit" meme, once we realise how vertical we can go then we can end up using significantly fewer resources in aggregate (as there is overhead in distributed systems of course).
I also think we are disincentivised from going vertical as most cloud providers prioritise splitting workloads, most people don't have 16TiB of RAM available to them, but they might have a credit card on file for a cloud provider/hyperscaler.
*EDIT*: Largest AWS Instance is, I think, the x2iedn.metal ith 128vCPU and 4TiB RAM
*EDIT2*: u-24tb1.metal seems larger; 448vCPU and 24TiB Memory, but I'm not sure if you can actually use it for anything that's not SAP HANA.
My laptop is a SPoF in exactly the same way.
If my laptop is closed then data collection will still happen, as collection and processing are different systems; but my ability to mutate the data hands-on is affected.
Thus any downtime of my laptop is not really a problem.
See also: Jupyter notebooks, Excel, etc;
I will also point out that robustness in distributed systems is not as cut and dry for two reasons:
1: These are not considered hot-path systems that are mission critical so will be neglected by SRE.
2: Complexity is increased in distributed systems, thus you have more likelihood of failure until you have a lot of effort put into it.
Hbase also ran on that infrastructure serving real-time workloads. Downtime on any of the Hbase clusters would be a high severity outage.
So minutes/mo of downtime would certainly have unacceptable business impact. Another important thing is replication. Drives do fail, and if a single drive failure brings down prod how long would that take to fix?
To be clear in general my opinions are aligned with the article, I think using the whole machine at high utilization is the only environmentally (and financially) responsible way. But I don't believe it's true that purely vertical scaling is realistic for most businesses.
EDIT: there are also security and compliance concerns that rule out the scenario of copying data onto an employee laptop. I guess what I'm trying to get at is the scenario seems a little contrived.
You already failed if thats happening.
Are we really at the degenerated level of sysadmin competence that we forgot even what RAID is?
Like all "additional components", RAID controllers come with their own quirks and I have heard of rare cases of RAID controllers being the cause of data loss, but RAID as a concept is designed to combat bit-rot and lossy hardware.
ZFS in the same vein is also designed around this concept and attempts to join RAID, an LVM and a filesystem to make "better" choices on how to handle blocks of data. Since RAID only sees raw blocks and is not volume or filesystem aware there are cases where it's slower.
That said, I have to also mention that when I was investigating HBASE there was no way to force consistency of data, there was no fsync() call in the code, it only writes to VFS and you have to pray your OS flushes the cache to disk before the system fails. HBASE Parity is configured by HDFS which is essentially doing exactly what RAID does. Except only to VFS and without parity bits.
If so, then both things in the gp are true: raid isn't enough, and can be a false sense of security.
HDFS is distributed across multiple machines, each one which can have RAID. It is unlikely that enough machines will fail to lose data.
At the risk of troll-feeding, what are you hoping to accomplish with this? Of course I haven't "forgot even what RAID is", and I'm confident my competence is not "degenerated".
In this magical world where we can fit the entire "data lake" on one box of course we can replicate with RAID, but you've still got a spof. So this only works if downtime is acceptable, which I'll concede maybe it could be iff this box is somehow, magically, detached from customer facing systems.
But they never really are. Assuming even that there aren't ever customer impacting reads from this system, downtime in the "data lake" means all the systems which write to it have to buffer (or shed) data during the outage. Random, frequent off-nominal behavior is a recipe for disaster IME. So this magic data box can't be detached, really.
I've only ever worked at companies which are "always on" and have multi-petabyte data sets. I guess if you can tolerate regular outages and/or your working data set is so small that copying it around willy-nilly is acceptably cheap go for it! I wish my life was that simple.
If you really have multi-petabyte datasets then probably you are at the scale where distributed storage and systems will be superior.
The point of this conversation is that most people are not at this scale but think they are. IE: they sincerely believe that a dataset does not fit in ram of a single box because it's 1TiB or they think because it doesn't sit on a single 16TiB drive then a distributed system is the only solution.
The original post is an argument about that; that a single node can outcompete a large cluster, so you should avoid clustering until it really cannot fit on a single box anymore.
Your addendum was reliability is a large factor. Mostly this does not bear resemblance with reality. You might be surprised to learn that reliability follows a curve where you get very close to high reliability with a single machine, you diminish it enormously with a distributed system and then start approaching higher reliability when you have a lot more effort into your distributed system.
My comment about RAID was simply because it's very obvious that a single drive failure should not be taking a single machine down, similarly a CPU fault or memory fault can also be configured to not take down a machine. That you didn't understand this was either a failing of our industry knowledge; or, if you did understand this then the comment was disingenuous and intentionally misleading- which is worse.
I've also only worked at companies that were "always on" but that's less true than you think also.
I have never worked anywhere that insisted that all machines are on all the time, which is really what you're arguing. There is no reason to have a processing box turned on when there's no processing that's required.
Storage and aggregation: sure, those are live systems and should be treated as such, but it is never a single system that both ingests and processes. Sometimes they have the same backing store, but usually there is an ETL process and that ETL process is elastic, bursty, etc. and its outputs are what people are actually doing reports based on.
For example, I think Dean & Ghemawat reasonably describe what were their incentives: saving capital by reusing an already distributed set of machines while conserving network bandwidth. In table 1 they write average job duration was around 10 minutes involving 150 computers and that on average 1.2 workers died per such job!
The computers had 2-4 GiB memory, 100megabit ethernet and ISA HDDs. In 2003 when they got map reduce going Google's total R&D budget was $90million. There was no cloud so if you wanted a large machine you had to pay up front.
What they did with Map Reduce is a great achievement.
But I would advise against scaling horizontally right from the start because we may need to scale horizontally at some time in future. If it will fit on one machine, do it on one.
People came up with horizontal scaling across COTS hardware, often called Beowulf clusters, to have more computing for cheaper. They’d run UNIX or Linux with a growing collection of open-source tools. They’d be able to get the most out of their compute by customizing it to their needs.
So, vertical scaling being exorbitantly expensive and less flexible at the time, too.
On the other hand they had 2GHz Xeon CPUs!
Table 1 in the paper suggests that average read throughput per worker for typical jobs was around 1MB/s.
(note, comparing a 2 node active-passive hot-spare setup with a bazillion node horizontal hellscape is not in scope.)
Reflecting on a decade in the industry I can say cut, sort, uniq, xargs, sed, etc etc have taken me farther than any programming language or ec2 instance.
https://news.ycombinator.com/item?id=30595026 - 1 year ago (166 comments)
https://news.ycombinator.com/item?id=22188877 - 3 years ago (253 comments)
https://news.ycombinator.com/item?id=17135841 - 5 years ago (222 comments)
https://news.ycombinator.com/item?id=12472905 - 7 years ago (171 comments)
It’s been 8 years and I think RDBMS is stronger than ever?
Colored the entire course with a “yeah right”. Frankly is Hadoop still popular? Sure, it’s still around but I don’t hear much about it anymore. Never ended up using it professionally, I do most of my heavy data processing in Go and it works great.
Most of the companies I have worked with that actively have spark deployed are using it on queries with less than 1TB of data at a time and boy howdy does it make no sense.
If the workload fits in memory and a single machine, DuckDb is so much more lightweight and faster.
This data is all read-only, I suspect a set of PostgreSQL servers would perform much better.
If your data is too big to fit into DuckDB, consider Clickhouse, which is also column-based and understands standard SQL.
Some of the projects I remember from the Joyent team were: dumping recordings of local mariokart games to manta and running analytics on the raw video to generate office kart racer stats, the bog standard dump all the logs and map/reduce/grep/count them, and I think there was one about running mdb postmortems on terabytes of core dumps.
Far too often we waste time optimising for problems we don’t have, and, most likely, will never have.
The reason we use these is for when we have a data set _larger_ than what can be done on a single machine.
I’ve heard this anecdote on HN before but without ever seeing actual evidence it happened, it reads like an old wives tale and I’m not sure I believe it.
I’ve worked on a Hadoop cluster and setting it up and running it takes quite serious technical skills and experience and those same technical skills and experience would mean the team wouldn’t be doing it unless they needed it.
Can you really imagine some senior data and infrastructure engineers setting up 100 nodes knowing it was for 60GB of data? Does that make any sense at all?
They had a team working on a Hadoop based solution and their biggest internal implementations was about what you're describing, in practice.
It makes sense because internal politics.
each node in our hadoop cluster had 64GiB of ram (which is the max amount you should have for a single node java application, where 32G is allocated for heap FWIW), we had I think 6 of these nodes for a total of 384GiB memory.
Our storage was something like 18TiB across all nodes.
It would be a big machine, but our entire cluster could easily fit. Largest machine on the market right now is something like 128CPU's and 20TiB of Memory.
384GiB was available in a single 1U rackmount server at least as early as 2014.
Storage is basically unlimited with direct-attached-storage controllers and rackmount units.
1Us are extremely commodity, basically as “low end” as it gets, so I like to use them as if they are a baseline.
A 1U that can take 1.5TiB of ram might be part of the same series of machines that might have a 4U machine that could do 10TiB. But those are hugely expensive. Both to buy and to run
I have to teach developers that yes, we can have a 500MB data cache in ram, and that’s actually not a lot at all.
I was a SQL Server DBA at Cox Automotive. Some director/VP caught the Hadoop around 2015 and hired a consultant to set us up. The consultant's brother worked at Yahoo and did foundational work with it.
Consultant made us provision 6 nodes for Hadoop in Azure (our infra was on Azure Virtual Machines) each with 1 TB of storage. The entire SQL Server footprint was 3 nodes and maybe 100 GB at the time, and most of that was data bloat. He complained about such a small setup.
The data going into Hadoop was maybe 10 GB, and consultant insisted we do a full load every 15 minutes "to keep it fresh". The delta for a 15 minute interval was less than 20 MB, maybe 50 MB during peak usage. Naturally his refresh script was pounding the primary server and hurting performance, so we spent additional money to set up a read replica for him to use.
Did I mention the loading process took 16-17 minutes on average?
You can quit reading now, this meets your request, but in case anyone wants a fuller story:
Hadoop was used to feed some kind of bespoke dashboard product for a customer. Everyone at Cox was against using Microsoft's products for this, while the entire stack was Azure/.Net/SQL Server...go figure. Apparently they weren't aware of PowerBI, or just didn't like it.
I asked someone at MS (might have been one of the GuyInACube folks, I know I mentioned it to him) to come in and demo PowerBI, and in a 15 minute presentation absolutely demolished everything they had been working on for a year. There was a new data group director who was pretty chagrined about it, I think they went into panic mode to ensure the customer didn't find out.
The customer, surprisingly, wasn't happy with the progress or outcome of this dashboard, and were vocally pointing out data discrepancies compared to the production system. Some of them days or even a week out of date.
Once the original contract was up, and time to renew, the Hadoop VP now had to pay for the project from his budget, and about 60 days later it was mysteriously cancelled. The infra group was happy, as our Azure expenses suddenly halved, and our database performance improved 20-25%.
The customer seemed to be happy, they didn't have to struggle with the prototype anymore, and wow, where did all these SSRS reports that were perfectly fine come from? What do you mean they were there all along?
At work, they run big jobs on lots of data on big clusters. The processing pipeline also includes small jobs. It makes sense to write them in Spark and run them in the same way on the same cluster. The consistency is a big advantage and that cluster is going to be running anyway.
And this isn’t even wrong, bc what they need is a long-term maintainable method that scales up IF needed (rarely), is documented and survives loss of institutional knowledge three layoffs down the line.
I’ve spent time in some of the largest distributed computing deployments and cost was always a constant factor we had to account for. The easiest promos were always “I saved X hundred million” because it was hard to argue against saving money. And these happened way more than you would guess.
Yeah obviously if you run hundreds or thousands of severs then efficiency matters a lot, but then there isn't really the option to use a single machine with a lot of RAM instead, is there?
I'm talking about the typical BigCorp whose core business is something else than IT, like insurance, construction, mining, retail, whatever. Saving a single AKS cluster just doesn't move the needle.
I think my original point was more in the “engineers want to do cool, scalable stuff” realm - and so any solution has to support scaling out to the n’th degree.
Organisational factors pull a whole new dimension into this.
Some times the killer feature of that data analytics pipeline isn't scalability, but robustness, reproducibility and consistency.
Sure, it's not.
But the only alternative to that is not building some monster cluster to process a few gigabytes.
You can write a good script (instead of hacking one together), put it in source control and pull it from there automatically to the production server and run it regularly from cron. Now you have your robustness, reproducibility and consistency as well as much higher performance, for about one-ten-thousandth of the cost.
> Hopefully this has illustrated some points about using and abusing tools like Hadoop for data processing tasks that can better be accomplished on a single machine with simple shell commands and tools.
Recently had an argument with a senior engineer on our team because a pipeline that processed several PB of data, scaled to +1000 machines and was all account a success was just a Python script using multiprocessing distributed with ECS and didn't use Spark.
That's not to say that fancy tools don't have their use; but people often forget how much you can do with a few simple commands if you understand a pipeline and how the commands work.
Another post I love is https://stackoverflow.com/questions/2908822/speed-up-the-loo... where the guy manages to speed up a function which will run for days to milliseconds.
- desktop computers are really powerful (Apple Mx, AMD Epyc etc.)
- software like Polars
> Consulting service: you bring your big data problems to me, I say "your data set fits in RAM", you pay me $10,000 for saving you $500,000.
> Small Data is when is fit in RAM. Big Data is when is crash because is not fit in RAM.
Our sarcastic team-motto is very much this: https://twitter.com/DEVOPS_BORAT/status/41587168870797312
[1]https://buy.hpe.com/us/en/compute/mission-critical-x86-serve...
[2]https://buy.hpe.com/us/en/compute/mission-critical-x86-serve...
IBM Power System E980 https://www.ibm.com/downloads/cas/VX0AM0EP (notably the E880 did 32TB in 2014)
SGI UV 300 https://www.uvhpc.com/sgi-uv-300
SGI UV 3000 https://www.uvhpc.com/sgi-uv-3000
Still: of course worthwhile to point out how oversized a compute cluster approach is when the whole dataset would actually fit into memory of a single machine.
We don't have hypotheses, experiments, and results published, of what a given computing system X, made up of Y, does or doesn't achieve. There are certainly research papers, algorithms and proof-of-concepts, but (afaict) no scientific evidence for most of the practices we follow and results we get.
We don't have engineering specifications or tolerances for what a given thing can do. We don't have calculations for how to estimate, given X computing power, and Y system or algorithm, how much Z work it can do. We don't even have institutional knowledge of all the problems a given engineering effort faces, and how to avoid those problems. When we do have institutional knowledge, it's in books from 4 decades ago, that nobody reads, and everyone makes the same mistakes again and again, because there is no institutional way to hold people to account to avoid these problems.
What we do have, is some tool someone made, that then millions of dollars is poured into using, without any realistic idea whatsoever what the result is going to be. We hope that we get what we want out of it once we're done building something with it. Like building a bridge over a river and hoping it can handle the traffic.
1) There are (practically) no consequences for bad software.
2) The rate of change is too high to introduce true software development standards.
Modern engineering best practice is "follow the standards". The standards were developed in blood -- people were either injured or killed, so the standard was developed to make sure it didn't happen again. In today's society, no software defects (except maybe aircraft and medical devices) are considered severe enough for anyone to call for the creation and enforcement of standards. Even Teslas full-self-driving themselves into parked fire trucks and killing the occupants doesn't seem enough.Engineers that design buildings and bridges also have an advantage not available to computers: physics doesn't change, at least not at scales and rates that matter. When you have a stable foundation it is far easier to develop engineering standards on that foundation. Programmers have no such luxury. Computers have only been around for less than 100 years, and the rate of change is so high in terms of architecture and capabilities that we are constantly having to learn "new physics" every few years.
Even when we do standardize (e.g. x86 ISA) there is always something bubbling in research labs or locked behind NDAs that is ready to overthrow that standard and force a generation of programmers into obsolescence so quickly there is no opportunity to realistically convey a "software engineering culture" from one generation to the next.
I look forward to the day when the churn slows down enough that a true engineering culture can develop.
Imagine what scenario we would be in if they laid down the Standards of Software Engineering (tm) 20 years ago. Most of us would likely be chafing against guidelines that make our lives much worse for negative benefit.
In 20 years we'll have a much better idea of how to write good software under economic constraints. Many things we try to nail down today will only get in the way of future advancements.
My hope is that we're starting to get close though. After all, 'general purpose' languages seem to be converging on ML* style features.
* - think standard ML not machine learning. Static types, limited inference, algebraic data types, pattern matching, no null, lambdas, etc.
There is no single development, in either technology or management technique,
which by itself promises even one order of magnitude improvement within a decade
in productivity, in reliability, in simplicity."
I like to over-simplify that quote down to: Humans are too stupid to write software any better than they do now.
We have been writing software for 70 years and the real world outcomes have not gotten a lot better than when we started. There are improvements in how the software is developed, but the end result is still unpredictable. Without thorough quality control - which is often disdained, and there is no requirement to perform - the result is often indistinguishable whether it was created by geniuses or amateurs.That's why I would much rather have "chafing guidelines" that control the morass, than to continue to wade through it and get deeper and deeper. If we can't make it "better", we can at least make it more predictable, and control for the many, many, many problems that we keep repeating over and over as if they're somehow new to us after 70 years.
"Guidelines" can't stop researchers from exploring new engineering materials and techniques. Just having standard measures, practices, and guidelines, does not stop the advancement of true science. But it does improve the real-world practice of engineering, and provides more reliable outcomes. This was the reason professional engineering was created, and why it is still used today.
But this is better than where we just came from. Not that long ago, you would build software by getting a bunch of wizards together in a basement and hope they produce something that you can sell.
If things feel worse (I hope) that's because the rest of us muggles aren't as good as the wizards that came before us. But at least we're working in a somewhat tractable fashion.
The mathematical frameworks for construction were first laid out ~1500s (iirc). And people had been doing it since time immemorial. The mathematics for computation started about 1920-30s. And there's currently no mathematics for the comprehensibility of blocks of code. [Sure there's cyclomatic complexity and Weyuker's 9 properties, but I've got zero confidence in either of them. For example, neither of them account for variable names, so a program with well named variables is just as 'comprehensible' as a program with names composed of 500MB of random characters. Similarly, some studies indicate that CC has worse predictive power of the presence of defects than lines of code. And from what I've seen in Weyuker, they haven't shown that there's any reason to assume that their output is predictive of anything useful.]
It may not have been clear in 2014, but it is now: Data scientists are not computer scientists or software engineers. So tarring software engineers with data scientists practices is really a low blow. Not that we're perfect by any means, but that data point you're drawing a line through isn't even on the graph you're trying to draw.
I was unlucky enough to brush that world about a year ago. I am grateful I bounced off of it. It was surreal how much infrastructure data science has put into place just to deal with their mistake of choosing Python as their fundamental language. They're so excited about the frameworks being developed over years to do streaming of things that a "real" compiled language can either easily do on a single node, or could easily stream. They simply couldn't process the idea that I was not excited about porting all my code to their streaming platform because my code was already better than that platform. A constant battle with them assuming I just must not Get It and must just not understand how awesome their new platforms were, and me trying to explain how much of a downgrade it was for me.
"We don't have calculations for how to estimate, given X computing power, and Y system or algorithm, how much Z work it can do."
Yeah, we do, actually. I use this sort of stuff all the time. Anyone who works competently at scale does, it's a basic necessity for such things. Part of the mismatch I had with the data scientists was precisely that I had this information and not only did they not, they couldn't even process that it does exist and basically seemed to assume I must just be lying about my code's performance. It just won't take the form you expect. It's not textbooks. It can't be textbooks. But that's not the criterion of whether such data exists.
Of course you have to spend some effort to actually get code this fast, and that probably isn't worth it for the one-shot job. But jobs like compression, video codecs, cryptography and that newfangled AI stuff all have experts that write code in this manner, for generally good reasons, and they can all ballpark how a job like this can be solved in a close to optimal fashion.
Money quote from https://www.bitecode.dev/p/hype-cycles:
> geeks think they are rational beings, while they are completely influenced by buzz, marketing, and their emotions. Even more so than the average person, because they believe they are less susceptible to it than normies, so they have a blind spot.
(The top reference I get for it is a spam-site that is trying to hype the name and sell a domain.)
Yet another trend in our fashion driven industry.
I guess I liked the other Mach better, but I expect neither to go anywhere.
Must be nostalgia. AI is much, much worse. And, even more importantly, not only it is a annoying buzzword, it is already threatening lives (see the mushroom guide written by AI) and democracies (see the "singing Modi" and "New Hampshire Officials to Investigate A.I. Robocalls").
Also both OpenAI and Anthrophic argued if licenses were required to train LLMs on copyrighted content, today’s general-purpose AI tools simply could not exist.
At least 90% when people mention wanting to use AI for something, I can at least see why they think AI will help them (even if I think it will be challenging in practice).
99% of the time when people talk about big data it is complete bullshit.
Let me know when I can do it locally.
there are 60tb ssd's out there.. you might even fit all of the 8tb on a given server
The point is not that there are no valid uses for Hadoop, but that most people who think they have big data do not have big data. Whereas your use case sounds like it (for the time being) genuinely is big data, or at least at a size where it is a reasonable tradeoff and judgement call.
To people's beliefs on this, here's a Forbes article on Big Data [1] (yes, I know Forbes is now a glorified blog for people to pay for exposure). It uses as example a company with 2.2 million pages of text and diagrams. Unless those are far above average, they fit in RAM on a single server, or on a small RAID array of NVMe drives.
That's not Big Data.
I've indexed more than that as a side-project on a desktop-class machine with spinning rust.
The people who think that is big data are the audience of this, not people with actual big data.
[1] https://www.forbes.com/sites/forbestechcouncil/2023/05/24/th...