Performance in Big Data Land: Every CPU cycle matters
eng.localytics.com
eng.localytics.com
So no, it's not "every" CPU cycle, it's the ones that scale with the highest dimension of your data that matter. Which is the same old story we have always had, save your energy for optimising the parts that matter, because the ones that matter probably matter orders of magnitudes more than the ones that don't.
This made me cringe. Whether a series of operations takes place in one transaction or many isn't something you can just turn on and off depending on what looks more expensive!
The article ended up suggesting more transactionality, which is generally good (although the reason given is not the important one, namely "you're less likely to have all your data completely ruined"), but if you make the process distributed and aren't careful about sharding you may end up trading average-case cost in network load for much worse worst-case cost due to lock contention and transaction failures.
Optimizing database access patterns at scale is hard, and blithely making major changes to things that impact correctness is not the way to do it.
In my experience, CPU is rarely the big issue when dealing with a lot of data (I am talking about tens of PB per day). IO is the main problem and designing systems that move the least amount of data is the real challenge.
You analyse algorithms in terms of IO access, and specifically access pattern. If you cannot make the algorithm in a scanning fashion, you're in for a bad time.
https://www.usenix.org/conference/nsdi15/technical-sessions/...
From the paper the following 3 quotes highlight exactly why they where CPU bound:
> We found that if we instead ran queries on uncompressed data, most queries became I/O bound
> is an artifact of the decision to write Spark in Scala, which is based on Java: after being read from disk, data must be deserialized from a byte buffer to a Java object
> for some queries, as much as half of the CPU time is spent deserializing and decompressing data
Yep, this is a great point. The data locality/reducing IO is huge, but the way things actually play out for us when data isn't segmented/partitioned properly, it chews up CPU/memory. This is a lot of why the post was geared around CPU usage: concurrency in Vertica can be a little tricky, and stabilizing compute across the cluster has paid more dividends than any storage or network subsystem tweaks we've made.
We're not at the PB/day mark, though, so there's definitely classes of problems we are blissfully ignorant on. :)
We humans are not very good at appreciating orders of magnitude. I usually explain it this way: if it takes you 1 hour to process 1M records, then 10M will take 10 hours, and 100M will take 4.2 days while 10B will take over a year.
Sorry about the attribution. I'm trying to find who controls the blog as we speak so I can have them add it. (I work at Localytics, but I'm not the author.)
We've gingerly explored flame graphs to understand Vertica behavior under load, and we still have a lot that we want to try and use it for. I'm not sure if it will make an appearance in a further post, but we've definitely used your perf_event/ftrace-based tooling. :)
Would love to see if the performance bump is highly significant on a much larger and complex data set.
Optimizing data types and minimizing locks seem like general optimization tips, I was hoping for more advanced techniques for 100B rows.
In reality, the change in data type probably optimized disk access more than it did number of CPU cycles. That can often be more of a bottleneck.
Reducing locking and using shorter data type seem inadequate for the "Big Data" scene.
http://codexpi.com/java-vs-cpp-performance-comparison-jit-co...
http://stackoverflow.com/questions/5641356/why-is-it-that-by...
http://beautynbits.blogspot.com/2013/01/performance-java-vs-...
A lot of big-data work involves pulling out struct fields from a deeply nested composite record, and then performing some manipulation on them.
50x is not unreasonable for C/C++ code that was OO and uses a data oriented approach instead.
JVM inlines virtual method calls as one of its optimizations. See: http://www.oracle.com/technetwork/java/whitepaper-135217.htm...
So, now they have a few million Java jockeys churning away and a few million person-decades of work put into their mud piles. When starting any new project, there isn't much question about how to build it: More Mud!
As an embedded developer where every cycle counts I have come up with the same question as the poster above why bother with such languages. If a switch processes packets at line rate with the use of ASIC's why not have some similar development in the world of big data.
The author actually indicates that every CPU cycle is important for code block that will be executed for each row. So once you optimize hot code blocks, you're good to go.
Modern CPUs have DRAM fetch time in the 100's of cycles. Any cache friendly algorithm is going to walk circles around something that plays pointer pinball instead.
Let's say we want to compile a predicate expression "bigintColumn > 4 and varcharColumn = 'str'". A generic interpreter would suffer from the addressed issues but if you generate bytecode for Java source "return longPrimitive > 5 && readAndCompare(buffer, 3, "str".getBytes(UTF8))" then you won't create even a single Java object the output is usually identical to C and C++.
Either way the average Java dev isn't going to be writing bytecode so I feel like C/C++ still has the advantage in performance cases.
Also bytecode has instruction sets for all primitive types, otherwise there wouldn't be any point to have these primitive types in Java language since it will also be converted to bytecode instructions.
There are solutions for all the addressed issues but they need to much work to implement in Java compared to C++. However, once you solve this specific problem (I admit that it's not a small one), there are lots of benefits of using Java compared to C++.
Hadoop core has some shockingly bad design choices (lots of disk IO), and no amount of layers on top of it is going to fix latency issues.
It has nothing to do with JVM "overhead" (which is mostly a myth, anyway).
But really I was using C++ as an example of something more fit for these types of projects than Java, it doesn't have to be only C++ of course.
The most commonly used definition deduced from reading online articles is, in terms of size
"More data than naively fits into the memory of my (midrange) laptop using a high overhead platform"
or in terms of speed
"More data per second than can be handled using the same naive database code we used in the 90s for our website's comment section"
For example 1TB of data won't fit in memory, but if all you need to do is a sequential read in under a day then it's not a problem.
BTW, some previous HN discussions along these lines:
"Don't use Hadoop - your data isn't that big " - https://www.chrisstucchio.com/blog/2013/hadoop_hatred.html and https://news.ycombinator.com/item?id=6398650
"Your data fits in RAM " - http://yourdatafitsinram.com/ - https://news.ycombinator.com/item?id=9581862
If it fits in memory, you can honestly apply normal algorithmic analysis and optimize for memory access and cpu cycles. Once it no longer fit in memory, you become severely limited by IO.
In fact, since vertica is column oriented, I don't think you can pad things easily.
Does the auto-commit add additional lock overhead for some reason ?