Online migrations at scale
stripe.com
stripe.com
We "allowed" small amounts of downtime previously (say 10-15 minutes at a time for migrations), but over the last year as our customer base has expanded, the window has shrunk smaller and smaller as to avoid service disruptions.
We're now at a stage where downtime is effectively not allowed anymore, so we've taken to this approach (make a new table, dual write new records, then write all the old records, remove references to the old table) to mitigate our downtime.
It's nice to know that other companies have taken this approach as well, my team honestly didn't know if we were doing it "properly" or not (if there is such a thing as properly)
Incidentally, I was thinking about doing a blog post about that too.
That's good. You're confirming there's an audience for my writing. Thank you :D
Event sourcing isn't without it's difficulties but if the domain is pretty well understood then it can be a very powerful pattern.
Stripe has an awesome product and seemingly an equally great engineering culture. I wonder if they have debated using CQRS/ES.
A company I worked for some years ago, LiveOps, also followed this methodology, because we had to.
LiveOps was and is a 'telephony in the cloud' provider, and the databases are in-line with some of the important call flows.
Last I looked, in the ten years previous, LiveOps had handled at least one call from something like 15% of all of the possible phone numbers in the United States, and all of that data was available in a huge replicated mysql cluster spanning multiple datacenters. We had mysql tables with tens of billions of rows, replicated across many servers.
And just to be crystal clear: there was absolutely no downtime allowed, ever, because downtime of any kind meant that phone calls would not go through.
We used extensive feature flags, along with the techniques described in the stripe.com article, to achieve exceptional uptime, relatively ok operational overhead and pretty good development velocity.
The point I'd like to drive home is that this is all possible even for medium sized organizations and moderately sized staffs. The important part is that everyone needs to keep these priorities, methods and techniques in mind, all the time, to make it work.
Mongo-modelling being a black art of sorts (to me, at least) I'd be more interested in what the tipping point was for them (data size and shape, usage patterns) than the relatively straightforward (conceptually at least, operatinally these things are never to be underestimated) table-doubling approach to changing a data model.
https://github.com/soundcloud/lhm
I've used it in the past a lot. You have to make sure to throttle the migration, or else it'll max out your CPU & memory.
Disclaimer: I worked on LHM a few years ago, still one of my favorite projects :)
Its main advantage seems to be that it doesn't require triggers on the original table (it works by subscribing to the replication log), thus not making write locks worse during the migration.
I haven't tried either yet, so if anybody has some experience with them, do please share a comparison!
It has a nice webapp to monitor and run schema migrations, including peer review. I use it all the time at Square, and it's pretty great.
1) Write code that works with current and new schema
2) Deploy that code
3) Run that version for long enough to know you will not roll it back
3) Run the migration, moving the db from current to new schema
4) Clean up code to only work with new schema
5) Deploy that code
Where step 3 uses an appropriate tool for online schema changes for your database.
Here a some writeups of that approach I found when trying to explain this to a friend who's a rails developer:
* http://blog.honeybadger.io/zero-downtime-migrations-of-large...
* https://pedro.herokuapp.com/past/2011/7/13/rails_migrations_...
* https://www.rainforestqa.com/blog/2014-06-27-zero-downtime-d...
* https://blog.engineyard.com/2011/zero-downtime-deploys-with-...
1. Create new data model and adjust reads to read from the new model first and fall back to the old model. Update ALL your readers first.
2. Change writes to use the new model.
3. Now you have different options. You can leave things as they stand. You can actively migrate data from old to new. You can do a lazy migration, for example, only when your read old data you migrate it.
There are some gotchas depending on your system specifics, esp. races between writes from the migration process and writes coming from the application. How you address those depends on the specifics of your situation...
1. Lisa has old tv_service subscription, saved to snapshot, but not live-copied to Subscription table.
2. Lisa removes tv_service subscription.
3. Before snapshot is updated, batch processing adds lisa_tv_service to the Subscription table.
4. Subscription now has lisa_tv_service, but Lisa.subscriptions does not include tv_service.
My only idea: wait a day, then a batch-process job looking for data that exists in Subscription table but not in Customer.subscriptions
Well, one solution would involve a log table used to log/lookup the migration status: Anytime the real-time system touches data (eg. deletes the subscription in the source table and has nothing to do in the new table) the event is logged. The backfilling routine uses that data to decide wether its data or the current data is "fresher".
If the old subscription data is never (or for a certain period) deleted but merely marked inactive (eg. with fields "timestamp_end" and "timestamp_modified" updated) then that information can be used during the backfilling to decide on what data is more recent. [edit: I guess this is what haldean is reffering to above by "tombstone records."]
That makes sense, I think. Does it? Elsewhere in the thread i referred to the procedure used as "conceptually relatively simple". Most things are conceptually simple, but always with pitfalls and dark things lurking around the edges.
- https://blog.codeship.com/rails-migrations-zero-downtime/
- https://blog.engineyard.com/2011/zero-downtime-deploys-with-...
In the Ruby on Rails world, it's definitely a thing (made easier with attribute reader / writer, and db/migrations/)
It's pretty normal for tech companies to do dozens of deploys a day, thanks to CI/CD. There are lots of ways to achieve this, but they all pretty boil down to the same idea: Boot up the new version before shutting down the old version.
if (migration_complete) {
write_new_data_model();
} else if (migration_in_progress) {
write_old_data_model();
write_new_data_model();
} else {
write_old_data_model();
}
Reading Code: if (migration_complete) {
read_new_data_model();
} else {
read_old_data_model();
}
The `migration_in_progress` feature flag is flipped first, then the migration jobs run to backfill the new data model from the old data model, then the `migration_complete` flag is flipped. As many times as necessary, `migration_in_progress` can be flipped to off at any time and the new data model can be blown away and new writing code deployed. If there's no confidence in the code for reading the new data model, then a certain percentage of requests/users can be directed to read from the new data model while double-writing continues.That sounds scary! how many time do you need the "once again"?
The other option is the on-demand option you mentioned. The biggest issue I see is the latency to fetch that information is suddenly a lot more random. If it's not in DB B then you have to fetch from A and write into B before returning the information.
You can't rewrite the pointers to the merged rows while the application is online.
https://stripe.com/assets/blog/animate-svg-fe67ad3b4fe397689...
Here they were modifying the data. But yes, it works the same way, and it's not really novel in any way.
A lot of migrations would be pretty trivial if you could have even a couple hours of downtime, but in an expectation of 24/7 availability that is no longer acceptable in most situations.
I'd like to hear about alternatives if you've had experience.
In the past I have either flagged records to say where the system should read the data from, or just built logic into the readers to say if there is no data in A, read from B.
Then you update the writers to migrate the data from A to B on every new update and remove the data from A. It is an expensive 1 time write to move the data, but then you don't have to worry about keeping data in sync across two storage locations.
What you end up with is all records that are actively getting worked on move first. At a later date you start migrating all those stale records from A to B with a background process. Once A is empty, remove the logic to read from A, remove the migration writes, remove the A datasource.
Such great pains come with huge systems. What's the alternative?
Taking the platform offline for a few hours? Management will say no. Or maybe Management will say yes once every three years, severely limiting your ability to refactor.
Doing a quick copy, and hope nobody complains about inconsistencies? Their reputation would suffer severely.
From the article I got the impression that both tables were being written to in the same database transaction, so this is not a possible failure scenario at all.