Object-store DAG bundles (S3/GCS) stage per-object on every fresh worker — support staging from a single archive object
- Dominant language
- Python
- Stars
- 46.9k
- Forks
- 17.8k
- Avg merge
- 2d 9h
- Merged PRs (30d)
- 472
Description
### Description
All object-store DAG bundles refresh by listing the prefix and downloading objects **one at a time, serially** (`S3Hook.sync_to_local_dir()`; the GCS bundle ports the same loop, and the proposed Azure Blob bundle in #67016 follows the same template). For a DAG project of N small files that costs ~N GET (+N HEAD from `download_file`) requests and N sequential round-trips — on every component that stages the bundle. Under KubernetesExecutor every task instance is a fresh pod with an empty bundle dir, so the full cost lands on the critical path of **every task start**, and the existing ETag/mtime skip logic never helps (nothing is cached in a fresh pod). Even in the steady state, discovering that *nothing changed* still costs a full per-object walk.
**Proposal:** an opt-in archive mode for object-store bundles — e.g. an `archive_key` kwarg:
```json
{"name": "dags", "classpath": "...S3DagBundle",
"kwargs": {"bucket_name": "my-dags", "archive_key": "bundle/dags.tar.gz", "prefix": "dags/"}}
```
`refresh()` then does: HEAD the archive and compare its ETag against the last staged one (unchanged ⇒ done, **1 request total**) → one GET → verify + unpack (`tarfile` with `filter="data"`) → atomic rename swap — and **falls back to the stock per-object sync on any archive problem**, so the worst case is exactly today's behavior and `prefix` keeps its meaning. Publishing the archive is the deployer's job (one `tar | aws s3 cp` in the same CI that syncs the prefix). The staging logic is provider-agnostic, so it could live in a shared helper that the S3/GCS/Azure bundles all use rather than being copy-pasted per provider like the current sync loop.
Beyond the request count, the archive model is structurally simpler in two ways:
- **Updates and deletions are implicit.** The archive is the complete desired state: extraction into a fresh directory plus a rename swap means a file absent from the archive simply doesn't exist in the new tree. No stale-file reconciliation logic at all (the reconciliation in the per-object sync is where #62622 lived), and the DAG processor never observes a half-updated directory — the swap is all-or-nothing, whereas the per-object sync mutates the live directory while it may be read.
- **The archive's ETag is a true bundle version.** The stock bundles are `supports_versioning = False` for a structural reason: a prefix of independently-mutable objects has no atomic snapshot to name. A single archive does — its ETag (or a content-addressed key such as `dag-bundle-.tar.gz`) identifies an exact, immutable bundle state. That would give object-store bundles a path to real versioning (e.g. a task staging the same bundle version its DAG run started with), which the current per-object design cannot express.
A lighter complementary option is a concurrency knob (`max_workers`) on `sync_to_local_dir` for users who can't add an archive-publishing step — that improves wall-clock ~10× but still issues N requests and keeps the reconciliation logic.
**Fully local, reproducible demo** (astro + MinIO, no AWS account; includes a working `S3DagBundle` subclass implementing the above): https://github.com/aoelvp94/airflow-s3-dag-bundle-staging-demo — `just e2e` reproduces, for the identical 400-file / ~3 MiB file set:
| staging strategy | wall-clock (median) | requests served | files staged |
|---|---:|---:|---:|
| stock per-object sync | 15.8 s | 802 (400 GET + 400 HEAD + listings) | 400 |
| single archive | 1.5 s | 5 | 400 |
That's against a *loopback* MinIO (~ms round-trips). Real S3 adds ~25–40 ms per request, which is how a 400-file bundle reaches 20+ seconds per pod. The cost scales linearly with file count and per-request latency, not with bytes.
### Use case/motivation
DAG projects are many tiny files, not few big ones — moving a few MB is a request-count problem, not a bandwidth one (multipart parallelism never engages at these sizes). What I'm trying to achieve: a fresh worker should be able to stage a DAG bundle in ~1 request and ~1 second, the way every other system ships code to ephemeral workers (Spark archives, Lambda zips, image layers — Airflow's own `GitDagBundle` gets this implicitly since a fetch is one packed transfer, and a commit SHA plays the same "atomic named snapshot" role the archive ETag would play here).
Concretely, for a typical mid-size deployment — 400-file bundle, `KubernetesExecutor`, 10k task instances/day — per-object staging means ~8M billable object-store requests/day and ~50 pod-hours/day of pure waiting, and ~20 s of latency in front of every task, which erodes the headroom of tight-cadence DAGs (a 5-min-schedule DAG loses a third of its budget before user code runs) and silently inflates every `dagrun_timeout`. With three providers now sharing this staging pattern, fixing it once at the shared level benefits all of them.
### Related issues
None tracking staging performance. Adjacent: [apache/airflow#62622](https://github.com/apache/airflow/issues/62622) (S3 stale-delete fix — a bug class the archive model eliminates by construction), [apache/airflow#66987](https://github.com/apache/airflow/issues/66987) / [apache/airflow#67016](https://github.com/apache/airflow/pull/67016) (Azure Blob bundle in review, same sync pattern), and the GCS bundle which ports the same loop.
### Are you willing to submit a PR?
- [x] Yes I am willing to submit a PR!
### Code of Conduct
- [x] I agree to follow this project's [Code of Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
Contributor guide
Research direction
Start with S3Hook.sync_to_local_dir() and the S3DagBundle refresh path, then compare the corresponding GCS loop and the proposed Azure Blob bundle in #67016. Done means an opt-in archive path can validate and stage a complete bundle atomically, skip unchanged archives, fall back to per-object sync on archive errors, and preserve the existing prefix behavior.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- aws, azure, gcp, python
- Domain
- backend, cloud
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Active
- Clarity
- Mostly clear
- Newbie friendliness
- 45/100