andygrove opened a new issue, #2500: URL: https://github.com/apache/datafusion-ballista/issues/2500
## Summary A few weeks ago Comet moved its expensive CI suites behind a GitHub merge queue (apache/datafusion-comet#5843, followed by apache/datafusion-comet#5939 and apache/datafusion-comet#5963). I think it has worked out really well. PRs get a light CI run while they're being iterated on, and the full suite still runs once, against the actual merge result, before anything lands on `main`. The background discussion is in apache/datafusion-comet#5830. I'd like to propose we do the same here. Our CI is much smaller than Comet's, but the same pattern shows up in our numbers. Most of the minutes we spend on PRs go to end-to-end jobs that very rarely catch anything the core Rust jobs don't, and those jobs run a second time on `main` after every merge. The ASF's GitHub Actions limits apply per project, so this all comes out of the same allocation as DataFusion core, Comet and the other subprojects. ## Where our CI minutes go today Over the last two weeks (2026-09-13 to 2026-09-27): - PR runs used about 17,000 runner-minutes across 54 PRs and 130 pushed commits. A typical commit that touches Rust code runs 24 jobs. - Post-merge runs on `main` used another 5,600 runner-minutes for 34 merges. - Ten of those 24 jobs account for 68% of all PR runner-minutes (11,600 of 17,000). To see how much those ten jobs tell us on a PR, I went through six weeks of PR runs (352 commits across 142 PRs since 2026-08-16). For each job I counted the commits where it failed while the core jobs (`test linux crates`, `test linux ballista`, `Clippy`, `cargo check`, `check linux workspace` and `cargo doc`) all passed. That's the only case where running it on the PR told us something new. | Job | Runner | PR runner-minutes (2 weeks) | Failing commits (6 weeks) | Not caught by the core jobs | | --- | --- | ---: | ---: | ---: | | TPC-H SF10 (3 legs) | Linux | 3,628 | 6 | 0 | | TPC-DS SF1 (2 legs) | Linux | 2,061 | 7 | 1 | | `macos test` | macOS | 1,698 | 25 | 0 | | Docker image build | Linux | 1,282 | 5 | 0 | | h2o window SMALL | Linux | 1,108 | 7 | 0 | | `verify benchmark results (amd64)` | Linux | 1,036 | 6 | 0 | | k8s chaos scenarios (kind) | Linux | 812 | 12 | 5 | A few observations: - The TPC-DS catch was a real bug in #2364, where range shuffle reads failed on invalid seeks. That's exactly what these suites are for. Under this proposal it would still have been caught before merge, just in the queue rather than on the PR. - Four of the five chaos catches were on #2425, which is changing the chaos harness itself. So a job should probably still run on the PR when the PR touches that job's own harness. - On `main`, the post-merge runs of TPC-H, TPC-DS, h2o and the Docker build passed all 356 times in those six weeks. At that point they're pure cost. There's a correctness argument too. On 2026-09-09, #2016 merged from a branch that had last synced with `main` on 2026-08-18. In between, #2354 had removed `SchedulerTest::await_completion`, and #2016 added a new call to it. Neither PR was wrong on its own and git merged them cleanly, but the `ballista-scheduler` tests didn't compile on `main` for about 11 hours until #2439 fixed it. A merge queue builds the merge result against the current `main`, so it would have sent #2016 back instead. ## Proposal Split CI into two tiers, following Comet's design. - **PR tier.** Runs on every push to a PR, for fast feedback while a change is being worked on. - **Queue tier.** Runs in the merge queue once a committer queues an approved PR. It's the PR tier plus the expensive suites, run against the merge result, and it's the gate for landing on `main`. It runs the PR tier again because the merged tree isn't the same as the PR head. The #2016 break above was a compile error that the PR tier would have caught in the queue. **Stays on PRs** (and runs again in the queue): - Lints: RAT, prettier, rustfmt, `Cargo.toml` formatting, `cargo-machete`, the vendored proto and config docs sync checks, and `Clippy` - Builds: `cargo check`, `check linux workspace`, `cargo doc` and the Web TUI WASM build - Tests: `test linux crates` and `test linux ballista` - The path-filtered checks, as today: MSRV and `python/Cargo.lock` checks, CodeQL, the ASF allowlist check, and Python ruff plus tests on one Python version **Moves to the merge queue:** - `macos test` - TPC-H SF10 (3 legs), TPC-DS SF1 (2 legs) and h2o window SMALL - `verify benchmark results (amd64)` - The Docker image build - The k8s chaos scenarios, which also keep their nightly run - Python tests on the other supported Python versions **Push to `main`:** the docs and Web TUI deploys, CodeQL and the allowlist check as today, and a cache refresh when dependencies change. Nothing the queue already tested on the same tree runs again. A queue-tier job can still run on a PR in two ways: - **When the PR changes the job's own harness.** For example its workflow file, `ci/scripts/ballista_cluster.sh` or `benchmarks/**` for the query suites, `chaos-testing/**` and the chaos Dockerfile for the chaos job, or `dev/docker/**` for the Docker build. The PR tier can't stand in for the job in that case. - **With a label.** Something like `run-query-suites` (TPC-H, TPC-DS, h2o and the benchmark verification), `run-macos-tests` and `run-full-ci` should cover most needs. If you're changing the planner, shuffle, AQE or task scheduling, adding `run-query-suites` early is a good idea so you don't find out in the queue. `cargo check` and `Clippy` already run on macOS in the PR tier, so a change that doesn't compile on macOS still fails on the PR. Only macOS-specific test failures would wait for the queue, and there were none in the six weeks above. ## Expected savings Using the same two weeks as a baseline: | | Today | With the queue (estimate) | | --- | ---: | ---: | | PR runs | 17,000 | about 5,400, plus label runs | | Merge queue | none | about 5,600 (one full run per merge) | | Post-merge runs on `main` | 5,600 | deploys and cache refreshes only | So PR minutes drop by well over half, and total CI minutes should come down by something like 40 to 50%, depending on how often labels get used. That's before counting any gains from fixing the cache (below). A Rust-touching commit goes from 24 jobs to 14, and PR feedback gets faster too, because the jobs moving out are the slowest ones on a PR today. It would also free up budget for coverage we've held back on because of cost, such as the partition-count sweep in #2201. ## How it would work This would mostly be a port of what Comet built, and Comet has written up the design and the lessons in its [CI contributor guide](https://datafusion.apache.org/comet/contributor-guide/ci.html) and [workflows README](https://github.com/apache/datafusion-comet/blob/main/.github/workflows/README.md). The main pieces: 1. **An umbrella workflow with a single `Required Checks` job.** A merge queue only waits on required status checks, and today `main` has none, so our CI is advisory. We can't just require the existing checks either. Our workflows use workflow-level `paths:` filters, and a workflow that doesn't trigger never reports its check, so a required check on it would sit at "Expected" forever. The `merge_group` event doesn't support `paths:` filters at all. Comet's answer is one `ci.yml` with a `changes` job that works out which suites a change needs (from the PR diff, or from the merge group's base and head), the existing jobs called as reusable workflows, and one flat `Required Checks` job that `needs:` all of them. That's the only required check. Skipped jobs count as a pass, and failed or cancelled jobs turn it red. apache/datafusion ran into the hanging-check problem when it required its checks directly (apache/datafusion#17538, reverted in apache/datafusion#17629). 2. **Merge queue config in `.asf.yaml`.** This is self-serve now. `.asf.yaml` rulesets accept a raw GitHub Rulesets payload, so we can copy Comet's `Merge Queue` ruleset: squash only, `ALLGREEN`, and `apache/root` as a bypass actor so a stuck queue can always be recovered. A queue run for us should take around 30 minutes, compared with about 2.5 hours for Comet, and we merge about 3 PRs a day (10 on the busiest day in the last six weeks). Comet's build concurrency of 2 is more than enough. 3. **A two-step rollout.** First merge the workflow changes and watch `Required Checks` report correctly on real PRs. Only then change `.asf.yaml` to require it and turn on the queue. Getting the order wrong is the one mistake that needs an INFRA ticket to undo, because a required check that never reports also blocks the PR that would fix it. 4. **Fixing the Actions cache at the same time.** Our cache is thrashing today. Every Rust job writes its own `rust-cache` entry, so a single PR run writes 18 entries totalling about 9.4 GB, which is nearly all of the 10 GB repository limit. When I checked today, 19 of the 20 cache entries belonged to PRs and `main` had just one. PR runs can only restore caches from `main` or from their own branch, so most runs start cold. Merge queue branches would make this worse, because no later run can reuse their caches. Comet ended up saving large caches only on push to `main` (apache/datafusion-comet#5973) and cutting the push run down to the cache-writing jobs (apache/datafusion-comet#5930). We should start there, and share cache keys between jobs that build the same thing. The TPC-H, TPC-DS and h2o jobs all run the same `cargo build --profile tpch-ci`, for example. 5. **Release branches.** Merge queue rules don't accept wildcards, so `branch-53` and `branch-54` would keep plain branch protection. PRs that target a release branch should run both tiers, since there's no queue behind them (apache/datafusion-comet#6218). 6. **Label runs in their own workflow.** Comet found that a label-triggered run inside the main workflow could leave `Required Checks` stuck at "Expected" (apache/datafusion-comet#6159). Running label-triggered jobs from a separate workflow avoids that. ## Trade-offs - **CI becomes a hard gate.** Today a committer can merge over a red check. With the queue, a failure removes the PR from the queue, and a flaky job blocks everyone's merges rather than one PR. We already have at least one flaky test. `executor_killed_after_shuffle_write_is_recovered::case_1_aqe_off` failed on `main` after #2408, which only changed CI scripts. We'd need to fix or quarantine flaky tests promptly. The good news is that the jobs moving to the queue have been very stable on `main`. - **Merging takes a little longer.** A committer queues the PR with "Merge when ready" and it lands around 30 minutes later if the queue run is green. In exchange, everything that lands has been tested against the current `main`. - **Some failures show up later.** A TPC-DS or macOS-only failure would appear in the queue rather than on the PR. Going by the last six weeks that should be rare, and the labels let authors opt in early for risky changes. ## Questions - Does the split look right? In particular, should `verify benchmark results (amd64)` or one TPC-H leg stay in the PR tier as a cheap end-to-end canary? - Should `macos test` gate every merge, or would a nightly run be enough? - Is anyone opposed to making CI a required check on `main`? That's the biggest change for committers. If there's support, I'm happy to do the implementation, starting with the umbrella workflow and the `Required Checks` job. -- 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]
