9 ms·
Migrating From MongoDB To Riak At Bump
- salsakran 14y agoI was reading along and nodding my head until I got to the 1000 line haskell program that handles issues stemming from a lack of consistency. I'm not exactly a SQL fanboy, but maybe ACID is kinda useful in situations like this and having to write your own application land 1000 liners for stuff that got solved in SQL land decades ago isn't the best use of time?
- timdoug 14y agoIt's a development vs. operations tradeoff; we'd rather write code during the day so that when a database machine falls over at night we don't get paged.
- salsakran 14y agoI sorta get that, but it's not like Postgres/MySQL/etc are operational nightmares. It's just that going from one barely proven DB to another barely proven DB in a migration that involves a boatload of application side code seems heavy handed from the outside.
- willbmoss 14y agoI'll agree they are not operational nightmares, but now that we're set up with Riak we can do things that I'm pretty sure your Postgres/MySQL/etc. setup cannot. 1. Add individual nodes to the cluster and have it automatically rebalance over the new nodes. 2. Have it function without intervention in the face of node failure and network partitions. 3. Query any node in the cluster for any data point (even if it doesn't have that data locally). I'm sure there's other things I'm missing, but the point made by timdoug is the key one. We're at a scale now where it's worth trading up-front application code for reliability and stability down the line.
- xb95 14y agoYes, Riak gives you certain flexibilities that make certain things easier. If a node dies, you generally don't have to worry about anything. Stuff Just Works. (Well, in the Riak case, this is true until you realize that you have to do some dancing around the issue that it doesn't automatically re-replicate your data to other nodes and relies instead on read repair. This puts a certain pressure on node loss situations that I find is very similar to traditional RDBMS.) But of your list, I have done all of these things in a MySQL system and for a comparable 1000 lines of code. 1. We implemented a system that tracks which servers are part of a "role" and the weights assigned to that relationship. When we put in a new server, it would start at a small weight until it warmed up and could fully join the pool. Misbehaving machines were set to weight=0 and received no traffic. 2. Node failure is easy given the above: set weight=0. This assumes a master/slave setup with many slaves. If you lose a master, it's slightly more complicated but you can do slave promotion easily enough: it's well enough understood. (And if you use a Doozer/Zookeeper central config/locking system, all of your users get notified of the change in milliseconds. It's very reliable.) Network partitions are hard to deal with for most applications more so than most databases. It is worth noting that in Riak, if you partition off nodes, you might not have all data available. Sure, eventual consistency means that you can still write to the nodes and be assured that eventually the data will get through, but this is a very explicit tradeoff you make. "My data may be entirely unavailable for reading, but I can still write something to somewhere." IMO it's a rare application that can continue to run without being able to read real data from the database. 3. In a master/slave MySQL environment you would be reading from the slaves anyway unless your tolerance for data freshness is such that you cannot allow yourself to read slightly stale data. I.e., payment systems for banks or things that would fit better in a master/master environment. Since the slaves are all in sync, broadly speaking, you can read from any of them. (But you should use the weighted random as mentioned in the first point.) ... Please also note that I am not trying to knock Riak. It's neat, it's cool, it does a great job. It's just a different system with a different set of priorities and tradeoffs that may or may not work in your particular application. :) But to say that it can do things the others can't is incorrect. Riak requires you to have lots of organizational knowledge about siblings and conflict resolution in addition to the operational knowledge you need. A similar MySQL system requires you to have a different set of knowledge -- how to design schemas, how to handle replication, load balancing, etc. Is one inherently better than the other? I don't think so. :)
- aphyr 14y agoit's not like Postgres/MySQL/etc are operational nightmares Every database is an operational nightmare. Just wait until MySQL segfaults every time it gets more than 4 SSL connections at once, or takes four hours to warm up its inno cache. Or Cassandra compaction locks up a node so hard its peers think its down, turning your cluster into a recursive tire fire of GC. Or a Riak node stops merging bitcask and consumes all your disk in a matter of hours. Or Mongo ... My point is that every database comes with terrifying failure modes. There is no magic bullet yet. You have to pick the system which fits your data and application.
- lwat 14y agoYou don't seem to have any experience with solid, mature RDBMSs like PostgreSQL or Oracle or SQL Server. MySQL sucks in so many ways and all the other names you mentioned are not RDBMSs. People running large numbers of (for example) MS SQL Servers will tell you that they are ridiculously robust and mature and do NOT randomly crash in the middle of the night. And even when your hardware lets you down your hot failover machine is ready to take over without skipping a beat. Magic bullet? These systems are the closest you can find in the world of software.
- aphyr 14y ago<sigh> Or when the Oracle SCN advances too quickly, causing random systems to refuse connections or crash outright. Just because software is large, supported, and mature doesn't mean it is free of serious design flaws. [Edit:] Don't get me wrong, Oracle's DBs are serious business; an incredible feat of engineering. If I could afford them I'd probably use them more often. But everything in this business is a cost tradeoff--in licenses, in dev time, in risk.
- aphyr 14y agoCan't reply any deeper, figure this is relevant... Databases are, in my experience and those of my peers, the aspect of a system most likely to crash and burn. On the surface, it looks like regular old bugs--every project has them. But there's a reason databases are so hard: they're the perfect storm of unreliable systems, bound together with complex and unpredictable code to try to present a coherent abstraction. Specifically: 1. DB's need to do disk IO for persistence, which pits them against one of the slowest, most-likely-to-fail components of modern hardware. 2. They must maintain massive, complex caches: one of the hardest problems in computer science in general. 3. They must be highly concurrent, both for performance on multicore architectures and to service multiple clients. 4. Query optimization is a dark art involving multiple algorithms and frankly scary heuristics; the supporting data structures can be quite complex in their own regard. 5. They need to send results over the network, an even higher latency and less reliable system. 6. They must present some order of events: whether ACID-transactional, partially ordered with vclocks, etc. Almost all these problems interact with each other: a logically optimal query plan can end up fighting the disk or network; concurrency makes data structures harder, networks destroy assumptions about event ordering, caches can break the operating system's assumptions, etc. Moreover, the logical role of a database puts it in a position to destroy other services: when it fails, shit hits the fan, everywhere. Recovery is hard: you have to replay logs, spool up caches, hand off data, resolve conflicts. Yes, even in mostly-CP systems like MSSQL. All of these operations are places where queues can overflow, networks can fail, disks can thrash, etc. Furthermore, the quest for reliability can involve distributing databases over multiple hosts or geographically disparate regions. This distribution causes huge performance and concurrency penalties: we're limited by the speed of light for starters, and furthermore by the CAP theorem. These limits will not be broken barring a revolution in physics; they can only be worked around. The only realistic approaches (speaking broadly) are CP (most traditional DBs) and AP (under research). CP is well-explored, AP less so. I expect the AP space to evolve considerably over the next ten years.
- lwat 14y agoThis doesn't really answer any of the questions I asked you but from what you wrote here you have to agree that going with a tried-and-true RDBMS that's been around for decades is a much better prospect than choosing one of the new, unproven NoSQL products.
- lwat 14y agoDo you mean the physical server crashing? You should have a hot spare replicated machine for that. If you mean the RDBMS failing randomly at night, that really doesn't happen on mature systems like PostgreSQL or MS SQL Server. I mean it could happen but it's as rare as hen's teeth.
- btilly 14y agoEvery time I hear someone make a claim like this I have to wonder how much experience they actually have with these systems in a demanding environment. Relational databases, all relational databases, have a disturbing tendency to fall over suddenly under load with little advance warning. They don't have a pleasant gradual failure mode. Instead some point of contention goes from 99% of capacity (at which point it takes very little load) to 101% of capacity (at which point everything falls apart). If you've never experienced this, then I'm pretty confident that you've never scaled a database to its limit. Which isn't that surprising. Most companies don't have sufficient load to make their database break a sweat. But once you've encountered the region where problems happen, life gets much more difficult.
- lwat 14y agoI've been using SQL Server 5 hours a day for almost 15 years now and yes, I've taken databases to the limits of the hardware many times. We've upgraded our SQL Servers' hardware many times due to increasing load. We've never had a scaling problem that we could not immediately and easily resolve. If your server falls apart it's usually because of a single bad query and with SQL Server it's super easy to determine which query is at fault and kill it, at which point everything is immediately back to normal with NO LOST DATA. You can even set limits to how much resources a query can consume if you want. As for the 'no warning' thing, all you need to do is monitor your server. It should not run at 99% CPU or IO capacity at peak times! If you know the limits of your hardware it's really not difficult to monitor the actual usage and plan your upgrades accordingly. It doesn't matter how bad things get you can rest assured that you'll end up with a consistent database once the dust settles. You can even do database restores up to an arbitrary point in time if you need to! We've fucked up in every way imaginable but we've never lost any data let alone an entire database. I have nothing but praise for SQL Server.
- coops 14y agofrom this answer and your blog post it appears you were not using mongodb replica sets. is this true?
- WALoeIII 14y agoEventual Consistency is a feature, not a limitation of Riak (and friends). It requires you to think about your application different, but it enables things that you could not do before. For example, you can now handle databases in multiple datacenters, reducing latency to the client.
- salsakran 14y agoUhm..... no. This is backwards. Multi-DC capability is a feature. Eventual Consistency is an explicit tradeoff in a desired characteristic (Consistency) to allow other features.
- jes5199 14y agowell, yes and no. When you violate Consistency in SQL, your write fails. If it's a rare race condition, then the error probably just bubbles up through your application as an exception. Perhaps, if resolving conflicts was not something that we avoided but something that we baked into our application design, then we would be more likely to write code that handled it gracefully.
- aphyr 14y agoIndeed, these problems were solved years ago, for single servers. Conflict resolution in distributed systems is a significantly more complex problem, and invariably requires tradeoffs specific to the app. Those 1000 lines are likely declarations like "Merge changes to this list via set union", "Merge changes to this set using 2P-Set CRDTs", and so forth.
- cbsmith 14y agoThe RDBMS space has been addressing how to do this with distributed systems for quite some time as well (a couple of decades at least), and at least the more sophisticated systems tend to support a fairly broad set of approaches to addressing this problem (and means of expressing those choices quite succinctly).
- rdtsc 14y ago> The RDBMS space has been addressing how to do this with distributed systems for quite some time as well So how do you do what they did with MySQL. Can you build an always available multi-master cluster that is available, partition tolerant and eventually consistent?
- aphyr 14y agoYes: SimpleDB, HBase, Cassandra, Riak, Voldemort.
- rdtsc 14y ago> Yes: SimpleDB, HBase, Cassandra, Riak, Voldemort. Flagged your post since you obviously didn't even bother reading 2 sentences before replying.
- aphyr 14y agoSorry, I must have misunderstood. You wrote: Can you build an always available multi-master cluster that is available, partition tolerant and eventually consistent? These are exactly the properties of Dynamo. Could you elaborate more on what you were looking for?
- wpietri 14y agoMySQL is circa 1 million lines of code. I like SQL engines for moderate data sets that fit nicely on one machine and well within the normal performance envelope. But even there I will often have to try a few different incantations and cross my fingers that one of them will perform reasonably because that's easier than trying to figure out what that 1 MLOC engine is up to. And I don't know anybody who does very large MySQL setups without a lot more hassle than that. For some things I'd much rather deal with 1KLOC that I had to write myself than the 1 MLOC that I'm scared to even start digging through.
- gbog 14y ago"MySQL is to database what PHP is to programming languages". Use PostgreSQL.
- AaronBBrown 14y agoBy that, you mean used effectively on some of the largest, most profitable websites in the entire world? ;)
- papsosouid 14y agoWhy do mysql and PHP apologists think "people have managed to succeed despite deliberately making things more difficult for themselves" is a compelling argument? Mysql and PHP didn't make them succeed, or even help them succeed.
- AaronBBrown 14y agoIt proves that the technologies are capable products at the most massive scales. Do you have evidence to support that Facebook, Tumblr, or Etsy would have been better off had they chosen different technologies (of course not)? Or that they made things more difficult for themselves? How would Facebook be improved by PostgreSQL? At scale, data is so massively partitioned that the fact that you don't have windowing functions is utterly irrelevant. I hate PHP as much as the next guy, but the facts are that it gets the job done. And that, in business, is what matters. The angst on here about MySQL is unfounded and is largely a symptom of groupthink.
- gbog 14y agoIt its entertaining to see those weekly stories about NoSQL disasters. Hopefully one or two people will learn one or two things in three process. Let's try one: don't judge technologies on their sex appeal: the SQL old lady will take better care of your data than these young siliconed dolls.
- _Lemon_ 14y agoI have decided on wanting to use riak as well. I was wondering if anyone had examples of how they used it with their data model? For example this article mentions "With appropriate logic (set unions, timestamps, etc) it is easy to resolve these conflicts" however timestamps are not an adequate way to do this due to distributed systems having partial ordering. The magicd may be serialising all requests to riak to mitigate this (essentially using the time reference of magicd) in which case they're losing out on the distributed nature of riak (magicd becomes a single point of failure / bottleneck). Insight into how others have approached this would be awesome.
- reiddraper 14y agoThere are a several ways to approach this. The simplest is to just take last-write-wins, which is the only option some distributed databases give you. For cases where this isn't ideal, you resolve write-conflicts in a couple ways. One way is to write domain-specific logic that knows how to resolve your values. For example, your models might have some state that only happen-after another state, so conflicts of this nature resolve to the 'later' state. Another approach is to use data-structures or a library designed for this, like CRDTs. Some resources below: A comprehensive study of Convergent and Commutative Replicated Data Types http://hal.archives-ouvertes.fr/inria-00555588/ http://hal.archives-ouvertes.fr/inria-00555588/ https://github.com/reiddraper/knockbox https://github.com/reiddraper/knockbox https://github.com/aphyr/meangirls https://github.com/aphyr/meangirls https://github.com/ericmoritz/crdt https://github.com/ericmoritz/crdt https://github.com/mochi/statebox https://github.com/mochi/statebox
- stock_toaster 14y agoAre there any connector libs that provide "simple" last-write-wins out of the box?
- aphyr 14y agoIt's not hard. def merge(siblings) { sort_by(siblings) { |s| s.timestamp } }.last Or, in Knockbox/Meangirls, strategies like LWW-set. Until your clocks get out of sync.
- timhaines 14y agoIf you're thinking about using Riak, make sure you benchmark the write (put) throughput for a sustained period before you start coding. I got burnt with this. I was using the LevelDB backend with Riak 1.1.2, as my keys are too big to fit in RAM. I ran tests on a 5 node dedicated server cluster (fast CPU, 8GB ram, 15k RPM spinning drives), and after 10 hours Riak was only able to write 250 new objects per second. Here's a graph showing the drop from 400/s to 300/s: http://twitpic.com/9jtjmu/full http://twitpic.com/9jtjmu/full The tests were done using Basho's own benchmarking tool, with the partitioned sequential integer key generator, and 250 byte values. I tried adjusting the ring_size (1024 and 128), and tried adjusting the LevelDB cache_size etc and it didn't help. Be aware of the poor write throughput if you are going to use it.
- rb2k_ 14y agoI had the same experience about throughput being a bit sub-par. For me it was a test on a single macbook pro with a regular 2.5" hdd. Which client did you use to write to riak? protobuf or http? Also: which language? did you use threading? Did you enable search?
- timhaines 14y agoWell, for the benchmark, I was using Basho's benchmarking tool which is erlang, and I was testing with protobuf. I had 5 concurrent clients running for the benchmark, but also tried with more and less, and got about the same results. Search wasn't in use on the test bucket. For my app, I'd integrated Riak using ruby.
- fsckin 14y agoThanks for mentioning Basho Bench. Looks slick. For anyone else interested, it's at: http://wiki.basho.com/Benchmarking.html http://wiki.basho.com/Benchmarking.html
- timhaines 14y agoThe benchmarking tool is very slick. Easy to configure for a variety of scenarios, and once you figure out how to install R it produces those pretty graphs.
- supo 14y agoRandom thought on proto buffers: OP is advocating using the "required" modifier for fields and touting it as an advantage in comparison to JSON. I would move the field value verification logic to the client, because it can cause backwards compatibility problems if you un-require it.
- clu3 14y ago@timdoug, could you share specific problems with Mongo that made|forced you switch to Riak please? "Operational qualities" are little vague
- timdoug 14y agoWe experienced some significant difficulties with sharding; the automatic methods of doing so only seemed to shard a single-digit percentage of our data. We've also encountered some wildly unexpected issues with master/slave replication and related nomination procedures. You're right that this post is vague with regard to those details; they would be a good candidate for a future blog post, but the desired takeaway from this one is that we're quite pleased with the performance and scalability that Riak provides.
- tlianza 14y agoI find these kinds of stories interesting, but without some feel for the size of the data, they're not very useful/practical. I've heard of Bump, and used it once or twice, but I don't actually know how big or popular it is. If we're talking about a database for a few million users, only a tiny percentage of which are actively "bumping" at any time, it's really hard for me to imagine this is an interesting scaling problem. Ex. If I just read an article about a "data migration" who's scale is something a traditional DBMS would yawn at, the newsworthiness would have to be re-evaluated.
- polynomial 14y agoI was actually going to make a joke about "if the number of people I know who actually use Bump is any indication, it's not clear they even need a large data store."
- yahelc 14y agoThey celebrated 80 million installations as of 2 months ago (up from 50 million 8 months ago). http://blog.bu.mp/introducing-bump-pay-a-new-project-and-app-fr http://blog.bu.mp/introducing-bump-pay-a-new-project-and-app... That's a growth rate of 5 million installs a month; if they kept up that pace, they're at 90 million installs. To put that in perspective, Instagram "only" has 50 million users. http://www.quora.com/Instagram/How-many-users-does-Instagram-have http://www.quora.com/Instagram/How-many-users-does-Instagram... More bump data here: http://bu.mp/static/images/infographic_9-2011_6.pdf http://bu.mp/static/images/infographic_9-2011_6.pdf I'm not a user, but it seems like they have serious data.
- heretohelp 14y agoEven at 90 million users, with anything approaching a reasonable level of activity, we're not talking about serious data. 90 million rows of denormalized data isn't a big deal, and if I had to guess, their ops per second is probably no higher than what a dedicated single, or maybe a small master-slave postgres deployment could handle. Again, something a DBA would yawn at. And I say this as someone who scaled up an API for a service that plugged into multiple ad networks concurrently for a total of billions of impressions per month with a high level of reliability. Using NoSQL and an RDBMS combined. People who want to preach the NoSQL message should probably have some actual experience. Otherwise, it just makes very viable NoSQL solutions look really bad.
- stephen 14y ago> During the migration, there were a number of fields that should have been set in Mongo but were not Imagine that...this fascination with schema-less datastores just baffles me: http://draconianoverlord.com/2012/05/08/whats-wrong-with-a-schema.html http://draconianoverlord.com/2012/05/08/whats-wrong-with-a-s... I'm sure schema-less datastores are a huge win for your MVP release when it's all greenfield development, but from my days working for enterprises, it seems like you're just begging for data inconsistencies to sneak into your data. Although, in the enterprise, data actually lives longer than 6 months--by which time I suppose most start ups are hoping to have been bought out. (Yeah, I'm being snarky; none of this is targeted at bu.mp, they obviously understand pros/cons of schemas, having used pbuffers and mongo, I'm more just talking about how any datastore that's not relational these days touts the lack of a schema as an obvious win.)
- timdoug 14y agoYeah -- that's one of the huge benefits of marking a field as ``required'' in a protobuf. The ability to enforce a contract prevents a ton of unexpected and incomplete data making it on disk (and also, e.g., across the wire to clients). Having strict types represented in the serialization format is also handy; when one pulls out an int32 from a protobuf it's going to be an int32, and not an integer that somehow found its way into being a string.
- gizzlon 14y agoSure, but it could still be wrong in sooo many other ways.. (not arguing either way, just saying ;)
- lucaspiller 14y agoCould you elaborate more on how you have used Protobuffs, as I'm not sure I fully understand. I've previously used Riak in Erlang and Ruby projects, so am fairly familiar with how it works. It exposes a HTTP and Protobuffs API which allows you to store objects of arbitrary types (JSON, Erlang binary terms, Images, Word Documents, etc). From the sounds of it you are serialising a Protobuffs packet and sending this as the content of the object. Why did you choose this, over say JSON, which MongoDB uses?
- gizzlon 14y agoWould be interesting to see a follow up in 6 months or so.. It doesn't seem fair to compare [old tech] with [new tech] when you've felt all the pitfalls with one but not the other.