3 ms·
I have been involved in using a system built like this. All I can say is... It feels like you're building a database out of an event stream. A shitty one at th
by slackingoff2017 9y ago
I have been involved in using a system built like this. All I can say is... It feels like you're building a database out of an event stream.
A shitty one at that... Basically the write log part, only without a way to apply that state reliably like a real database. So you have to keep the log around basically forever. It's like you're in the middle of a DB recovery all the time.
After insane amounts of research and deep thought my personal opinion is that this is the wrong way to do scalable systems. Event sourcing and eventual consistency are taking industry for a ride in the wrong direction.
In my quest to find a better way I found some research/leaks/opinions of Googlers, and I think they're right. Even Netflix admits that using eventual consistency means they have to build systems that go around and "fixup" data that ends up in bad states. Ew. Service RPC loops in any such systems are Pandora's box. Are these calls getting the most recently updated version of the data? Nobody knows. Even replaying the event log can't save you, the log may be strongly ordered but the data state between services that call each other is party determined by timing. Undefined behavior.
You'll notice that LinkedIn/Netflix/Uber etc all seem to be building their systems using this pattern. Who is conspicuously absent? Google. The father of containers, VM's, and highly distributed systems is mum.
Researching Google's systems gives some fascinating answers to the problem of distributed consistency, a solution I'm stunned hasn't seen more attention. Google decided as early as 2005 that eventually consistent systems were too hard to use and manage. All of their databases, BigTable, MegaStore, Spanner, F1... They're all strongly consistent in certain ways.
How does Google do it? They make the database the source of truth. Service RPC calls either fail or succeed immediately. Service call loops, while bad for performance, produce consistent results. Failures are easy to find because data updates either succeed or fail immediately, not in some unbouded future time.
The rest of the industry is missing the point of microservices IMO. Google's massively distributed systems are enabled largely by their innovative database designs. The rest of the industry is trying to replicate the topography of Google's internal systems without understanding what makes them work well.
For microservices to be realistically usable for most use cases we need someone to come up with decent competition to Google's database systems. When you have a transactional distributed database all the problems with data spread across multiple services goes away.
HBase was a good attempt but doesn't get enough love. A point missed in the creation of HBase, that becomes clear when reading the papers about MegaStore and Spanner, is that it wasn't designed to be used as a data store by itself. Instead, it has the minimal features needed to build a MegaStore on top of it. The weirder features of HBase/BigTable (like keeping around 3 copies of changed data, and row level atomicity without transactions) are clearly designed to make it possible to build a database on top of it.
Unfortunately nobody thus far has taken up that challenge, and outside Google were all stuck with shitty databases that Google tossed away a decade ago.
- jamesblonde 9y agoGreat insightful comment. I came to the same conclusion a number of years ago. We did something about it - we built a new Hadoop platform around a not very well known distributed, in-memory, open-source database - MySQL Cluster (NDB). It is not the MySQL Server you think you know. It is an in-memory OLTP engine used by most network operators as a call subscriber DB. It can handles millions reads or writes/sec on commodity hardware (it has been benched at 200m reads/sec, about 80m writes/sec). It has transactions (read committed isolation level) and row-level locks. It supports efficient cross-partition transactions using one transaction coordinator per database node (up to 48 of them). You can build scalable apps with strong consistency if you can write apps with primary key ops and partition-pruned index scans. We managed to scale out HDFS by 16X with this technique. Since then, we have been doing like you suggested - we built a microservices architecture for Hadoop called Hopsworks around the transactional distributed database. All the evils of eventually consistency go away - systems like Apache Ranger/Sentry become just tables in the DB. More reading is available here: http://www.hops.io/?q=content/news-events http://www.hops.io/?q=content/news-events
- teacpde 9y agohttps://dataworkssummit.com/munich-2017/sessions/breaking-the-1-million-opssec-barrier-in-hops-hadoop/ https://dataworkssummit.com/munich-2017/sessions/breaking-th...
- psandersen 9y agoHopsworks looks like it might be exactly what I need, I do typical data science work for small to small-medium data and wanted to start properly playing with spark on a HDFS store. Currently most work is just done in R/Python in VM's on a small proxmox cluster (where only 1 node is always on) but I'd like start gently moving to spark, run the stack on a single node and scale on demand. Is Hopsworks for me, does this approach even make sense for such small data or am I crazy? Thanks for your response!
- jamesblonde 9y agoYes, Hopsworks can run on anything from 1 server to 1000s. We are finalizing the first proper release now - Jupyter support, tensorflow, pyspark, sparkr, python-kernel for jupyter too,