Skip to content

Object-store DAG bundles (S3/GCS) stage per-object on every fresh worker — support staging from a single archive object #73010

Description

@aoelvp94

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:

{"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 S3DagBundle does not delete stale dag recursively #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-<sha>.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-demojust 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 (S3 stale-delete fix — a bug class the archive model eliminates by construction), apache/airflow#66987 / apache/airflow#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?

  • Yes I am willing to submit a PR!

Code of Conduct

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    kind:featureFeature Requestsneeds-triagelabel for new issues that we didn't triage yet

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions