This is an automated email from the ASF dual-hosted git repository.
vincbeck pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 92c919de773 Stop reporting components as degraded after an instance
restart (#73222)
92c919de773 is described below
commit 92c919de773a7c8b652fc30b10594a40780b4ddd
Author: Vincent <[email protected]>
AuthorDate: Fri Sep 18 12:03:31 2026 -0400
Stop reporting components as degraded after an instance restart (#73222)
The health endpoint's ``detailed_status`` compared live job rows against all
unfinished job rows, but a job row is only marked finished by a cooperative
shutdown. An instance lost to SIGKILL, an out-of-memory kill, or a node
eviction leaves its row unfinished forever, so restarting any scheduler,
triggerer, or Dag processor made the component look permanently ``degraded``
even though the deployment was fully healthy. Those same abandoned rows were
also listed in ``instances``, advertising hosts that no longer exist.
The question ``detailed_status`` should answer is whether every part of a
component's work has a live instance covering it, which needs the expected
set
of work to be declared rather than inferred from whatever rows happen to be
in
the table. Dag bundle configuration already declares it for the components
that
partition their work, and schedulers do not partition theirs at all.
---
.../logging-monitoring/check-health.rst | 44 ++-
.../src/airflow/api/common/airflow_health.py | 174 +++++++---
.../api_fastapi/core_api/datamodels/monitor.py | 17 +-
.../core_api/openapi/v2-rest-api-generated.yaml | 55 ++--
.../src/airflow/dag_processing/bundles/manager.py | 36 +-
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 88 ++---
.../airflow/ui/openapi-gen/requests/types.gen.ts | 27 +-
.../tests/unit/api/common/test_airflow_health.py | 361 +++++++++++++++------
.../bundles/test_dag_bundle_manager.py | 40 ++-
.../src/airflowctl/api/datamodels/generated.py | 50 ++-
10 files changed, 589 insertions(+), 303 deletions(-)
diff --git
a/airflow-core/docs/administration-and-deployment/logging-monitoring/check-health.rst
b/airflow-core/docs/administration-and-deployment/logging-monitoring/check-health.rst
index f42ea8702f4..5289e7f042c 100644
---
a/airflow-core/docs/administration-and-deployment/logging-monitoring/check-health.rst
+++
b/airflow-core/docs/administration-and-deployment/logging-monitoring/check-health.rst
@@ -51,7 +51,6 @@ including per-instance details when multiple schedulers,
triggerers, or Dag proc
"detailed_status": "healthy",
"instances": [
{
- "status": "healthy",
"hostname": "scheduler-1.example.com",
"latest_scheduler_heartbeat": "2018-12-26T17:15:11+00:00"
}
@@ -63,16 +62,9 @@ including per-instance details when multiple schedulers,
triggerers, or Dag proc
"detailed_status": "degraded",
"instances": [
{
- "status": "healthy",
"hostname": "triggerer-1.example.com",
"latest_triggerer_heartbeat": "2018-12-26T17:16:12+00:00",
"team_name": "team-a"
- },
- {
- "status": "unhealthy",
- "hostname": "triggerer-2.example.com",
- "latest_triggerer_heartbeat": "2018-12-26T17:10:00+00:00",
- "team_name": null
}
]
},
@@ -82,7 +74,6 @@ including per-instance details when multiple schedulers,
triggerers, or Dag proc
"detailed_status": "healthy",
"instances": [
{
- "status": "healthy",
"hostname": "dag-processor-1.example.com",
"latest_dag_processor_heartbeat": "2018-12-26T17:16:12+00:00",
"bundle_names": ["dags-team-a"]
@@ -100,11 +91,28 @@ including per-instance details when multiple schedulers,
triggerers, or Dag proc
* ``status`` (legacy aggregate): ``"healthy"`` if **any** running instance
is alive, otherwise ``"unhealthy"``
(including when no running jobs exist for that component).
- * ``detailed_status``: reflects the full set of running instances:
-
- * ``"healthy"`` — every running instance is alive
- * ``"degraded"`` — some instances are alive and some are not
- * ``"down"`` — no running instance is alive (including when no jobs exist)
+ * ``detailed_status``: whether every part of that component's work is being
covered by a live instance.
+ What counts as "every part" differs per component, because only some of
them divide their work up:
+
+ * **Dag processor** — the parts are the Dag bundles in ``[dag_processor]
dag_bundle_config_list``.
+ A processor started without ``--bundle-name`` covers every configured
bundle;
+ one started with it covers only the bundles it was given.
+ ``"healthy"`` when every configured bundle has a live processor,
``"degraded"`` when only some do,
+ ``"down"`` when none do.
+ * **Triggerer** — with ``[core] multi_team`` enabled, the parts are the
teams those bundles are scoped to
+ (plus the unscoped bundles), because a triggerer only picks up triggers
for its own team.
+ ``"healthy"`` when every team scope has a live triggerer, ``"degraded"``
when only some do,
+ ``"down"`` when none do. With multi-team disabled, no team filtering
applies, so any live triggerer
+ covers everything: ``"healthy"`` if one is alive, ``"down"`` if none is.
+ * **Scheduler** — schedulers are symmetric and share no partitioned work,
so there is nothing partial
+ to report: ``"healthy"`` if at least one is alive, ``"down"`` if none
is. ``"degraded"`` is never
+ returned for the scheduler. Use ``instances`` to see how many replicas
are up, and your orchestrator
+ or the ``scheduler_heartbeat`` metric to alert on reduced scheduling
throughput.
+
+ Because the expected set of work comes from configuration rather than from
the job table,
+ ``detailed_status`` is unaffected by how instances come and go. Restarting
an instance — including after
+ a ``SIGKILL``, an out-of-memory kill, or a node eviction, none of which
let Airflow mark the old job row
+ as finished — does not report the component as ``"degraded"``.
* ``latest_*_heartbeat``: the most recent heartbeat among running jobs of
that type (ordered by heartbeat descending),
or ``null`` when there are none.
@@ -113,9 +121,13 @@ including per-instance details when multiple schedulers,
triggerers, or Dag proc
``[triggerer] triggerer_health_check_threshold``,
``[dag_processor] health_check_threshold``).
- * ``instances``: one entry per **running** job of that type (``null`` when
there are none). Each entry includes:
+ * ``instances``: one entry per **live** instance of that type — ``null``
when none is live. Airflow cannot
+ mark a job as finished when its process is killed abruptly (``SIGKILL``,
an out-of-memory kill, a node
+ eviction), so the job row of such an instance is never closed; listing it
would show a host that is gone
+ and, under an orchestrator that assigns a fresh hostname on each restart,
will never come back.
+ ``status`` and ``latest_*_heartbeat`` are still derived from every
unfinished job, so how long a dead
+ component has been silent stays visible after it drops out of
``instances``. Each entry includes:
- * ``status``: ``"healthy"`` or ``"unhealthy"`` for that instance
* ``hostname``: host where the component is running
* the corresponding ``latest_*_heartbeat`` for that instance
* ``team_name`` (triggerer only): team the triggerer is scoped to, or
``null`` when unscoped
diff --git a/airflow-core/src/airflow/api/common/airflow_health.py
b/airflow-core/src/airflow/api/common/airflow_health.py
index e95a2966b3a..def85190e44 100644
--- a/airflow-core/src/airflow/api/common/airflow_health.py
+++ b/airflow-core/src/airflow/api/common/airflow_health.py
@@ -16,10 +16,14 @@
# under the License.
from __future__ import annotations
+import logging
+from enum import Enum
from typing import TYPE_CHECKING, Any
from sqlalchemy import select
+from airflow.configuration import conf
+from airflow.dag_processing.bundles.manager import
_get_configured_bundle_team_names
from airflow.jobs.dag_processor_job_runner import DagProcessorJobRunner
from airflow.jobs.job import Job
from airflow.jobs.scheduler_job_runner import SchedulerJobRunner
@@ -29,10 +33,22 @@ from airflow.utils.session import NEW_SESSION,
provide_session
if TYPE_CHECKING:
from sqlalchemy.orm import Session
-HEALTHY = "healthy"
-UNHEALTHY = "unhealthy"
-DEGRADED = "degraded"
-DOWN = "down"
+log = logging.getLogger(__name__)
+
+
+class HealthStatus(str, Enum):
+ """Aggregate health of a component: whether it has at least one live
instance."""
+
+ HEALTHY = "healthy"
+ UNHEALTHY = "unhealthy"
+
+
+class DetailedHealthStatus(str, Enum):
+ """How much of a component's work has a live instance covering it."""
+
+ HEALTHY = "healthy"
+ DEGRADED = "degraded"
+ DOWN = "down"
@provide_session
@@ -54,14 +70,26 @@ def _job_instance_health(job: Job, heartbeat_field_name:
str) -> dict[str, Any]:
heartbeat = job.latest_heartbeat.isoformat() if job.latest_heartbeat else
None
return {
"hostname": job.hostname,
- "status": HEALTHY if job.is_alive() else UNHEALTHY,
heartbeat_field_name: heartbeat,
}
-def _legacy_status(jobs: list[Job]) -> str:
+def _live_jobs(jobs: list[Job]) -> list[Job]:
+ """
+ Narrow unfinished job rows down to the replicas that are actually running.
+
+ ``instances`` describes the deployment as it is now, so a row left behind
by a replica that
+ never got to write its ``end_date`` must not be listed: the host it names
is gone, and under an
+ orchestrator that assigns a fresh hostname per restart it will never come
back. The unfinished
+ rows are still what ``status`` and ``latest_*_heartbeat`` are derived
from, so how long a dead
+ component has been silent stays visible even once it drops out of
``instances``.
+ """
+ return [job for job in jobs if job.is_alive()]
+
+
+def _legacy_status(jobs: list[Job]) -> HealthStatus:
"""Top-level status: healthy if any instance is alive."""
- return HEALTHY if any(job.is_alive() for job in jobs) else UNHEALTHY
+ return HealthStatus.HEALTHY if any(job.is_alive() for job in jobs) else
HealthStatus.UNHEALTHY
def _triggerer_instance_health(job: Job) -> dict[str, Any]:
@@ -78,19 +106,74 @@ def _dag_processor_instance_health(job: Job) -> dict[str,
Any]:
}
-def _aggregate_detailed_status(jobs: list[Job]) -> str:
- """detailed_status: healthy (all alive), degraded (some alive), down (none
alive)."""
- alive_count = sum(1 for job in jobs if job.is_alive())
- if alive_count == 0:
- return DOWN
- if alive_count == len(jobs):
- return HEALTHY
- return DEGRADED
+# ``detailed_status`` answers "is every part of this component's work being
done", which needs a
+# denominator. Counting job rows cannot supply one: ``end_date`` is only
written by a cooperative
+# shutdown, so a replica lost to SIGKILL, an OOM kill, or a node eviction
leaves an unfinished row
+# behind forever and a restarted replica adds a second one. The denominator is
therefore taken from
+# the declared work partition instead, which is unaffected by how replicas
come and go:
+#
+# * Dag processor -- the bundles in ``[dag_processor] dag_bundle_config_list``.
+# * Triggerer -- the team scopes those bundles declare, since a triggerer only
picks up triggers for
+# its own team (see ``Trigger.ids_for_triggerer``).
+# * Scheduler -- schedulers are symmetric, so there is no partition and no
partial state to report.
+
+
+def _configured_bundle_teams() -> dict[str, str | None]:
+ """Map every configured Dag bundle to the team owning it, empty when the
config is unreadable."""
+ try:
+ return _get_configured_bundle_team_names()
+ except Exception:
+ # A health probe must not fail on malformed bundle config; callers
fall back to liveness only.
+ log.warning("Could not read the Dag bundle configuration",
exc_info=True)
+ return {}
+
+
+def _liveness_status(jobs: list[Job]) -> DetailedHealthStatus:
+ """Status for a component with no declared work partition: one live
replica covers everything."""
+ return DetailedHealthStatus.HEALTHY if any(job.is_alive() for job in jobs)
else DetailedHealthStatus.DOWN
+
+
+def _coverage_status(expected: set[Any], covered: set[Any]) ->
DetailedHealthStatus:
+ """Status from how much of a component's declared work partition its live
replicas cover."""
+ if not expected - covered:
+ return DetailedHealthStatus.HEALTHY
+ if expected & covered:
+ return DetailedHealthStatus.DEGRADED
+ return DetailedHealthStatus.DOWN
+
+
+def _dag_processor_detailed_status(jobs: list[Job]) -> DetailedHealthStatus:
+ """Status from bundle coverage, which partitions processor work with or
without multi-team mode."""
+ expected = set(_configured_bundle_teams())
+ if not expected:
+ return _liveness_status(jobs)
+
+ covered: set[str] = set()
+ for job in jobs:
+ if job.is_alive():
+ # A processor started without ``--bundle-name`` parses every
configured bundle.
+ covered |= expected if not job.bundle_names else expected &
set(job.bundle_names)
+ return _coverage_status(expected, covered)
+
+
+def _triggerer_detailed_status(jobs: list[Job]) -> DetailedHealthStatus:
+ """Status from team coverage, since a triggerer only picks up triggers for
its own team."""
+ if not conf.getboolean("core", "multi_team"):
+ # Outside multi-team mode no team filter is applied, so any live
triggerer serves every trigger.
+ return _liveness_status(jobs)
+
+ # ``None`` is a scope in its own right: triggers from bundles that declare
no team are only picked
+ # up by a triggerer started without ``--team-name``, so it must stay in
the expected set.
+ expected = set(_configured_bundle_teams().values())
+ if not expected:
+ return _liveness_status(jobs)
+
+ return _coverage_status(expected, {job.team_name for job in jobs if
job.is_alive()})
def get_airflow_health() -> dict[str, Any]:
"""Get the health for Airflow metadatabase, scheduler, triggerer, and dag
processor."""
- metadatabase_status = HEALTHY
+ metadatabase_status = HealthStatus.HEALTHY
latest_scheduler_heartbeat = None
latest_triggerer_heartbeat = None
@@ -100,55 +183,52 @@ def get_airflow_health() -> dict[str, Any]:
triggerer_instances: list[dict[str, Any]] | None = None
dag_processor_instances: list[dict[str, Any]] | None = None
- scheduler_status = UNHEALTHY
- triggerer_status = UNHEALTHY
- dag_processor_status = UNHEALTHY
+ scheduler_status = HealthStatus.UNHEALTHY
+ triggerer_status = HealthStatus.UNHEALTHY
+ dag_processor_status = HealthStatus.UNHEALTHY
- scheduler_detailed_status = DOWN
- triggerer_detailed_status = DOWN
- dag_processor_detailed_status = DOWN
+ scheduler_detailed_status = DetailedHealthStatus.DOWN
+ triggerer_detailed_status = DetailedHealthStatus.DOWN
+ dag_processor_detailed_status = DetailedHealthStatus.DOWN
try:
scheduler_jobs = get_jobs_health(SchedulerJobRunner)
scheduler_status = _legacy_status(scheduler_jobs)
- scheduler_detailed_status = _aggregate_detailed_status(scheduler_jobs)
- if scheduler_jobs:
+ scheduler_detailed_status = _liveness_status(scheduler_jobs)
+ if scheduler_jobs and scheduler_jobs[0].latest_heartbeat:
+ latest_scheduler_heartbeat =
scheduler_jobs[0].latest_heartbeat.isoformat()
+ if live_scheduler_jobs := _live_jobs(scheduler_jobs):
scheduler_instances = [
- _job_instance_health(job, "latest_scheduler_heartbeat") for
job in scheduler_jobs
+ _job_instance_health(job, "latest_scheduler_heartbeat") for
job in live_scheduler_jobs
]
-
- if scheduler_jobs[0].latest_heartbeat:
- latest_scheduler_heartbeat =
scheduler_jobs[0].latest_heartbeat.isoformat()
except Exception:
- metadatabase_status = UNHEALTHY
+ metadatabase_status = HealthStatus.UNHEALTHY
try:
triggerer_jobs = get_jobs_health(TriggererJobRunner)
triggerer_status = _legacy_status(triggerer_jobs)
- triggerer_detailed_status = _aggregate_detailed_status(triggerer_jobs)
- if triggerer_jobs:
- triggerer_instances = [_triggerer_instance_health(job) for job in
triggerer_jobs]
-
- if triggerer_jobs[0].latest_heartbeat:
- latest_triggerer_heartbeat =
triggerer_jobs[0].latest_heartbeat.isoformat()
+ triggerer_detailed_status = _triggerer_detailed_status(triggerer_jobs)
+ if triggerer_jobs and triggerer_jobs[0].latest_heartbeat:
+ latest_triggerer_heartbeat =
triggerer_jobs[0].latest_heartbeat.isoformat()
+ if live_triggerer_jobs := _live_jobs(triggerer_jobs):
+ triggerer_instances = [_triggerer_instance_health(job) for job in
live_triggerer_jobs]
except Exception:
- metadatabase_status = UNHEALTHY
- triggerer_status = UNHEALTHY
- triggerer_detailed_status = DOWN
+ metadatabase_status = HealthStatus.UNHEALTHY
+ triggerer_status = HealthStatus.UNHEALTHY
+ triggerer_detailed_status = DetailedHealthStatus.DOWN
try:
dag_processor_jobs = get_jobs_health(DagProcessorJobRunner)
dag_processor_status = _legacy_status(dag_processor_jobs)
- dag_processor_detailed_status =
_aggregate_detailed_status(dag_processor_jobs)
- if dag_processor_jobs:
- dag_processor_instances = [_dag_processor_instance_health(job) for
job in dag_processor_jobs]
-
- if dag_processor_jobs[0].latest_heartbeat:
- latest_dag_processor_heartbeat =
dag_processor_jobs[0].latest_heartbeat.isoformat()
+ dag_processor_detailed_status =
_dag_processor_detailed_status(dag_processor_jobs)
+ if dag_processor_jobs and dag_processor_jobs[0].latest_heartbeat:
+ latest_dag_processor_heartbeat =
dag_processor_jobs[0].latest_heartbeat.isoformat()
+ if live_dag_processor_jobs := _live_jobs(dag_processor_jobs):
+ dag_processor_instances = [_dag_processor_instance_health(job) for
job in live_dag_processor_jobs]
except Exception:
- metadatabase_status = UNHEALTHY
- dag_processor_status = UNHEALTHY
- dag_processor_detailed_status = DOWN
+ metadatabase_status = HealthStatus.UNHEALTHY
+ dag_processor_status = HealthStatus.UNHEALTHY
+ dag_processor_detailed_status = DetailedHealthStatus.DOWN
airflow_health_status = {
"metadatabase": {"status": metadatabase_status},
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/monitor.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/monitor.py
index 5bf60bd83bd..442a5e9f44a 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/monitor.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/monitor.py
@@ -16,23 +16,26 @@
# under the License.
from __future__ import annotations
+from airflow.api.common.airflow_health import DetailedHealthStatus,
HealthStatus
from airflow.api_fastapi.core_api.base import BaseModel
class BaseInfoResponse(BaseModel):
"""Base info serializer for responses."""
- status: str | None
+ status: HealthStatus | None
-class SchedulerInstanceInfoResponse(BaseInfoResponse):
+# Instances carry no status of their own: only running replicas are listed,
and how healthy the set
+# of them is together is what the component's ``status`` and
``detailed_status`` report.
+class SchedulerInstanceInfoResponse(BaseModel):
"""Scheduler instance info serializer for responses."""
hostname: str | None
latest_scheduler_heartbeat: str | None
-class TriggererInstanceInfoResponse(BaseInfoResponse):
+class TriggererInstanceInfoResponse(BaseModel):
"""Triggerer instance info serializer for responses."""
hostname: str | None
@@ -40,7 +43,7 @@ class TriggererInstanceInfoResponse(BaseInfoResponse):
team_name: str | None
-class DagProcessorInstanceInfoResponse(BaseInfoResponse):
+class DagProcessorInstanceInfoResponse(BaseModel):
"""Dag processor instance info serializer for responses."""
hostname: str | None
@@ -52,7 +55,7 @@ class SchedulerInfoResponse(BaseInfoResponse):
"""Scheduler info serializer for responses."""
latest_scheduler_heartbeat: str | None
- detailed_status: str | None
+ detailed_status: DetailedHealthStatus | None
instances: list[SchedulerInstanceInfoResponse] | None = None
@@ -60,7 +63,7 @@ class TriggererInfoResponse(BaseInfoResponse):
"""Triggerer info serializer for responses."""
latest_triggerer_heartbeat: str | None
- detailed_status: str | None
+ detailed_status: DetailedHealthStatus | None
instances: list[TriggererInstanceInfoResponse] | None = None
@@ -68,7 +71,7 @@ class DagProcessorInfoResponse(BaseInfoResponse):
"""DagProcessor info serializer for responses."""
latest_dag_processor_heartbeat: str | None
- detailed_status: str | None
+ detailed_status: DetailedHealthStatus | None
instances: list[DagProcessorInstanceInfoResponse] | None = None
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index e6c0274e61d..6b6bae9b1dc 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -11820,9 +11820,8 @@ components:
properties:
status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/HealthStatus'
- type: 'null'
- title: Status
type: object
required:
- status
@@ -14117,9 +14116,8 @@ components:
properties:
status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/HealthStatus'
- type: 'null'
- title: Status
latest_dag_processor_heartbeat:
anyOf:
- type: string
@@ -14127,9 +14125,8 @@ components:
title: Latest Dag Processor Heartbeat
detailed_status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/DetailedHealthStatus'
- type: 'null'
- title: Detailed Status
instances:
anyOf:
- items:
@@ -14146,11 +14143,6 @@ components:
description: DagProcessor info serializer for responses.
DagProcessorInstanceInfoResponse:
properties:
- status:
- anyOf:
- - type: string
- - type: 'null'
- title: Status
hostname:
anyOf:
- type: string
@@ -14170,7 +14162,6 @@ components:
title: Bundle Names
type: object
required:
- - status
- hostname
- latest_dag_processor_heartbeat
- bundle_names
@@ -14452,6 +14443,14 @@ components:
This is the set of allowable values for the ``warning_type`` field
in the DagWarning model.'
+ DetailedHealthStatus:
+ type: string
+ enum:
+ - healthy
+ - degraded
+ - down
+ title: DetailedHealthStatus
+ description: How much of a component's work has a live instance covering
it.
DryRunBackfillCollectionResponse:
properties:
backfills:
@@ -14964,6 +14963,14 @@ components:
- triggerer
title: HealthInfoResponse
description: Health serializer for responses.
+ HealthStatus:
+ type: string
+ enum:
+ - healthy
+ - unhealthy
+ title: HealthStatus
+ description: 'Aggregate health of a component: whether it has at least
one live
+ instance.'
ImportErrorCollectionResponse:
properties:
import_errors:
@@ -15703,9 +15710,8 @@ components:
properties:
status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/HealthStatus'
- type: 'null'
- title: Status
latest_scheduler_heartbeat:
anyOf:
- type: string
@@ -15713,9 +15719,8 @@ components:
title: Latest Scheduler Heartbeat
detailed_status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/DetailedHealthStatus'
- type: 'null'
- title: Detailed Status
instances:
anyOf:
- items:
@@ -15732,11 +15737,6 @@ components:
description: Scheduler info serializer for responses.
SchedulerInstanceInfoResponse:
properties:
- status:
- anyOf:
- - type: string
- - type: 'null'
- title: Status
hostname:
anyOf:
- type: string
@@ -15749,7 +15749,6 @@ components:
title: Latest Scheduler Heartbeat
type: object
required:
- - status
- hostname
- latest_scheduler_heartbeat
title: SchedulerInstanceInfoResponse
@@ -16871,9 +16870,8 @@ components:
properties:
status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/HealthStatus'
- type: 'null'
- title: Status
latest_triggerer_heartbeat:
anyOf:
- type: string
@@ -16881,9 +16879,8 @@ components:
title: Latest Triggerer Heartbeat
detailed_status:
anyOf:
- - type: string
+ - $ref: '#/components/schemas/DetailedHealthStatus'
- type: 'null'
- title: Detailed Status
instances:
anyOf:
- items:
@@ -16900,11 +16897,6 @@ components:
description: Triggerer info serializer for responses.
TriggererInstanceInfoResponse:
properties:
- status:
- anyOf:
- - type: string
- - type: 'null'
- title: Status
hostname:
anyOf:
- type: string
@@ -16922,7 +16914,6 @@ components:
title: Team Name
type: object
required:
- - status
- hostname
- latest_triggerer_heartbeat
- team_name
diff --git a/airflow-core/src/airflow/dag_processing/bundles/manager.py
b/airflow-core/src/airflow/dag_processing/bundles/manager.py
index 292bb440a21..2c20f923739 100644
--- a/airflow-core/src/airflow/dag_processing/bundles/manager.py
+++ b/airflow-core/src/airflow/dag_processing/bundles/manager.py
@@ -130,6 +130,32 @@ def _parse_bundle_config(config_list) ->
list[_ExternalBundleConfig]:
return list(bundles.values())
+def _read_bundle_config_list() -> list[_ExternalBundleConfig]:
+ config_list = conf.getjson("dag_processor", "dag_bundle_config_list")
+ if not config_list:
+ return []
+ if not isinstance(config_list, list):
+ raise AirflowConfigException(
+ "Section `dag_processor` key `dag_bundle_config_list` "
+ f"must be list but got {config_list.__class__}"
+ )
+ return _parse_bundle_config(config_list)
+
+
+def _get_configured_bundle_team_names() -> dict[str, str | None]:
+ """
+ Get the team owning each explicitly configured Dag bundle.
+
+ This reads the config rather than going through ``DagBundlesManager`` so
that callers who only
+ need the declared bundle partition neither import every bundle class nor
see the example-Dag
+ bundles that ``DagBundlesManager.parse_config`` injects when ``[core]
load_examples`` is set --
+ those are added by Airflow, not declared by the deployment.
+
+ :return: mapping of bundle name to team name, ``None`` for bundles that
are not team scoped.
+ """
+ return {cfg.name: cfg.team_name for cfg in _read_bundle_config_list()}
+
+
def _add_example_dag_bundle(bundle_config_list: list[_ExternalBundleConfig]):
from airflow import example_dags
@@ -274,15 +300,9 @@ class DagBundlesManager(LoggingMixin):
if self._bundle_config:
return
- config_list = conf.getjson("dag_processor", "dag_bundle_config_list")
- if not config_list:
+ bundle_config_list = _read_bundle_config_list()
+ if not bundle_config_list:
return
- if not isinstance(config_list, list):
- raise AirflowConfigException(
- "Section `dag_processor` key `dag_bundle_config_list` "
- f"must be list but got {config_list.__class__}"
- )
- bundle_config_list = _parse_bundle_config(config_list)
if conf.getboolean("core", "LOAD_EXAMPLES"):
_add_example_dag_bundle(bundle_config_list)
_add_provider_example_dags_to_bundle(bundle_config_list)
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 075af8fe558..4771e7de12e 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -1023,13 +1023,12 @@ export const $BaseInfoResponse = {
status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/HealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Status'
+ ]
}
},
type: 'object',
@@ -4523,13 +4522,12 @@ export const $DagProcessorInfoResponse = {
status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/HealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Status'
+ ]
},
latest_dag_processor_heartbeat: {
anyOf: [
@@ -4545,13 +4543,12 @@ export const $DagProcessorInfoResponse = {
detailed_status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/DetailedHealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Detailed Status'
+ ]
},
instances: {
anyOf: [
@@ -4576,17 +4573,6 @@ export const $DagProcessorInfoResponse = {
export const $DagProcessorInstanceInfoResponse = {
properties: {
- status: {
- anyOf: [
- {
- type: 'string'
- },
- {
- type: 'null'
- }
- ],
- title: 'Status'
- },
hostname: {
anyOf: [
{
@@ -4625,7 +4611,7 @@ export const $DagProcessorInstanceInfoResponse = {
}
},
type: 'object',
- required: ['status', 'hostname', 'latest_dag_processor_heartbeat',
'bundle_names'],
+ required: ['hostname', 'latest_dag_processor_heartbeat', 'bundle_names'],
title: 'DagProcessorInstanceInfoResponse',
description: 'Dag processor instance info serializer for responses.'
} as const;
@@ -4957,6 +4943,13 @@ This is the set of allowable values for the
\`\`warning_type\`\` field
in the DagWarning model.`
} as const;
+export const $DetailedHealthStatus = {
+ type: 'string',
+ enum: ['healthy', 'degraded', 'down'],
+ title: 'DetailedHealthStatus',
+ description: "How much of a component's work has a live instance covering
it."
+} as const;
+
export const $DryRunBackfillCollectionResponse = {
properties: {
backfills: {
@@ -5730,6 +5723,13 @@ export const $HealthInfoResponse = {
description: 'Health serializer for responses.'
} as const;
+export const $HealthStatus = {
+ type: 'string',
+ enum: ['healthy', 'unhealthy'],
+ title: 'HealthStatus',
+ description: 'Aggregate health of a component: whether it has at least one
live instance.'
+} as const;
+
export const $ImportErrorCollectionResponse = {
properties: {
import_errors: {
@@ -6830,13 +6830,12 @@ export const $SchedulerInfoResponse = {
status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/HealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Status'
+ ]
},
latest_scheduler_heartbeat: {
anyOf: [
@@ -6852,13 +6851,12 @@ export const $SchedulerInfoResponse = {
detailed_status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/DetailedHealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Detailed Status'
+ ]
},
instances: {
anyOf: [
@@ -6883,17 +6881,6 @@ export const $SchedulerInfoResponse = {
export const $SchedulerInstanceInfoResponse = {
properties: {
- status: {
- anyOf: [
- {
- type: 'string'
- },
- {
- type: 'null'
- }
- ],
- title: 'Status'
- },
hostname: {
anyOf: [
{
@@ -6918,7 +6905,7 @@ export const $SchedulerInstanceInfoResponse = {
}
},
type: 'object',
- required: ['status', 'hostname', 'latest_scheduler_heartbeat'],
+ required: ['hostname', 'latest_scheduler_heartbeat'],
title: 'SchedulerInstanceInfoResponse',
description: 'Scheduler instance info serializer for responses.'
} as const;
@@ -8679,13 +8666,12 @@ export const $TriggererInfoResponse = {
status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/HealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Status'
+ ]
},
latest_triggerer_heartbeat: {
anyOf: [
@@ -8701,13 +8687,12 @@ export const $TriggererInfoResponse = {
detailed_status: {
anyOf: [
{
- type: 'string'
+ '$ref': '#/components/schemas/DetailedHealthStatus'
},
{
type: 'null'
}
- ],
- title: 'Detailed Status'
+ ]
},
instances: {
anyOf: [
@@ -8732,17 +8717,6 @@ export const $TriggererInfoResponse = {
export const $TriggererInstanceInfoResponse = {
properties: {
- status: {
- anyOf: [
- {
- type: 'string'
- },
- {
- type: 'null'
- }
- ],
- title: 'Status'
- },
hostname: {
anyOf: [
{
@@ -8778,7 +8752,7 @@ export const $TriggererInstanceInfoResponse = {
}
},
type: 'object',
- required: ['status', 'hostname', 'latest_triggerer_heartbeat',
'team_name'],
+ required: ['hostname', 'latest_triggerer_heartbeat', 'team_name'],
title: 'TriggererInstanceInfoResponse',
description: 'Triggerer instance info serializer for responses.'
} as const;
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index f74d0a87cab..dc08fe4d80e 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -303,7 +303,7 @@ export type BackfillResponse = {
* Base info serializer for responses.
*/
export type BaseInfoResponse = {
- status: string | null;
+ status: HealthStatus | null;
};
/**
@@ -1181,9 +1181,9 @@ export type DagBundleResponse = {
* DagProcessor info serializer for responses.
*/
export type DagProcessorInfoResponse = {
- status: string | null;
+ status: HealthStatus | null;
latest_dag_processor_heartbeat: string | null;
- detailed_status: string | null;
+ detailed_status: DetailedHealthStatus | null;
instances?: Array<DagProcessorInstanceInfoResponse> | null;
};
@@ -1191,7 +1191,6 @@ export type DagProcessorInfoResponse = {
* Dag processor instance info serializer for responses.
*/
export type DagProcessorInstanceInfoResponse = {
- status: string | null;
hostname: string | null;
latest_dag_processor_heartbeat: string | null;
bundle_names: Array<(string)> | null;
@@ -1311,6 +1310,11 @@ export type DagVersionResponse = {
*/
export type DagWarningType = 'asset conflict' | 'duplicate dag id' |
'non-existent pool' | 'runtime varying value';
+/**
+ * How much of a component's work has a live instance covering it.
+ */
+export type DetailedHealthStatus = 'healthy' | 'degraded' | 'down';
+
/**
* Backfill collection serializer for responses in dry-run mode.
*/
@@ -1510,6 +1514,11 @@ export type HealthInfoResponse = {
dag_processor?: DagProcessorInfoResponse | null;
};
+/**
+ * Aggregate health of a component: whether it has at least one live instance.
+ */
+export type HealthStatus = 'healthy' | 'unhealthy';
+
/**
* Import Error Collection Response.
*/
@@ -1783,9 +1792,9 @@ export type ReprocessBehavior = 'failed' | 'completed' |
'none';
* Scheduler info serializer for responses.
*/
export type SchedulerInfoResponse = {
- status: string | null;
+ status: HealthStatus | null;
latest_scheduler_heartbeat: string | null;
- detailed_status: string | null;
+ detailed_status: DetailedHealthStatus | null;
instances?: Array<SchedulerInstanceInfoResponse> | null;
};
@@ -1793,7 +1802,6 @@ export type SchedulerInfoResponse = {
* Scheduler instance info serializer for responses.
*/
export type SchedulerInstanceInfoResponse = {
- status: string | null;
hostname: string | null;
latest_scheduler_heartbeat: string | null;
};
@@ -2140,9 +2148,9 @@ export type TriggerResponse = {
* Triggerer info serializer for responses.
*/
export type TriggererInfoResponse = {
- status: string | null;
+ status: HealthStatus | null;
latest_triggerer_heartbeat: string | null;
- detailed_status: string | null;
+ detailed_status: DetailedHealthStatus | null;
instances?: Array<TriggererInstanceInfoResponse> | null;
};
@@ -2150,7 +2158,6 @@ export type TriggererInfoResponse = {
* Triggerer instance info serializer for responses.
*/
export type TriggererInstanceInfoResponse = {
- status: string | null;
hostname: string | null;
latest_triggerer_heartbeat: string | null;
team_name: string | null;
diff --git a/airflow-core/tests/unit/api/common/test_airflow_health.py
b/airflow-core/tests/unit/api/common/test_airflow_health.py
index 79648ac99ac..cb10bcb79ac 100644
--- a/airflow-core/tests/unit/api/common/test_airflow_health.py
+++ b/airflow-core/tests/unit/api/common/test_airflow_health.py
@@ -23,10 +23,11 @@ import pytest
from airflow._shared.timezones import timezone
from airflow.api.common.airflow_health import (
- DEGRADED,
- DOWN,
- HEALTHY,
- UNHEALTHY,
+ DetailedHealthStatus,
+ HealthStatus,
+ _configured_bundle_teams,
+ _dag_processor_detailed_status,
+ _triggerer_detailed_status,
get_airflow_health,
get_jobs_health,
)
@@ -35,6 +36,7 @@ from airflow.jobs.scheduler_job_runner import
SchedulerJobRunner
from airflow.jobs.triggerer_job_runner import TriggererJobRunner
from airflow.utils.session import provide_session
+from tests_common.test_utils.config import conf_vars
from tests_common.test_utils.db import clear_db_jobs
pytestmark = pytest.mark.db_test
@@ -42,6 +44,14 @@ pytestmark = pytest.mark.db_test
STALE_HEARTBEAT_AGE = timedelta(minutes=5)
+def _patch_bundles(bundle_teams: dict[str, str | None]):
+ """Patch the declared bundle/team partition that ``detailed_status`` is
measured against."""
+ return patch(
+ "airflow.api.common.airflow_health._configured_bundle_teams",
+ return_value=bundle_teams,
+ )
+
+
def _mock_job(
*,
hostname: str,
@@ -61,9 +71,9 @@ def _mock_job(
def _empty_component(heartbeat_field: str) -> dict:
return {
- "status": UNHEALTHY,
+ "status": HealthStatus.UNHEALTHY,
heartbeat_field: None,
- "detailed_status": DOWN,
+ "detailed_status": DetailedHealthStatus.DOWN,
"instances": None,
}
@@ -131,7 +141,7 @@ def test_get_airflow_health_no_jobs(mock_get_jobs_health):
health_status = get_airflow_health()
assert health_status == {
- "metadatabase": {"status": HEALTHY},
+ "metadatabase": {"status": HealthStatus.HEALTHY},
"scheduler": _empty_component("latest_scheduler_heartbeat"),
"triggerer": _empty_component("latest_triggerer_heartbeat"),
"dag_processor": _empty_component("latest_dag_processor_heartbeat"),
@@ -143,7 +153,7 @@ def
test_get_airflow_health_metadatabase_unhealthy(mock_get_jobs_health):
health_status = get_airflow_health()
assert health_status == {
- "metadatabase": {"status": UNHEALTHY},
+ "metadatabase": {"status": HealthStatus.UNHEALTHY},
"scheduler": _empty_component("latest_scheduler_heartbeat"),
"triggerer": _empty_component("latest_triggerer_heartbeat"),
"dag_processor": _empty_component("latest_dag_processor_heartbeat"),
@@ -156,15 +166,14 @@ def
test_get_airflow_health_one_alive_job(mock_get_jobs_health):
health_status = get_airflow_health()
assert health_status == {
- "metadatabase": {"status": HEALTHY},
+ "metadatabase": {"status": HealthStatus.HEALTHY},
"scheduler": {
- "status": HEALTHY,
+ "status": HealthStatus.HEALTHY,
"latest_scheduler_heartbeat":
ALIVE_SCHEDULER_JOB_MOCK.latest_heartbeat.isoformat(),
- "detailed_status": HEALTHY,
+ "detailed_status": DetailedHealthStatus.HEALTHY,
"instances": [
{
"hostname": ALIVE_SCHEDULER_JOB_MOCK.hostname,
- "status": HEALTHY,
"latest_scheduler_heartbeat":
ALIVE_SCHEDULER_JOB_MOCK.latest_heartbeat.isoformat(),
}
],
@@ -183,24 +192,15 @@ def
test_get_airflow_health_mixed_alive_and_stale_jobs(mock_get_jobs_health):
]
health_status = get_airflow_health()
- assert health_status["scheduler"]["status"] == HEALTHY
- assert health_status["scheduler"]["detailed_status"] == DEGRADED
+ assert health_status["scheduler"]["status"] == HealthStatus.HEALTHY
+ # Schedulers are symmetric, so a stale row alongside a live one is not a
partial outage.
+ assert health_status["scheduler"]["detailed_status"] ==
DetailedHealthStatus.HEALTHY
+ # The rows the two dead replicas left behind name hosts that are gone, so
they are not listed.
assert health_status["scheduler"]["instances"] == [
{
"hostname": ALIVE_SCHEDULER_JOB_MOCK.hostname,
- "status": HEALTHY,
"latest_scheduler_heartbeat":
ALIVE_SCHEDULER_JOB_MOCK.latest_heartbeat.isoformat(),
},
- {
- "hostname": STALE_SCHEDULER_JOB_MOCK.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat":
STALE_SCHEDULER_JOB_MOCK.latest_heartbeat.isoformat(),
- },
- {
- "hostname": STALE_SCHEDULER_JOB_MOCK_2.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat":
STALE_SCHEDULER_JOB_MOCK_2.latest_heartbeat.isoformat(),
- },
]
assert (
health_status["scheduler"]["latest_scheduler_heartbeat"]
@@ -219,25 +219,14 @@ def
test_get_airflow_health_all_stale_jobs(mock_get_jobs_health):
]
health_status = get_airflow_health()
- assert health_status["scheduler"]["status"] == UNHEALTHY
- assert health_status["scheduler"]["detailed_status"] == DOWN
- assert health_status["scheduler"]["instances"] == [
- {
- "hostname": STALE_SCHEDULER_JOB_MOCK.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat":
STALE_SCHEDULER_JOB_MOCK.latest_heartbeat.isoformat(),
- },
- {
- "hostname": STALE_SCHEDULER_JOB_MOCK_2.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat":
STALE_SCHEDULER_JOB_MOCK_2.latest_heartbeat.isoformat(),
- },
- {
- "hostname": STALE_SCHEDULER_JOB_MOCK_3.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat":
STALE_SCHEDULER_JOB_MOCK_3.latest_heartbeat.isoformat(),
- },
- ]
+ assert health_status["scheduler"]["status"] == HealthStatus.UNHEALTHY
+ assert health_status["scheduler"]["detailed_status"] ==
DetailedHealthStatus.DOWN
+ assert health_status["scheduler"]["instances"] is None
+ # A component with no live replica still reports when it was last heard
from.
+ assert (
+ health_status["scheduler"]["latest_scheduler_heartbeat"]
+ == STALE_SCHEDULER_JOB_MOCK.latest_heartbeat.isoformat()
+ )
@patch("airflow.api.common.airflow_health.get_jobs_health")
@@ -245,21 +234,15 @@ def
test_get_airflow_health_mixed_triggerers_include_team_name(mock_get_jobs_hea
mock_get_jobs_health.side_effect = [[], [ALIVE_TRIGGERER_JOB_MOCK,
STALE_TRIGGERER_JOB_MOCK], []]
health_status = get_airflow_health()
- assert health_status["triggerer"]["status"] == HEALTHY
- assert health_status["triggerer"]["detailed_status"] == DEGRADED
+ assert health_status["triggerer"]["status"] == HealthStatus.HEALTHY
+ # Multi-team is off here, so every live triggerer serves every trigger
regardless of team.
+ assert health_status["triggerer"]["detailed_status"] ==
DetailedHealthStatus.HEALTHY
assert health_status["triggerer"]["instances"] == [
{
"hostname": ALIVE_TRIGGERER_JOB_MOCK.hostname,
- "status": HEALTHY,
"latest_triggerer_heartbeat":
ALIVE_TRIGGERER_JOB_MOCK.latest_heartbeat.isoformat(),
"team_name": ALIVE_TRIGGERER_JOB_MOCK.team_name,
},
- {
- "hostname": STALE_TRIGGERER_JOB_MOCK.hostname,
- "status": UNHEALTHY,
- "latest_triggerer_heartbeat":
STALE_TRIGGERER_JOB_MOCK.latest_heartbeat.isoformat(),
- "team_name": STALE_TRIGGERER_JOB_MOCK.team_name,
- },
]
assert (
health_status["triggerer"]["latest_triggerer_heartbeat"]
@@ -267,35 +250,34 @@ def
test_get_airflow_health_mixed_triggerers_include_team_name(mock_get_jobs_hea
)
+@_patch_bundles({"bundle-a": None})
@patch("airflow.api.common.airflow_health.get_jobs_health")
-def
test_get_airflow_health_triggerer_and_dag_processor_healthy(mock_get_jobs_health):
+def
test_get_airflow_health_triggerer_and_dag_processor_healthy(mock_get_jobs_health,
_mock_bundles):
mock_get_jobs_health.side_effect = [[], [ALIVE_TRIGGERER_JOB_MOCK],
[ALIVE_DAG_PROCESSOR_JOB_MOCK]]
health_status = get_airflow_health()
assert health_status == {
- "metadatabase": {"status": HEALTHY},
+ "metadatabase": {"status": HealthStatus.HEALTHY},
"scheduler": _empty_component("latest_scheduler_heartbeat"),
"triggerer": {
- "status": HEALTHY,
+ "status": HealthStatus.HEALTHY,
"latest_triggerer_heartbeat":
ALIVE_TRIGGERER_JOB_MOCK.latest_heartbeat.isoformat(),
- "detailed_status": HEALTHY,
+ "detailed_status": DetailedHealthStatus.HEALTHY,
"instances": [
{
"hostname": ALIVE_TRIGGERER_JOB_MOCK.hostname,
- "status": HEALTHY,
"latest_triggerer_heartbeat":
ALIVE_TRIGGERER_JOB_MOCK.latest_heartbeat.isoformat(),
"team_name": ALIVE_TRIGGERER_JOB_MOCK.team_name,
}
],
},
"dag_processor": {
- "status": HEALTHY,
+ "status": HealthStatus.HEALTHY,
"latest_dag_processor_heartbeat":
ALIVE_DAG_PROCESSOR_JOB_MOCK.latest_heartbeat.isoformat(),
- "detailed_status": HEALTHY,
+ "detailed_status": DetailedHealthStatus.HEALTHY,
"instances": [
{
"hostname": ALIVE_DAG_PROCESSOR_JOB_MOCK.hostname,
- "status": HEALTHY,
"latest_dag_processor_heartbeat":
ALIVE_DAG_PROCESSOR_JOB_MOCK.latest_heartbeat.isoformat(),
"bundle_names": ALIVE_DAG_PROCESSOR_JOB_MOCK.bundle_names,
}
@@ -304,6 +286,198 @@ def
test_get_airflow_health_triggerer_and_dag_processor_healthy(mock_get_jobs_he
}
+def _processor(*, alive: bool, bundle_names: list[str] | None) -> MagicMock:
+ return _mock_job(
+ hostname=f"dag-processor-{'alive' if alive else 'stale'}",
+ heartbeat=datetime(2024, 2, 1),
+ alive=alive,
+ bundle_names=bundle_names,
+ )
+
+
+def _triggerer(*, alive: bool, team_name: str | None) -> MagicMock:
+ return _mock_job(
+ hostname=f"triggerer-{team_name}-{'alive' if alive else 'stale'}",
+ heartbeat=datetime(2024, 2, 1),
+ alive=alive,
+ team_name=team_name,
+ )
+
+
+class TestDagProcessorDetailedStatus:
+ """``detailed_status`` measures which configured bundles a live processor
is parsing."""
+
+ BUNDLES = {"bundle-a": None, "bundle-b": None}
+
+ @pytest.mark.parametrize(
+ ("jobs", "expected"),
+ [
+ pytest.param(
+ [_processor(alive=True, bundle_names=["bundle-a",
"bundle-b"])],
+ DetailedHealthStatus.HEALTHY,
+ id="one_processor_parses_every_bundle",
+ ),
+ pytest.param(
+ [
+ _processor(alive=True, bundle_names=["bundle-a"]),
+ _processor(alive=True, bundle_names=["bundle-b"]),
+ ],
+ DetailedHealthStatus.HEALTHY,
+ id="a_processor_per_bundle",
+ ),
+ pytest.param(
+ [_processor(alive=True, bundle_names=None)],
+ DetailedHealthStatus.HEALTHY,
+ id="processor_without_bundle_name_parses_all",
+ ),
+ pytest.param(
+ [_processor(alive=True, bundle_names=[])],
+ DetailedHealthStatus.HEALTHY,
+ id="empty_bundle_names_parses_all",
+ ),
+ pytest.param(
+ [_processor(alive=True, bundle_names=["bundle-a"])],
+ DetailedHealthStatus.DEGRADED,
+ id="one_bundle_left_unparsed",
+ ),
+ pytest.param(
+ [_processor(alive=False, bundle_names=["bundle-a",
"bundle-b"])],
+ DetailedHealthStatus.DOWN,
+ id="only_processor_is_stale",
+ ),
+ pytest.param(
+ [_processor(alive=True,
bundle_names=["bundle-removed-from-config"])],
+ DetailedHealthStatus.DOWN,
+ id="processor_parses_nothing_configured",
+ ),
+ pytest.param([], DetailedHealthStatus.DOWN,
id="no_processor_at_all"),
+ ],
+ )
+ def test_bundle_coverage(self, jobs, expected):
+ with _patch_bundles(self.BUNDLES):
+ assert _dag_processor_detailed_status(jobs) == expected
+
+ def test_restarted_processor_leaves_the_component_healthy(self):
+ """A hard-killed replica keeps an unfinished job row; its replacement
must still read healthy."""
+ jobs = [
+ _processor(alive=True, bundle_names=["bundle-a", "bundle-b"]),
+ _processor(alive=False, bundle_names=["bundle-a", "bundle-b"]),
+ ]
+
+ with _patch_bundles(self.BUNDLES):
+ assert _dag_processor_detailed_status(jobs) ==
DetailedHealthStatus.HEALTHY
+
+ @pytest.mark.parametrize(
+ ("jobs", "expected"),
+ [
+ pytest.param(
+ [_processor(alive=True, bundle_names=None)],
+ DetailedHealthStatus.HEALTHY,
+ id="a_processor_is_alive",
+ ),
+ pytest.param(
+ [_processor(alive=False, bundle_names=None)],
+ DetailedHealthStatus.DOWN,
+ id="no_processor_is_alive",
+ ),
+ ],
+ )
+ def test_falls_back_to_liveness_without_configured_bundles(self, jobs,
expected):
+ with _patch_bundles({}):
+ assert _dag_processor_detailed_status(jobs) == expected
+
+
+class TestTriggererDetailedStatus:
+ """Under multi-team a triggerer only serves its own team, so every team
scope needs one."""
+
+ BUNDLES = {"bundle-a": "team-a", "bundle-b": "team-b", "bundle-shared":
None}
+
+ @pytest.mark.parametrize(
+ ("jobs", "expected"),
+ [
+ pytest.param(
+ [
+ _triggerer(alive=True, team_name="team-a"),
+ _triggerer(alive=True, team_name="team-b"),
+ _triggerer(alive=True, team_name=None),
+ ],
+ DetailedHealthStatus.HEALTHY,
+ id="every_team_scope_covered",
+ ),
+ pytest.param(
+ [
+ _triggerer(alive=True, team_name="team-a"),
+ _triggerer(alive=True, team_name=None),
+ ],
+ DetailedHealthStatus.DEGRADED,
+ id="one_team_has_no_triggerer",
+ ),
+ pytest.param(
+ [
+ _triggerer(alive=True, team_name="team-a"),
+ _triggerer(alive=True, team_name="team-b"),
+ ],
+ DetailedHealthStatus.DEGRADED,
+ id="unscoped_bundles_have_no_triggerer",
+ ),
+ pytest.param(
+ [_triggerer(alive=False, team_name="team-a")],
+ DetailedHealthStatus.DOWN,
+ id="no_triggerer_is_alive",
+ ),
+ pytest.param(
+ [_triggerer(alive=True, team_name="team-of-a-removed-bundle")],
+ DetailedHealthStatus.DOWN,
+ id="triggerer_serves_no_configured_team",
+ ),
+ pytest.param([], DetailedHealthStatus.DOWN,
id="no_triggerer_at_all"),
+ ],
+ )
+ def test_team_coverage(self, jobs, expected):
+ with conf_vars({("core", "multi_team"): "True"}),
_patch_bundles(self.BUNDLES):
+ assert _triggerer_detailed_status(jobs) == expected
+
+ def test_restarted_triggerer_leaves_the_component_healthy(self):
+ jobs = [
+ _triggerer(alive=True, team_name="team-a"),
+ _triggerer(alive=False, team_name="team-a"),
+ ]
+
+ with conf_vars({("core", "multi_team"): "True"}),
_patch_bundles({"bundle-a": "team-a"}):
+ assert _triggerer_detailed_status(jobs) ==
DetailedHealthStatus.HEALTHY
+
+ @pytest.mark.parametrize(
+ ("jobs", "expected"),
+ [
+ pytest.param(
+ [_triggerer(alive=True, team_name="team-a"),
_triggerer(alive=False, team_name="team-b")],
+ DetailedHealthStatus.HEALTHY,
+ id="a_triggerer_is_alive",
+ ),
+ pytest.param(
+ [_triggerer(alive=False, team_name=None)],
+ DetailedHealthStatus.DOWN,
+ id="no_triggerer_is_alive",
+ ),
+ ],
+ )
+ def test_ignores_teams_without_multi_team(self, jobs, expected):
+ with conf_vars({("core", "multi_team"): "False"}),
_patch_bundles(self.BUNDLES):
+ assert _triggerer_detailed_status(jobs) == expected
+
+ def test_falls_back_to_liveness_without_configured_bundles(self):
+ with conf_vars({("core", "multi_team"): "True"}), _patch_bundles({}):
+ assert (
+ _triggerer_detailed_status([_triggerer(alive=True,
team_name="team-a")])
+ == DetailedHealthStatus.HEALTHY
+ )
+
+
+@conf_vars({("dag_processor", "dag_bundle_config_list"): "not json"})
+def test_configured_bundle_teams_survives_unreadable_config():
+ assert _configured_bundle_teams() == {}
+
+
class TestAirflowHealthFromDb:
@pytest.fixture(autouse=True)
def cleanup_jobs(self):
@@ -346,7 +520,7 @@ class TestAirflowHealthFromDb:
health_status = get_airflow_health()
assert health_status == {
- "metadatabase": {"status": HEALTHY},
+ "metadatabase": {"status": HealthStatus.HEALTHY},
"scheduler": _empty_component("latest_scheduler_heartbeat"),
"triggerer": _empty_component("latest_triggerer_heartbeat"),
"dag_processor":
_empty_component("latest_dag_processor_heartbeat"),
@@ -367,13 +541,12 @@ class TestAirflowHealthFromDb:
health_status = get_airflow_health()
- assert health_status["metadatabase"]["status"] == HEALTHY
- assert health_status["scheduler"]["status"] == HEALTHY
- assert health_status["scheduler"]["detailed_status"] == HEALTHY
+ assert health_status["metadatabase"]["status"] == HealthStatus.HEALTHY
+ assert health_status["scheduler"]["status"] == HealthStatus.HEALTHY
+ assert health_status["scheduler"]["detailed_status"] ==
DetailedHealthStatus.HEALTHY
assert health_status["scheduler"]["instances"] == [
{
"hostname": job.hostname,
- "status": HEALTHY,
"latest_scheduler_heartbeat": heartbeat.isoformat(),
}
]
@@ -389,34 +562,21 @@ class TestAirflowHealthFromDb:
alive = _create_job(
session, SchedulerJobRunner, hostname="scheduler-alive",
heartbeat=alive_heartbeat
)
- stale = _create_job(
- session, SchedulerJobRunner, hostname="scheduler-stale",
heartbeat=stale_heartbeat
- )
- older_stale = _create_job(
+ _create_job(session, SchedulerJobRunner, hostname="scheduler-stale",
heartbeat=stale_heartbeat)
+ _create_job(
session, SchedulerJobRunner, hostname="scheduler-stale-2",
heartbeat=older_stale_heartbeat
)
session.commit()
health_status = get_airflow_health()
- assert health_status["scheduler"]["status"] == HEALTHY
- assert health_status["scheduler"]["detailed_status"] == DEGRADED
+ assert health_status["scheduler"]["status"] == HealthStatus.HEALTHY
+ assert health_status["scheduler"]["detailed_status"] ==
DetailedHealthStatus.HEALTHY
assert health_status["scheduler"]["instances"] == [
{
"hostname": alive.hostname,
- "status": HEALTHY,
"latest_scheduler_heartbeat": alive_heartbeat.isoformat(),
},
- {
- "hostname": stale.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat": stale_heartbeat.isoformat(),
- },
- {
- "hostname": older_stale.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat":
older_stale_heartbeat.isoformat(),
- },
]
assert health_status["scheduler"]["latest_scheduler_heartbeat"] ==
alive_heartbeat.isoformat()
@@ -425,25 +585,17 @@ class TestAirflowHealthFromDb:
first = timezone.utcnow() - STALE_HEARTBEAT_AGE
second = first - timedelta(minutes=1)
third = second - timedelta(minutes=1)
- jobs = [
- _create_job(session, SchedulerJobRunner, hostname="stale-1",
heartbeat=first),
- _create_job(session, SchedulerJobRunner, hostname="stale-2",
heartbeat=second),
- _create_job(session, SchedulerJobRunner, hostname="stale-3",
heartbeat=third),
- ]
+ _create_job(session, SchedulerJobRunner, hostname="stale-1",
heartbeat=first)
+ _create_job(session, SchedulerJobRunner, hostname="stale-2",
heartbeat=second)
+ _create_job(session, SchedulerJobRunner, hostname="stale-3",
heartbeat=third)
session.commit()
health_status = get_airflow_health()
- assert health_status["scheduler"]["status"] == UNHEALTHY
- assert health_status["scheduler"]["detailed_status"] == DOWN
- assert health_status["scheduler"]["instances"] == [
- {
- "hostname": job.hostname,
- "status": UNHEALTHY,
- "latest_scheduler_heartbeat": job.latest_heartbeat.isoformat(),
- }
- for job in jobs
- ]
+ assert health_status["scheduler"]["status"] == HealthStatus.UNHEALTHY
+ assert health_status["scheduler"]["detailed_status"] ==
DetailedHealthStatus.DOWN
+ assert health_status["scheduler"]["instances"] is None
+ assert health_status["scheduler"]["latest_scheduler_heartbeat"] ==
first.isoformat()
@provide_session
def test_get_airflow_health_mixed_triggerers_include_team_name(self,
testing_team, *, session):
@@ -456,7 +608,7 @@ class TestAirflowHealthFromDb:
heartbeat=alive_heartbeat,
team_name=testing_team.name,
)
- stale = _create_job(
+ _create_job(
session,
TriggererJobRunner,
hostname="triggerer-stale",
@@ -467,20 +619,13 @@ class TestAirflowHealthFromDb:
health_status = get_airflow_health()
- assert health_status["triggerer"]["status"] == HEALTHY
- assert health_status["triggerer"]["detailed_status"] == DEGRADED
+ assert health_status["triggerer"]["status"] == HealthStatus.HEALTHY
+ assert health_status["triggerer"]["detailed_status"] ==
DetailedHealthStatus.HEALTHY
assert health_status["triggerer"]["instances"] == [
{
"hostname": alive.hostname,
- "status": HEALTHY,
"latest_triggerer_heartbeat": alive_heartbeat.isoformat(),
"team_name": testing_team.name,
},
- {
- "hostname": stale.hostname,
- "status": UNHEALTHY,
- "latest_triggerer_heartbeat": stale_heartbeat.isoformat(),
- "team_name": testing_team.name,
- },
]
assert health_status["triggerer"]["latest_triggerer_heartbeat"] ==
alive_heartbeat.isoformat()
diff --git
a/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py
b/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py
index 08e63c3f77d..7d23af9afef 100644
--- a/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py
+++ b/airflow-core/tests/unit/dag_processing/bundles/test_dag_bundle_manager.py
@@ -29,7 +29,11 @@ import pytest
from sqlalchemy import func, select, update
from airflow.dag_processing.bundles.base import BaseDagBundle
-from airflow.dag_processing.bundles.manager import DagBundlesManager,
_guess_best_bundle_for_fileloc
+from airflow.dag_processing.bundles.manager import (
+ DagBundlesManager,
+ _get_configured_bundle_team_names,
+ _guess_best_bundle_for_fileloc,
+)
from airflow.exceptions import AirflowConfigException
from airflow.models.dag import DagModel
from airflow.models.dag_version import DagVersion
@@ -141,6 +145,40 @@ OTHER_BUNDLE_CONFIG = [
}
]
+TEAM_BUNDLE_CONFIG = [
+ {
+ "name": "team-bundle",
+ "classpath":
"unit.dag_processing.bundles.test_dag_bundle_manager.BasicBundle",
+ "kwargs": {"refresh_interval": 1},
+ "team_name": "team-a",
+ },
+ {
+ "name": "unscoped-bundle",
+ "classpath":
"unit.dag_processing.bundles.test_dag_bundle_manager.BasicBundle",
+ "kwargs": {"refresh_interval": 1},
+ },
+]
+
+
[email protected]("load_examples", ["False", "True"])
+def test_get_configured_bundle_team_names(load_examples):
+ with conf_vars(
+ {
+ ("core", "load_examples"): load_examples,
+ ("core", "multi_team"): "True",
+ ("dag_processor", "dag_bundle_config_list"):
json.dumps(TEAM_BUNDLE_CONFIG),
+ }
+ ):
+ assert _get_configured_bundle_team_names() == {
+ "team-bundle": "team-a",
+ "unscoped-bundle": None,
+ }
+
+
+@conf_vars({("dag_processor", "dag_bundle_config_list"): "[]"})
+def test_get_configured_bundle_team_names_without_config():
+ assert _get_configured_bundle_team_names() == {}
+
def test_get_bundle():
"""Test that get_bundle builds and returns a bundle."""
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 60282ddef5d..5e340d7667d 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -130,14 +130,6 @@ class AsyncConnectionTestResponse(BaseModel):
created_at: Annotated[datetime, Field(title="Created At")]
-class BaseInfoResponse(BaseModel):
- """
- Base info serializer for responses.
- """
-
- status: Annotated[str | None, Field(title="Status")]
-
-
class BulkActionNotOnExistence(str, Enum):
"""
Bulk Action to be taken if the entity does not exist.
@@ -548,7 +540,6 @@ class DagProcessorInstanceInfoResponse(BaseModel):
Dag processor instance info serializer for responses.
"""
- status: Annotated[str | None, Field(title="Status")]
hostname: Annotated[str | None, Field(title="Hostname")]
latest_dag_processor_heartbeat: Annotated[str | None, Field(title="Latest
Dag Processor Heartbeat")]
bundle_names: Annotated[list[str] | None, Field(title="Bundle Names")]
@@ -705,6 +696,16 @@ class DagWarningType(str, Enum):
RUNTIME_VARYING_VALUE = "runtime varying value"
+class DetailedHealthStatus(str, Enum):
+ """
+ How much of a component's work has a live instance covering it.
+ """
+
+ HEALTHY = "healthy"
+ DEGRADED = "degraded"
+ DOWN = "down"
+
+
class DryRunBackfillResponse(BaseModel):
"""
Backfill serializer for responses in dry-run mode.
@@ -806,6 +807,15 @@ class HTTPExceptionResponse(BaseModel):
detail: Annotated[str | dict[str, Any], Field(title="Detail")]
+class HealthStatus(str, Enum):
+ """
+ Aggregate health of a component: whether it has at least one live instance.
+ """
+
+ HEALTHY = "healthy"
+ UNHEALTHY = "unhealthy"
+
+
class ImportErrorResponse(BaseModel):
"""
Import Error Response.
@@ -1030,7 +1040,6 @@ class SchedulerInstanceInfoResponse(BaseModel):
Scheduler instance info serializer for responses.
"""
- status: Annotated[str | None, Field(title="Status")]
hostname: Annotated[str | None, Field(title="Hostname")]
latest_scheduler_heartbeat: Annotated[str | None, Field(title="Latest
Scheduler Heartbeat")]
@@ -1244,7 +1253,6 @@ class TriggererInstanceInfoResponse(BaseModel):
Triggerer instance info serializer for responses.
"""
- status: Annotated[str | None, Field(title="Status")]
hostname: Annotated[str | None, Field(title="Hostname")]
latest_triggerer_heartbeat: Annotated[str | None, Field(title="Latest
Triggerer Heartbeat")]
team_name: Annotated[str | None, Field(title="Team Name")]
@@ -1541,6 +1549,14 @@ class BackfillResponse(BaseModel):
dag_display_name: Annotated[str, Field(title="Dag Display Name")]
+class BaseInfoResponse(BaseModel):
+ """
+ Base info serializer for responses.
+ """
+
+ status: HealthStatus | None
+
+
class BulkCreateActionConnectionBody(BaseModel):
model_config = ConfigDict(
extra="forbid",
@@ -2006,9 +2022,9 @@ class DagProcessorInfoResponse(BaseModel):
DagProcessor info serializer for responses.
"""
- status: Annotated[str | None, Field(title="Status")]
+ status: HealthStatus | None
latest_dag_processor_heartbeat: Annotated[str | None, Field(title="Latest
Dag Processor Heartbeat")]
- detailed_status: Annotated[str | None, Field(title="Detailed Status")]
+ detailed_status: DetailedHealthStatus | None
instances: Annotated[list[DagProcessorInstanceInfoResponse] | None,
Field(title="Instances")] = None
@@ -2181,9 +2197,9 @@ class SchedulerInfoResponse(BaseModel):
Scheduler info serializer for responses.
"""
- status: Annotated[str | None, Field(title="Status")]
+ status: HealthStatus | None
latest_scheduler_heartbeat: Annotated[str | None, Field(title="Latest
Scheduler Heartbeat")]
- detailed_status: Annotated[str | None, Field(title="Detailed Status")]
+ detailed_status: DetailedHealthStatus | None
instances: Annotated[list[SchedulerInstanceInfoResponse] | None,
Field(title="Instances")] = None
@@ -2320,9 +2336,9 @@ class TriggererInfoResponse(BaseModel):
Triggerer info serializer for responses.
"""
- status: Annotated[str | None, Field(title="Status")]
+ status: HealthStatus | None
latest_triggerer_heartbeat: Annotated[str | None, Field(title="Latest
Triggerer Heartbeat")]
- detailed_status: Annotated[str | None, Field(title="Detailed Status")]
+ detailed_status: DetailedHealthStatus | None
instances: Annotated[list[TriggererInstanceInfoResponse] | None,
Field(title="Instances")] = None