Designing Data-Intensive Applications
dataintensive.net
dataintensive.net
The architecture seems to solve two big problems
* Scaling RDBMS (there are solutions like Cloud Spanner but they rely heavily on Google's proprietary network for low latency and are expensive as balls)
* Keeping downstream systems in sync. A lot of companies have gone with a "Tail the RDBMS" of some kind, e.g. writing the MySQL binlog to a Kafka queue and having downstream consumers read from that, but this seems like a more elegant solution.
Are there any examples or experiences of people working with systems like this? What are some downsides, challenges, and actual benefits?
It's a balancing act between two extremes - locking everything down and ensuring the tx has been committed on all nodes/propagated to all consumers on one hand and sending an "ack" back to the client with a loose promise of eventual consistency on the other.
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.
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!
The second major issue is technical. Building something like Spanner requires a very accurate time source that absolutely will not skew. Ever. This is how Google avoids partitions and essentially breaks CAP theory. Perfect time gives you globally accurate timestamps without exchanging data. Distributed transactions without locks or shared state, just usually benign contention.
They're not that hard to build, just no demand. Possibly some issues with ITAR preventing such accurate clocks from becoming commodity hardware? It could easily lead to extremely accurate IMU's which are definitely limited. Not sure, but that's what I ran into researching fibre optic gyros. Atomic clocks could probably be built on SMT scale for a few cents IMO.
As usual, it seems Google is already doing this and has been for years. We either need to wait for the trickle down that Google thankfully does after about 5 years... Or get some deep pocketed tech behemoth to foot the bill for everyone else
From their 2017 paper on Spanner and CAP:
> To the extent there is anything special, it is really Google’s wide-area network, plus many years of operational improvements, that greatly limit partitions in practice, and thus enable high availability.
[0] https://static.googleusercontent.com/media/research.google.c...
I dislike to lose ACID. Mainly because my apps are all about financial/business stuff.
My ideal DB right know will be alike:
Incoming data ->
- Write WAL (disk)
- Convert to Logical JSON-like structure (memory). For Consisten API
- Pre-Triggers (BLOCK) <- Validations
- Persist on Table(s) (only the data necesary to perform validations later?. Cache?)
- Persist on READ-LOG (all the data!)
- Post-Triggers (NON-BLOCK!) Read by:
Log Listener(s) ->
- Build caches, (secondary?) indexes, etc (NON-BLOCK)Many software applications, especially web application servers, are effectively a set of data structures updated incrementally from an incoming data stream, and then served to users. The analogy of many applications to a database (or an interpreter) is an accurate one and, in my personal experience, useful as well.
> So you have to keep the log around basically forever.
Different logs can have different retention periods, and you linearize across different logs by using a single writer per timeline. Many domains allow for enough splitting of timelines to enable effective parallel processing without compromising consistency (consider the case of distinct customer organizations using a time-tracking product - there's no reason they need to be able to write to the same database, and volume/contention within any individual organization is likely to be low, allowing for maintained performance).
Any component sends its data into kafka-esque (we're using combination of NATS and PubSub) pipeline where series of workers read, process and write data into our RDBMS which is the ultimate source of the truth.
This allowed us to run a double scale system, where all components are running on their own pace and the RDBMS is running on its own. There is some inherited delay in the data propagation, but it works for service like our (search engine) that doesn't require real-time exposure of newly acquired data.
This design also allows for a frugality as the RDBMS cluster is only scaled based on long term trends and now short term bursts. We also are able to buy committed usage for the cluster as we've great predictability in its growth.
The transaction log encourages small logical "patches" (set a field, increment a number, replace a substring, move an array element, etc.) that are applied in sequence but can be disentangled by clients to generate a consistent UI, and also used to resolve conflicts between multiple distributed writers. You can also follow the log through the gRPC and HTTP APIs, and you can register "watches" on queries that cause matching changes to trigger an event.
While the transaction log is the underlying data model, we also maintain a consistent view of the current version so that you can use it as a document database with classical CRUD operations. So on the surface it's a lot like Firebase or CouchDB, except you get things like joins and schemas.
Drop me an email (see profile) and I can send you some links.
*[_type == "blogpost" && published == true] {
_id, title,
"bodyExcerpt": regexReplace(body, "(.+?)(\. |$)", "\\1"),
author -> { name, category },
"sectionNames": sections -> name
} | order(createdAt desc)[0..20]
I don't know much about N1QL, I will have to read more about it.I can imagine it may not dive deep enough for people who really understand the internals of a given data store and the content is probably available elsewhere. However this book is a thoughtful and engaging curation of a lot of information that I may have missed otherwise
That's a great review - I learned about the book recently, and it sounds like exactly what I need right now, to make a more informed decision about database choices.
It will not appeal to the absolute novice to be sure. But for anyone else who has worked on systems for moving data (ETL, streams) and storing data (databases and other data stores), this book will show you how things (probably stuff you've done bits of pieces of) fit together and expound on the few foundational big ideas that makes everything cohere. Once you've understood that, you are on your way to designing data systems that are much cleaner and more scalable.
My experience reading this book is a bit like that of a tradesperson going back to school to learn theory, and after being enlightened, coming away with a new understanding of how to put together theory and practice to better his craft.
I chanced upon this book through an excellent interview Martin Kleppmann did on Software Engineering Daily podcast. If you want the talk-show-host cliff notes version of what the book is about, you should listen to this particular episode:
https://softwareengineeringdaily.com/2017/05/02/data-intensi...
Does anyone know books that are similar in style? (conceptual, showcasing different solutions to problems and their tradeoffs, high signal-to-noise)
A few shorter books have come out that try to touch different approaches that I've liked: "Seven Languages in Seven weeks" -- and the series has also gained 7 databases, concurrency models and web frameworks.
Finally, there's this anthology where OSS authors described what they did in their applications, so there's a ton of practical information http://aosabook.org/en/index.html
Edit: I picked it up based off the recommendation of another HN commenter who said they picked it up after listening to Martin on the SWE Daily Podcast[0]
[0] https://softwareengineeringdaily.com/2017/05/02/data-intensi...
I was excited about this book because there is a gap in distributed systems books. I feel like there's are a large amount of blogs but most of the books available on amazon are text books and/or include heavy math.
Great book if you'd like to finally understand Jepsen.io articles.
I tend to read on the tube a lot and lugging a hefty O'Reilly tome around on my commute isn't ideal.
Currently the book is just sitting on a shelf, I'll get round to it one day!
I wrote a more detailed review at http://horia141.com/designing-data-intensive-applications-re... for the interested.
Thankfully, this book cites its sources and extensively documents references, and the author even maintains the reference links[0].
Would very much appreciate if someone has some resources to link or case studies.
As with another poster, Edward Tufte's books came to mind - though it's about visual presentation of information, not user interface/experience design.
I've also felt that there's an unmet demand for books that provide a thorough overview of UI/UX design patterns, especially the way this book (Designing Data-Intensive Applications) does for its domain.
It is an overloaded term for sure, but the title of the book caught my attention, it is only when I read the article that I realized it was talking about the design of application implementations, not the UX.