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]

Reply via email to