Runaway complexity in Big Data... and a plan to stop it
slideshare.net
slideshare.net
Append only immutable DBs are inherently better just because they essentially add a 4th dimension to your data (all history, all data, all the time), and hook up to the fact that Von Neuman machines love splitting things up into constrained sub-streams/problems and crunching through the entire data set in memory across clusters using discrete memory frames/units of work (think Google search index/GPU framebuffers/Integer arithmetic).
Storage is now unlimited and random seeks are very expensive. The more you serialize your processing and split them into in memory units that utilize cache locality - the faster you'll perform. You'll turn performance problems into throughput problems - all you have to do is to keep the data flowing.
RDMS databases are dead for the same reason I don't manage my YouTube bandwith use or bother managing files or bother managing RAM - I have cable now with terabyte hard disks and 32-64 GBs of RAM.
The future is immutable with periodic/continuous historical stream compaction (pre-computation) with queries running across hundreds of clusters and reducing all searches to essentially linear map-reduce (sometimes with b-trees/other speed ups) + a real time dumb/inaccurate stream layer.
1 machine cannot store the Internet - a million machines can. 1 machine cannot process the Internet - but a million machines operating on one 64GB slice of it located in RAM and operating with cache friendly memory-local discrete operations can run through the Internet millions of times a second.
1) I can't afford a million machines or hundereds of clusters or keeping petabytes of data around in a way that doesn't make the data completely useless. Yes mutable state causes complexity, but I can't afford to rid myself of that problem by using brute force.
2) Performance is not replaceable by throughput if you have users waiting for answers on questions they just dreamt up a second ago and you have new data coming in all the time.
3) Cache locality and immutability don't go well together. Many indexes will always have to be updated, not just replaced wholesale with an entirely new version of the index.
2) What kind of queries do you have that take that long? If it's scientific computing - no way around it. If it's just DB slicing/aggregating - if you use last gen structures you'll get last gen perf.
3 ) Cache locality on repeatable computable units on local in memory data do go better - continuous batched background updates to indices with a real time layer like the OP suggested will become the new standard.
I am indeed operating in an environment of limited resources. It's called "The Real World". Reading marketing language like "last gen" makes me lose interest in a debate very quickly.
Just because you happen to be constrained by your resources it does not follow that your world is any more "Real" than mine - it's just different - which is what I said.
What part of "last gen" is marketing speak? Original Xbox is "last gen", the iPod is "last gen" - quite literally the last generation.
The original Xbox is no longer available in the stores. That's because it has been superseded by the new model, which does everything the old one did, only better! Poor me who just doesn't have "the capital" to get myself the shiny new one just yet.
Apparently, that's the way you want me to think about immutability versus mutability in data structures. Makes no sense.
2) I think you misunderstood what was meant here. His point was that your limiting factor really becomes throughput if you have the right architecture.
3) Indexes and cache locality also tend not to play well together. ;-)
> "Activities of users at terminals and most application programs should remain unaffected when the internal representation of data is changed and even when some aspects of the external representation are changed." http://www.seas.upenn.edu/~zives/03f/cis550/codd.pdf
It's not clear from the slides how RDBMSs manage to conflate these two concerns in practice.
It's also unclear to me how the "Lambda Architecture" differs from what we've been calling [soft real-time] materialized view maintenance for decades.
It should be easier than ever, particularly if you cheat and declare rotational media an unsupported legacy format.
This word, I don't think it means what you think it means.
They will not always work for big data, which often means "max out our storage with crap, don't worry disk is cheap" eg web metrics
When a decision maker goes "don't worry disk is unlimited", the resulting application is prone to maxing out storage.
Whenever your application maxes out your storage, you have no space for previous versions of the data.
I'm not going to say anything better than it was said in the slides / book draft, so I'll just encourage you to take these techniques seriously... they're born out of necessity, not because they sound like fun, and real people are using them to solve problems that are hell to solve any other way.
That said, these are not problems that everyone has. If you're not nodding you're head along with the mutability / sharding / whatever complaints at the beginning of the deck, you can probably still get by with a more traditional architecture.
(Also, rereading... I should probably note that not everything needs to be kept forever; only the source data, since the views can be recomputed from them at any time. That makes things a bit cheaper.)
'Bigness' of data != data size
'Bigness' of data == data size / budget
Twitter isn't a typical company. I assume they have both a budget and competent management that will let them get away with something like the Lambda architecture.
I reckon it's a lot harder to scale to even a terabyte under the constraints of a grubby setting like a datawarehouse for some instrument monitoring company.
Those guys will allow at best MS SQL for storage, and won't mind putting their developers through hell.
Having worked on this kind of stuff myself, I'd have to argue the exact opposite. I've always ended up building something precisely like what is described when trying to tackle those kinds of problems at large scale.
I know that in places I worked, Materialized views were mostly limited in the application realm because too many of them over enough data brought the DB to its knees.
Of course, to be fair, you're probably already looking at a rearchitecting effort at this point!
However, looking at the architecture diagram, much of the complexity could be hidden behind the scenes with a devoted toolset built on top of Postgres, rather than trying to cobble together Kafka, Thrift, Hadoop, and Storm. Sometimes, one big tool beats a lot of small tools.
For that matter, I wonder if a series of Postgres-XC (or Storm) servers couldn't do the same thing without learning a series of complex tools.
Step 1: Everything goes through a stored procedure. Deletes, updates, and creates are code generated through a DSL. Migrations would be a pain under this system, but mitigated by the fact that the complexity is being handled by the tool.
This stored procedure then writes to the database of record and then sends a series of updates to, essentially a system of real-time materialized views that serve the same purpose as the standard NoSQL schema.
The lambda architecture purposes would still be fulfilled with lower data-side complexity. After I read through Big Data I'll revisit this with a more nuanced view, but I wonder if hadoop and nosql really give you anything.
Basically, you store all of the data as events (event sourcing) and create separate query views (projections) which are populated from event streams.
tldr - the hand me down from the 90s can be good enough if you are not big yet
At my previous^2 job, we had a distributed event store that captured all data and could give you all events (or events matching a very limited set of possible filters) either in a given time range, or streaming from a given time onwards. For any given view we'd have four instances of the database containing it, populated their own streaming inserters; if we discovered a bug in the view creation logic, we'd delete one database and re-run the (newly updated) inserter until it caught up, then repeat with the others (queries automatically went to the most up to date view, so this was transparent to the query client - they'd simply see bugged data (some of the time) until all the views were rebuilt).
The events system guaranteed consistency at the cost of a bit of latency (generally <1 second in practice, good enough for all our query workloads); if an event source hard-crashed then all event streams would stop until it was manually failed (ops were alerted if events got out of date by more than a certain threshold). This could also happen if someone forgot to manually close the event stream after taking a machine out of service (but at least that only happened in office hours); hard-crashes were thankfully pretty rare. Rebuilding the full view after discovering a bug was obviously quite slow, but there's no way to avoid that (and again this was quite rare).
In use it was a very effective architecture; we handled consistency of the events stream in one place, and building the view in another. We only had to write the build-the-view code once, we only had one view store and one event store to maintain. And we built it all on mysql.
Secondarily, if this were a service, wow!
One of the problems is servicing two workloads, one for "transactional" processing and another for analysis.
For transactional systems you need to be able to change things quickly and consistently. For analysis you need to be able to query lots of data quickly.
For decades people have realized that these are two separate workloads, so have built 2 systems, a transactional system (on an RDBMS) and a data warehouse (generally on an RDBMS). Data is then shipped between the 2 in batch jobs.
The transactional system is normalized, and the data warehouse is normalized. Within the data warehouse you make denormalized copies of the data that fit the reporting workload required so that as much of the workload is pre-computed as possible.
The problem is that as reporting requirements change, you need to modify these pre-computed stores, as they are very heavily tuned for the particular reporting requirement. Building pre-computed stores is generally done in SQL and can be challenging as you are generally shifting a lot of data and you are trusting the RDBMS to get its optimizations right.
There is a trend now to use Hadoop for the building of these pre-computed stores (and even to use Hadoop for the entire data warehouse). However, writing map-reduce jobs for queries is cumbersome compared to SQL so your productivity suffers. But you don't have to pay Oracle or IBM for licenses.
The key problem is that you need to pre-compute stuff to do fast aggregation, but you can't pre-compute everything. So what you pre-compute is dependent on what your users want, and that changes all the time.
So what you want is a system that lets you change what is pre-computed easily and efficiently.
Why don't you just: 1) Master your data in "normal form" (graph/entity-based data models like Datomic or Neo4j are perfect for this) in a scalable database (e.g., Cassandra) 2) Use Hadoop for offline analytics 3) Index all of the data needed for realtime queries in a search engine (e.g., elasticsearch), smartly partitioning the data for performance
As for "schema" changes -- you don't need them in an entity-based model because you can always be additive with the attributes in your schema. (You may wish to have a deprecation policy for some attributes but doing this allows you to as slowly as you'd like update any applications and batch jobs that depend on deprecated attributes.)
The Lambda Architecture is the only design I know of that:
1) Is based upon immutability at the core.
2) Cleanly separates how data is stored and modeled from how queries are satisfied.
3) Is general enough for computing arbitrary functions on arbitrary data (which encapsulates every possible data system)
For example, if you want to know the sum of the elements in a tree, but you've already calculated the sum of the elements of another tree that differs on only one path, then you won't need to do much work, you can use the precomputed work for most of the tree and just perform calculations for the parts that are different.
To find the median, you need the data sorted. You can express something like mergesort as a catamorphism. You take the lists of the child nodes and the singleton list of the value at the node and merge them together. However you have to do a linear amount of work at each node and store a linear amount of data, so it's not really within the spirit of mapreduce because it doesn't scale well, although it is possible.
Catamorphisms are operations on trees, so the first problem with neural networks is that you'd have to embed your neural network in a tree. In the worst possible case, a fully connected network, to find the next value of each neuron, you'll have to find the current value of every other neuron. The problem is that you don't have access to "cousin nodes'" data, so it isn't possible.
My initial reaction is to close the window immediately due to the font on the slides. I don't know if that is the fault of my browser (chrome 21.x on this PC), slideshare, or the presenter who put together the slides, but the staggered characters drive me bonkers.
Edit: scratch that - each letter shows up as a separate span in the generated HTML. There's no sane way to typeset that.
What I got out of it was "Just store the raw data", always compute results/queries, BigTable and Map/Reduce are cool. I felt like I missed something perhaps someone here can help.
Basically the 'lamba architecture' he refers to is event sourcing, or write-ahead logging, but with scalability in mind, and some cool hooks for maintaining correctness.
You use your hdfs store as your event log, and a couple layers to handle making the batch processing (map-reduce jobs) into real-time queryable databases (along with disposable caches that can be updated real-time to handle the stuff that comes in between batch-jobs in a real-time way).
The goal is to never lose the raw actions - so even updates to various layers (including the batch processing!) don't result in data corruption, just some time to reprocess all the raw inputs again.
As to the content, see the papers on data flow architectures from the 70's and 80's [1]. They are very cool. We've done something similar at Blekko where we store raw data in a table structure and build in pre-computed results with combinators [2]. The Map/Reduce paper [3] is an excellent introduction to a number of these concepts. This is all good stuff and something that is helpful for people to have in their toolboxes. The title of the post gave me the impression that there was something new here (I'm always on the prowl for new stuff on these problems) and I didn't see what the new stuff was, it seemed like the stuff we know just presented more coherently rather than as a collection of links. Perhaps that is more clear, perhaps not.
[1] https://en.wikipedia.org/wiki/Dataflow
[2] http://highscalability.com/blog/2012/4/25/the-anatomy-of-sea...
In particular, what do I do when there is erroneous historical data that violates the new schema (the newly discovered constraint that was there all along)?
Before any data is written to the "data" files of the database, it is written to a "log" file sequentially. Once that has succeeded, it will write to the data file, then write to the log file again to say that the write was successful.
That way if the database breaks during the write, the system knows about it. It's also useful because it allows you to restore the database to a point in time, as you have a full log of all the transactions run (and you know how to un-run them).
In this case, they just make the transaction log a little more accessible.
For your example, erroneous historical data would remain in the system but with a new record being added indicating that at a certain date and time the old record was replaced with the new one.
Basically if you think of a store processing a purchase for $100 and refund of $100 as 2 transactions (for $100 and -$100) rather than simply deleting the first transaction.
This is important because things may have happened between the deposit of money and its refunding (e.g. interest payments).
No, MapReduce is a framework for computing catamorphisms over trees.
The Lambda Architecture, which we will be introducing later in this chapter, provides a general purpose approach to implementing an arbitrary function on an arbitrary dataset and having the function return its results with low latency.
Thanks for the post...slideshare is annoying to read (on an iPad) but the content was digestible. Mostly, it helped remind me to download the update of the OP's book