8 ms·
DoorDash manages high-availability CockroachDB clusters at scale
- rickreynoldssf 3y agoI'm not really seeing why DoorDash needs all their operational data in one monster clustered database. I would think its so much simpler to shard the data by region for operational queries and aggregate in the background for long-term storage.
- esafak 3y agoManual sharding is a crutch and a pain. Just use a distributed database, and let the database company worry about it.
- killingtime74 3y agoResume driven development
- mbyio 3y agoCockroach automates the sharding of data by region and provides tools that let you use and manage it more like a traditional database. If they didn't use cockroach, they would have to write/setup tools and adapters to do all that anyway. It would probably be more familiar to developers conceptually if they used traditional sharding, but why build and maintain all that when you can just use off the shelf software?
- karmakaze 3y agoCockroachDB had "Follow the sun" multi-region data balancing, which then got generalized to "Follow the workload"[0]. [0] https://www.cockroachlabs.com/docs/stable/topology-follow-the-workload https://www.cockroachlabs.com/docs/stable/topology-follow-th...
- JohnBooty 3y agoI've never done geographic sharding but it seems kind of hard. How do you pick shard boundaries? How do you deal with entities who are near the boundaries and whose current operational data therefore spans >1 shards? (Imagine somebody at near the geographic intersection of like, five shards looking for pizza in a 10 miles radius or w/e) Also the majority of entities they're tracking (users, drivers) do not have fixed locations. Maybe it's not as hard as I'm thinking. I guess you just have to accept that any query can span an arbitrary number of shards and the results need to be union'd. I'm sure a lot of smart people have tackled this at the DoorDashes and Ubers of the world and maybe there's some optimal way of handling it. I would love to hear about that.
- jfim 3y ago> I've never done geographic sharding but it seems kind of hard. How do you pick shard boundaries? How do you deal with entities who are near the boundaries and whose current operational data therefore spans >1 shards? (Imagine somebody at near the geographic intersection of like, five shards looking for pizza in a 10 miles radius or w/e) You could do it by market (eg. SFBA, Los Angeles, San Diego) or by state.
- jordanthoms 3y agoThey would have to have many shards per city to keep up with the level of write traffic though. And what happens when a user from SFBA goes down to LA?
- michaelt 3y agoWould they? I mean, I've seen conventional SQL databases handle ten million orders per hour on a single host. I find it hard to believe DoorDash is processing more than ten million orders per hour, even in a large city. I suppose they might exceed what a single host can handle if they're, I don't know, recording every driver's location once per second?
- JohnBooty 3y ago
- sciurus 3y agoThe article says they have 300+ clusters, not one monster one.
- jordanthoms 3y agoSharding is anything but simple. A single shard per region wouldn't have enough write capacity so they'd have to be managing likely 100+ shards in each region - you'd have to build a lot of infrastructure to automate setting those up, rebalancing traffic to avoid hot spots and underutilized shards, in sync with schema migrations etc. Even after that, now your applications using the DB have to be aware of the sharding - interactions between users who are housed on different shards etc could require a lot of work at the application layer. If your customers can be easily be split into tenants which never interact with each other this isn't so bad but for a consumer app like DoorDash there isn't clear tenant boundaries. We looked at all this for Kami and realised that it would be much easier for us to move from PostgreSQL to CockroachDB (we had exceeded the write capacity of a single PostgreSQL primary) than to shard Postgres, and it'd make future development much faster. We could have made sharding work if we had to... but it's not 2013 any more and we have distributed SQL databases, why not use them?
- cellularmitosis 3y agoThat’s surprising — the education market seems like an even better fit for sharding: students and teachers generally stay within the context of a single school.
- jordanthoms 3y agoIt does seem like there would be a clean boundary between each school district, but actually there's plenty of sharing and collaboration on Kami that happens with users between districts, teachers and students move schools, parents can have children in different school districts, etc. Even a single classroom assignment can cross those, e.g. when someone external comes in to do a training session. We could have modified our application layer to handle those cases, but it's a lot of extra complexity and room for error, and we'd have had to consider and solve for all of these cross-tenant situations as we add new functionality, so I was really keen to avoid that. Also, there are some really big districts - NYCDOE has >1.1 million students and 1,800 schools. Even with them on a dedicated shard, it's quite possible that it'd get overloaded and we'd be spending more dev effort figuring out how to safely split them onto multiple shards. When we looked at using distributed SQL database instead it was a clear win - from the application's perspective, it just looks like a really, really big PostgreSQL box, so we didn't need to change much. (the SQL support is very close to PG - The most annoying thing for us was the lack of trigram indexes, and Cockroach has now added those now). And in terms of the operational side, upgrading and maintaining CRDB has actually been easier than PG - version upgrades are easier to do without downtime, and schema migrations don't lock tables.
- mike_d 3y agoBecause CockroachDB is a vendor that abstracts away all the thinking parts of running a database cluster. They do regional sharding, clustering, consistency, etc. for you. They could have just as easily dropped in Oracle. You pay for expensive DB up front, and can hire cheaper junior DBAs and developers going forward.
- ravenstine 3y agoBut it's at scale!!!!!
- snihalani 3y agointeresting. curious if anyone has benchmarked it relative to other dbs. like: https://benchmark.clickhouse.com/ https://benchmark.clickhouse.com/
- jordanthoms 3y agoClickhouse is a totally different use case - Cockroach is OLTP, Clickhouse is OLAP. We use both Cockroach and Clickhouse at scale and they are both great but not competing products - Cockroach is great for the types of reads and writes you do when serving user requests, processing transactions etc, but isn't optimal for analytics queries where you are going do things like read and aggregate data on a 50TB table. Clickhouse eats those kinds of aggregate queries for breakfast, and is fast for some types of small read queries too, but it's not built to handle random writes or frequently updating rows of data.
- karmakaze 3y agoCrDB is not about many many low latency queries, like say MySQL. It's designed more for getting your workload processes down to making as few large queries as every one incurs quorum latencies. You don't want to prototype something in Rails and hope there's no hidden lazy queries happening along the way. It wouldn't be a good idea to take a large working PostgreSQL app and try to switch over to using CrDB. You'd spend all your time (unwittingly rewriting the entire app) speeding up and grouping a few queries at a time.
- namibj 3y agoYou can have small queries, they just have to be be sent before you block on the results from the first of each group.
- jordanthoms 3y agoWe moved took a large working PostgreSQL app and switched it over to CRDB and that doesn't match my experience. Our existing schemas and query patterns moved over nicely - latency for small indexed reads and writes did increase from ~1ms to ~3ms, but the max throughput now effectively unlimited since we could add capacity by adding new nodes into the cluster and letting CRDB automatically rebalance the workload. There was an increase in cost as it will need more cores, disk etc compared to a single-primary PostgreSQL, but that makes sense when you consider that every bit of data is getting stored on 5 different nodes and there are overheads to maintain the consistency. For the highest throughput endpoints we did make some changes to be more optimal on CRDB so we could run a smaller cluster, but it didn't require anything close to a rewrite.
- cebert 3y agoThis reads like a long form advertisement.
- al_borland 3y agoCase studies hosted on a company's own website generally are. It's kinds of an, "it worked for them, so it will work for you," thing.
- candiddevmike 3y ago"Art of the possible" (YMMV)
- gizmo 3y agoDoorDash has about 35 million users, and there is zero interaction between users. The median user uses doordash maybe once a week. So 5 million sessions a day, all happening in the same 3 hour window. That's 2 million sessions per hour at peak times. How does DoorDash get to 1.2 million queries per second. 1.2mqps * 10000 seconds in 3 hours = 12 billion queries to process 5 million orders? That's wild. Is it all analytics? This is highly suspect. 35m users isn't nothing, but it isn't exactly Facebook scale either.
- BrentOzar 3y agoI’m not excusing the wild number, but just tossing out some additional load: * Drivers checking in for work, especially if the apps poll automatically * Drivers phoning home with live location updates * Restaurants sending automated updates on order status * Push notifications to users with status changes on their orders * Users with multiple devices (like I have at least 5 devices with the UberEats app)
- deleted 3y ago[deleted]
- jakjak123 3y agoYes, our server had 120k queries/sec, but 80% of that traffic was driver heartbeats or connection verification. We halved it by disabling the connection verification query.
- javawizard 3y agoHold up, what do you mean by "our server"? Do you work for DoorDash?
- jakjak123 3y ago[dead]
- cdchn 3y ago
- joshstrange 3y agoDo you know what DoorDash doesn’t manage? A staging/test environment. All testing for API integrations is done in prod on shared account. The docs and the API endpoints themselves leave a lot to be desired as well.
- jvans 3y agoI've been advocating for this approach for a long time. At some level of size it is so brutally difficult to maintain an environment that mirrors production that the effort isn't worth it. With enough tooling in place you can mitigate the risk to customers significantly
- stingraycharles 3y agoSo I assume they use feature flags instead, and staggered rollout of new features? As that’s a common alternative to heavy up-front testing.
- xyst 3y ago> About 1.2 million queries per second at daily peak hours. > About 2,300 total nodes spread across 300+ clusters. > About 1.9 petabytes of data on disk. > Close to 900 changefeeds. > Largest cluster is currently 280 TB in size (but has peaked above 600 TB), with a single table that is 122 TB. all of this yet my food still arrives cold af kidding aside, I wonder if DD has the same problems as Uber or Lyft except with food delivery. Each new "change feed" is a specific region, county/municipality, or city. Federal, state, and local laws all handled delicately.
- orangechairs 3y agoDoorDash's engineering blog has a much more indepth look at their architecture: https://doordash.engineering/2023/02/07/how-we-scaled-new-verticals-fulfillment-backend-with-cockroachdb/ https://doordash.engineering/2023/02/07/how-we-scaled-new-ve...
- beembeem 3y ago> my food still arrives cold af Ha. The first thing I noticed and you almost got to it in your summary: at 1.2MM/2300 = 520 qps per node, this isn't a wild setup. I'm wrapping my head around how they're generating that amount of load. Seems like an easy task for any database to handle.
- sean0- 3y agoHa! Amazing (didn’t know this was being written or put up). This is a summary of a recent conference talk: https://youtu.be/jCjrfpF64Kc?si=Gf-gp_ixX2V6Qz8V https://youtu.be/jCjrfpF64Kc?si=Gf-gp_ixX2V6Qz8V This was my team. We did and lived this. AMA.
- sverhagen 3y agoWell, it looks like the sibling threads are very interested to know where the need comes from to even _have_ 1.2 million queries per second. How does that break down? How much is that just core functionality versus analytics and tracking?
- sean0- 3y agoThat's core functionality, not analytics. Nearly all of that is browsing, ordering, location updates, etc. The dismissive comments are amusing and show a lack of understanding of how the business works and, subsequently, the technology required to power the end-to-end flow for users.
- sverhagen 3y agoOf course we don't know your business; we do back-of-envelope math based on our own experiences and of course a lot of assumptions about yours. Those numbers are just... impressive... or unbelievable... depending on where you're coming from.