andygrove opened a new issue, #2525:
URL: https://github.com/apache/datafusion-ballista/issues/2525
**Is your feature request related to a problem or challenge? Please describe
what you are trying to do.**
Operators want to upgrade a running Ballista deployment without an outage.
Today we only partly support that.
What works on `main`:
- Executors can roll without losing queries. On SIGTERM an executor stops
taking new tasks, tells the scheduler it's stopping, and finishes its running
tasks. The scheduler then re-runs any work whose shuffle output went away with
it, and if the last executor disappears it waits a grace period for new ones
before failing jobs (#2029).
- `BALLISTA_PROTOCOL_VERSION` (#2071, #2088) keeps a new scheduler from
handing work to old executors, and in the documented Kubernetes setup the
client-facing Service only routes to a scheduler that passes `/readyz`. The
steps are in the Kubernetes guide under "Rolling both Deployments together".
What doesn't:
- **Replacing the scheduler fails every query in flight.** Scheduler state
is in memory only (`ClusterStorage::Memory` is the only backend), so the new
scheduler starts empty. The Kubernetes guide says so: "a short window with no
scheduler, so queries in flight fail". Running several scheduler replicas keeps
new queries flowing, but queries on the replaced scheduler still fail, because
replicas don't share state.
- **Clients aren't part of the upgrade.** Within a major version, clients
and schedulers work together in practice. Across a major version nothing
guarantees that, because DataFusion's plan encoding and Ballista's
client-facing protocol both change between majors (#2367, #2370, #2376). Today
a mismatched client is accepted and can get wrong results. #2514 (draft) would
make that an error, which means clients and the cluster have to move to a new
major together, and any client on the other side gets errors until it catches
up.
This is getting more pressing. Users run clusters around the clock (#2071),
and now that the Python client is released separately from the Rust crates
(#2512), Python users will often be a major version behind the cluster.
**Describe the solution you'd like**
This needs a design, and I don't think we know the answer yet. Some
questions it should cover:
- **Scheduler handover.** Can an outgoing scheduler drain, finishing the
jobs it has while a new scheduler takes new ones? Or do we need shared job
state so another scheduler can pick up a running job? #2349 is adding the
public APIs for scheduler failover, which could be a building block. But it
expects an incompatible execution graph version to fail with a clear error, so
handing over between versions isn't in scope there.
- **Clients across a major version.** Could a scheduler accept clients from
the previous major for a while? That only works if it can decode their plans
with their original meaning, which DataFusion doesn't guarantee. The
alternatives are running old and new clusters side by side and moving clients
over, or some form of version negotiation.
- **What we promise.** Which surfaces stay compatible across which versions,
and for how long. #2261 covers the executor protocol, the REST API, and the
event log. Clients of the scheduler's gRPC API would be a new surface.
**Describe alternatives you've considered**
- Upgrade clients and the cluster together and accept an outage window. This
is what happens today, and #2514 would make it explicit for major upgrades.
- Blue/green at the deployment level: stand up a second cluster on the new
version, move clients over, then retire the old one. This needs no Ballista
changes, but it doubles the cluster for the length of the upgrade and leaves
the client cutover to each user.
**Additional context**
The compatibility table under "Rolling both Deployments together" describes
the first release with the protocol check ("old scheduler pre-dates the
check"). From the next `BALLISTA_PROTOCOL_VERSION` bump on, an old scheduler
rejects new executors too, so that table needs updating after 55.0.0.
Related: #2029, #2071, #2088, #2261, #2349, #2512, #2514.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]