9 ms·
Carving the scheduler out of our orchestrator
- plaidfuji 4y ago> With strict bin packing, we end up with Katamari Damacy scheduling, where a couple overworked servers in our fleet suck up all the random jobs they come into contact with.
- anyas 4y agoWhat is the rationale for having the warm spares? If you're already able to spin up new VMs within an HTTP request, that seems pretty fast. Is the extra complexity just to make that path even faster?
- tptacek 4y agoThe "warm spare" isn't running, but it's pre-loaded onto a particular worker. The slowest part of starting an app from scratch is checkout out the container; "warming up" the worker with the "spare" machine eliminates that.
- nitwit005 4y agoI suppose I'm more curious how you tested this, than the basic idea. For example, if Nomad was doing unwanted movement of apps between machines, how do you prove the replacement code won't?
- dijksterhuis 4y ago$CURRENT_JOB is breaking apart and refactoring an nightmarish "orchestration application" where about 70% of the codebase solely in django. will be referencing and sharing this fairly frequently I reckon so massive thank you @tptacek
- yencabulator 4y agoThis made me follow the link to https://fly.io/blog/ipv6-wireguard-peering/ https://fly.io/blog/ipv6-wireguard-peering/ and I think you have a copy-pasto: > Our WireGuard mesh sees IPv6 addresses that look like fdaa:host:host::/48 but the rest of our system sees fdaa:net:net::48. That was probably mean to be net:host -> host:net
- tptacek 4y agoI'm worried that I'm meaner about K8s in this than I mean to be, which would be a problem not least because I don't have enough K8s experience to justify meanness. I'm really more just surprised at how path-dependent the industry is; even systems that were consciously built not to echo Borg, even greenfield systems like Flynn that were reimaginactiments of orchestration, all seem to follow the same model of central, allocating schedulers based on distributed consensus about worker inventory.
- ilyt 4y ago> all seem to follow the same model of central, allocating schedulers based on distributed consensus about worker inventory. Because frankly, it's easy to code. You just look thru list of what you have, list of what you put on it, and assign based on this and that algorithm. It keeps workers simple - register to cluster, send status (or wait for healthckeck), wait for request to start the job, done, you got a worker. It keeps the scheduler to just be scheduler - take state, apply list of "changes" (jobs to start/stop), execute. It keeps client simple - submit to scheduler, observe progress. It is also easy to debug. You just have one place doing the scheduling (and not a bunch of workers playing the "stock market" of resources), running on one node without distributed mess to worry about. Any test is just "feed this state to code and see whether it reacted as we need". Any job that got "rejected" for scheduling have clear reason. It also allows bigger optimizations to happen from point of view of a cluster. Some designs are reinvented so many times because they are just good or easy to implement first approximation of solution.
- vidarh 4y agoBe mean about K8s, it can take it. I don't think any of the criticism of it is unreasonable. It's a huge, convoluted beast. To me, K8s makes most sense in terms of creating a swiss army knife that lots of people know. You can do far better than K8s for specialised cases like yours and/or with people who has the right skills. But for a lot of people with simpler needs it's better to just stick to what is easy to hire for than look for an optimal solution.
- intelVISA 4y ago
- schmichael 4y agoAs the Nomad Team Lead, this article is a gift - thank you Fly! - even if they're transitioning off of Nomad. Their description of Nomad is exactly what I would love people to hear, and their reasons for DIYing their own orchestration layer seem totally reasonable to me. Nomad has never wanted people to think we're The One True Way to run all workloads. I hope Nomad covers cases like scaling-from-zero better in the future, but to do that within the latency requirements of a single HTTP request is quite the feat of design and implementation. There's a lot of batching Nomad does for scale and throughput that conflict with the desire for minimal placement+startup latency, and it's yet to be seen whether "having our cake and eating it too" is physically possible, much less whether we can package it up in a way operators can understand what tradeoffs they're choosing. I've had the pleasure of chatting with mrkurt in the past, and I definitely intend to follow fly.io closely even if they're no longer a Nomad user! Thanks again for yet another fantastic post, and I wish fly.io all the best.
- tptacek 4y agoThis whole article started with me rewatching your Nomad deep dive video and then chasing papers. :)
- Dowwie 4y agoAre you referring to, "Nomad Under the Hood" [1]? Which papers did you find useful? [1] https://www.youtube.com/watch?v=m6DnmVqoXvw https://www.youtube.com/watch?v=m6DnmVqoXvw
- sitkack 4y agoI haven't watched the whole video so I don't know about papers referenced, but between chatgpt and papers mentioned in the source repo, here is a good starting list. Large-scale cluster management at Google with Borg https://research.google/pubs/pub43438/ https://research.google/pubs/pub43438/ Omega: flexible, scalable schedulers for large compute clusters https://research.google/pubs/pub41684/ https://research.google/pubs/pub41684/ The Chubby lock service for loosely-coupled distributed systems https://disco.ethz.ch/courses/hs08/seminar/papers/osdi06-google-chubby.pdf https://disco.ethz.ch/courses/hs08/seminar/papers/osdi06-goo... Apache Mesos https://scholar.google.com/scholar?hl=en&as_sdt=0%2C5&q=Apache+Mesos&btnG= https://scholar.google.com/scholar?hl=en&as_sdt=0%2C5&q=Apac... mentioned in the nomad repo and docs site Sparrow: Distributed, Low Latency Scheduling https://cs.stanford.edu/~matei/papers/2013/sosp_sparrow.pdf https://cs.stanford.edu/~matei/papers/2013/sosp_sparrow.pdf SWIM: Scalable Weakly-consistent Infection-style Process Group Membership Protocol https://www.cs.cornell.edu/projects/Quicksilver/public_pdfs/SWIM.pdf https://www.cs.cornell.edu/projects/Quicksilver/public_pdfs/... Raft: In search of an Understandable Consensus Algorithm https://raft.github.io/raft.pdf https://raft.github.io/raft.pdf ---- Finally, take a look at the papers referenced on semanticscholar https://www.semanticscholar.org/search?q=nomad%20orchestration%20scheduler&sort=relevance https://www.semanticscholar.org/search?q=nomad%20orchestrati...
- kalev 4y agoI’m annoyed by the way this is written. The topic is super interesting but the author tried to hard to be funny and being a non-native reader it’s difficult to determine if certain words are technical jargon or trying to be funny.
- Dowwie 4y agoI think it succeeded at being funny. tptacek is a Nümad lad.
- tptacek 4y agoI'm mostly just curious about which terms the reader thinks are made up.
- titanomachy 4y agoI particularly liked "Katamari Damacy scheduling", but I would guess that at least some of the cultural references (dilithium crystals, Marcellus Wallace) could be confusing to folks not steeped in a particularly nerdy subset of American-ish culture. I think if references and jokes are landing for (say) 80% of your target audience then you're doing pretty well. I liked the article. Interesting approach to scheduling as kind of a market exchange problem. It's cool to see what an infra service can look like when you ditch the constraints and path dependence that led to the "standard" cloud architecture and build something totally new. The user experience seems like multi-tenant Borg, but nimbler.
- romantomjak 4y ago
- coredog64 4y agoAlthough it’s currently archived, there’s another open source orchestrator that’s similar to Borg: Treadmill. AFS solves for Colossus, with packages being distributed into AFS. [0] https://github.com/morganstanley/treadmill https://github.com/morganstanley/treadmill
- korijn 4y agoI can't shake the sense that they seem to just prefer building something new (and blog about it) over figuring out how to configure a/the scheduler properly. Anyway, they did get it done and made it work, so whatever.
- mochomocha 4y agoI have my own fair share of complaints about k8s, but I can't say the author articulates clearly what is wrong with k8s scheduler exactly. IMO it's one of the better part of k8s. The core scheduler is pretty well written and extensible through scheduling plugins, to implement whatever policies you heart desires (which we extensively make use of at Netflix). The main issue I have with it is the lack of built-in observability, which makes it non-trivial to A/B test scheduling policies in large scale deployment setups because you want to be able to log the various subscores of your plugins. But it's so extensible through NodeAffinity and PodAffinity plugins that you can even delegate part (or all!) of the scheduling decisions outside of it if you want. Besides observability, one issue we've had to overcome with k8s scheduling is the inheritance of the Borg design decisions around pod shape immutability, which makes implementing things like oversubscription less easy in a "native" way.
- tptacek 4y agoNothing is wrong with the k8s scheduler! Really, nothing is wrong with k8s at all, beyond our more general problem of "our users want to run Linux apps, not k8s apps". K8s, Borg, Omega, Flynn, Nomad, and to some extent Mesos all share a common high-level scheduler architecture: a logically centralized, possibly distributed server process that functions like an allocator and is based on a consistent view of available cluster resources. It's a straightforward and logical way to design an orchestrator. It's probably the way you'd decide to do it by default. Why wouldn't you? It's an approach that works well in other domains. And: it works well for clusters too. The point of the post is that it's not the only way to design an orchestrator. You can effectively schedule without a centralized consistent allocator scheduler, and when you do that, you get some interesting UX implications. They're not necessarily good implications! If you're running a cluster for, like, Pixar, they're probably bad. You probably want something that works like Borg or Omega did. You have a (relatively) small number of (relatively) huge jobs, you want optimal placement†, and you probably want to minimize your hardware costs. We have the opposite constraints, so the complications of keeping a globally consistent real-time inventory of available resources and scheduling decisions don't pay their freight in benefits. That's just us, though. It's probably not anything resembling most k8s users. † In fact, going back even before Borg but especially once Borg came on the scene, mainstream schedulers have been making this distinction --- between service jobs and batch jobs, where batch jobs are less placement sensitive and more delay sensitive. So one way to think about the design approach we're taking is, what if the whole platform scheduler thought in terms of a batch-friendly notion of jobs, and then you built the service placement logic on top of it, rather than alongside it?
- filereaper 4y agoMesos always had these notions of two level scheduling that let you build your own orchestration. Aurora, Marathon, etc... would add the flavor of Orchestration that's needed. Mesos provided the resources requested. https://mesos.apache.org/documentation/latest/architecture/ https://mesos.apache.org/documentation/latest/architecture/
- GauntletWizard 4y agoThat's precisely the reason Mesos lost - It didn't come batteries included with a scheduler, and nobody ever packed Aurora+Nomad nicely.
- necubi 4y agoI remain sad that Mesos never really took off, and then k8s ate the market. It had a lot of really clever ideas and could do stuff well that k8s still can't (although it has been catching up, ten years later). In particular, the scheduler architecture allowed it to work well for both batch workloads and service-like workloads within a single resource pool.
- ilyt 4y agoYeah it was neat and oh so much less complex to setup than k8s from scratch. But k8s is not just job scheduler, it comes with (at hefty complexity cost) with all the piping and plumbing around it so it appealed to developers that could pop out a single .yaml that described their whole architecture.
- tptacek 4y agoBorg, Omega, K8s, and Nomad all have a separate scheduling pathway for batch jobs, don't they? My understanding is that this is the big complexifier for all distributed schedulers: that once you have two different service models for scheduling (placement sensitive delay insensitive, and placement insensitive delay sensitive), you now have two scheduler process, and you have to work out concurrency between their claims on resources. The Omega design in particular is, like, ground up meant to address this problem; the whole paper is basically "why the Mesos approach is suboptimal for this problem".
- jallmann 4y agoThis was a great article. While I was at Livepeer (distributed video transcoding on Ethereum [1]), we converged onto a very similar architecture, disaggregating scheduling into the client itself. The key piece is to have a registry [2] with a (somewhat) up-to-date view of worker resources. This could actually be a completely static list that gets refreshed once in a while. Whenever a client has a new job, they can look at their latest copy of the registry, select workers that seem suitable, submit jobs to workers directly, handle re-tries, etc. One neat thing about this "direct-to-worker" architecture is that it allows for backpressure from the workers themselves. Workers can respond to a shifting load profile almost instantaneously without having to centrally deallocate resources, or wait for healthchecks to pick up the latest state. Workers can tell incoming jobs, "hey sorry but I'm busy atm" and the client will try elsewhere. This also allows for rich worker selection strategies on the client itself; eg it can preemptively request the same job on multiple workers, keep the first one that is accepted, and cancel the rest, or favor workers that respond fastest, and so forth. [1] We were more of a "decentralized job queue" than "distributed VM scheduler" with the corresponding differences, eg shorter-lived jobs with fluctuating load profiles and our clients could be thicker. But many of the core ideas are shared, even our workers were called "orchestrators" which in turn could use similar ideas to manage jobs on GPUs attached to it... schedulers all the way down! [2] Here the registry seems to be constructed via the Corrosion gossip protocol; we used the blockchain with regular healthcheck probes.