Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,9 @@ Airflow supports multiple types of Dag Bundles, each catering to specific use ca
These bundles integrate with Git repositories, allowing Airflow to fetch Dags directly from a repository. The `GitDagBundle` does support versioning.

**airflow.providers.amazon.aws.bundles.s3.S3DagBundle**
These bundles reference an S3 bucket containing Dag files. They do not support versioning of the bundle, meaning tasks always run using the latest code.
These bundles reference an S3 bucket containing Dag files. They use the latest code by default. On Airflow 3.4
and later, an optional publisher-managed, content-addressed manifest protocol supports versioned Dag runs and
atomic deployments.

**airflow.providers.google.cloud.bundles.gcs.GCSDagBundle**
These bundles reference a GCS bucket containing Dag files. They do not support versioning of the bundle, meaning tasks always run using the latest code.
Expand Down Expand Up @@ -140,7 +142,9 @@ For an S3 Dag bundle, the required kwarg is ``bucket_name``. You can optionally

.. note::

``S3DagBundle`` does not support versioning. Tasks always run against the latest code in the bucket.
``S3DagBundle`` uses the latest code unless ``manifest_key`` is configured. Manifest mode requires Airflow 3.4
or later and an S3 Versioning-enabled bucket. Each Dag run then records a content-addressed release and can
retrieve the exact object versions used when the run was created.

See :doc:`apache-airflow-providers-amazon:bundles/index` for the full list of kwargs and more examples.

Expand Down
82 changes: 40 additions & 42 deletions airflow-core/src/airflow/dag_processing/bundles/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
import tempfile
import warnings
from abc import ABC, abstractmethod
from contextlib import contextmanager
from contextlib import contextmanager, suppress
from dataclasses import dataclass, field
from datetime import timedelta
from fcntl import LOCK_SH, LOCK_UN, flock
Expand All @@ -33,8 +33,8 @@
from typing import TYPE_CHECKING, Any

import pendulum
from pendulum.parsing import ParserError

from airflow._shared.timezones import timezone
from airflow.configuration import conf

if TYPE_CHECKING:
Expand Down Expand Up @@ -96,13 +96,6 @@ class BundleUsageTrackingManager:
:meta private:
"""

def _parse_dt(self, val) -> DateTime | None:
try:
dt = pendulum.parse(val)
return dt if isinstance(dt, pendulum.DateTime) else None
except ParserError:
return None

@staticmethod
def _filter_for_min_versions(val: list[TrackedBundleVersionInfo]) -> list[TrackedBundleVersionInfo]:
min_versions_to_keep = conf.getint(
Expand Down Expand Up @@ -134,15 +127,10 @@ def _find_all_tracking_files(self, bundle_name) -> list[TrackedBundleVersionInfo
for file in tracking_dir.iterdir():
log.debug("found bundle tracking file, path=%s", file)
version = file.name
dt_str = file.read_text()
dt = self._parse_dt(val=dt_str)
if not dt:
log.error(
"could not parse val as datetime bundle_name=%s val=%s version=%s",
bundle_name,
dt_str,
version,
)
try:
dt = timezone.from_timestamp(file.stat().st_mtime)
except FileNotFoundError:
# A concurrent cleanup may remove a tracking file after iterdir() observed it.
continue
found.append(TrackedBundleVersionInfo(lock_file_path=file, version=version, dt=dt))
return found
Expand All @@ -169,9 +157,10 @@ def log_info(msg):
with open(info.lock_file_path, "a") as f:
flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB) # exclusive lock, do not wait
# remove the actual bundle copy
shutil.rmtree(bundle_version_path)
with suppress(FileNotFoundError):
shutil.rmtree(bundle_version_path)
# remove the lock file
os.remove(info.lock_file_path)
info.lock_file_path.unlink(missing_ok=True)
except BlockingIOError:
log_info("could not obtain lock. stale bundle will not be removed.")
return
Expand Down Expand Up @@ -285,8 +274,9 @@ class BaseDagBundle(ABC):
that bundle version. This also means, that on a single worker, it's possible that multiple versions of the same
bundle are used at the same time.

In contrast, the DAG processor uses a bundle to keep the DAGs from that bundle up to date. There will not be
multiple versions of the same bundle in use at the same time. The DAG processor will always use the latest version.
In contrast, the DAG processor uses a bundle to keep the DAGs from that bundle up to date. It discovers the
latest version while allowing in-flight and version-pinned callback work to finish against its original
generation.

:param name: String identifier for the DAG bundle
:param refresh_interval: How often the bundle should be refreshed from the source in seconds
Expand All @@ -297,6 +287,8 @@ class BaseDagBundle(ABC):
"""

supports_versioning: bool = False
refreshes_to_versioned_paths: bool = False
"""Whether refreshing publishes a new immutable path while older published paths remain usable."""

_locked: bool = False

Expand Down Expand Up @@ -452,6 +444,9 @@ class BundleVersionLock:
"""
Lock version of bundle when in use to prevent deletion.

Acquisition failures propagate because running without the lease could allow cleanup to remove code in use.
Release failures are logged so they do not mask the outcome of the protected operation.

:meta private:
"""

Expand All @@ -476,28 +471,35 @@ def _log_exc(self, msg):
self.lock_file_path,
)

def _update_version_file(self):
"""Create a version file containing last-used timestamp."""
if TYPE_CHECKING:
assert self.lock_file_path
self.lock_file_path.parent.mkdir(parents=True, exist_ok=True)

with tempfile.TemporaryDirectory() as td:
temp_file = Path(td, self.lock_file_path)
now = pendulum.now(tz=pendulum.UTC)
temp_file.write_text(now.isoformat())
os.replace(temp_file, self.lock_file_path)

def acquire(self):
if not self.version:
return
if self.lock_file:
return
self._update_version_file()
if TYPE_CHECKING:
assert self.lock_file_path
self.lock_file = open(self.lock_file_path)
flock(self.lock_file, LOCK_SH)
self.lock_file_path.parent.mkdir(parents=True, exist_ok=True)
while True:
lock_file = open(self.lock_file_path, "a+")
flock(lock_file, LOCK_SH)
try:
path_stat = self.lock_file_path.stat()
except FileNotFoundError:
flock(lock_file, LOCK_UN)
lock_file.close()
continue
descriptor_stat = os.fstat(lock_file.fileno())
if (descriptor_stat.st_dev, descriptor_stat.st_ino) != (
path_stat.st_dev,
path_stat.st_ino,
):
flock(lock_file, LOCK_UN)
lock_file.close()
continue
now = pendulum.now(tz=pendulum.UTC).timestamp()
os.utime(self.lock_file_path, (now, now))
self.lock_file = lock_file
return

def release(self):
if self.lock_file:
Expand All @@ -506,11 +508,7 @@ def release(self):
self.lock_file = None

def __enter__(self) -> Self:
# wrapping in try except here is just extra cautious since this is in task execution path
try:
self.acquire()
except Exception:
self._log_exc("error when attempting to acquire lock")
self.acquire()
return self

def __exit__(self, exc_type, exc_val, exc_tb):
Expand Down
Loading