Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/common/partition_helpers.py: 31%
29 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
1# Licensed to the Apache Software Foundation (ASF) under one
2# or more contributor license agreements. See the NOTICE file
3# distributed with this work for additional information
4# regarding copyright ownership. The ASF licenses this file
5# to you under the Apache License, Version 2.0 (the
6# "License"); you may not use this file except in compliance
7# with the License. You may obtain a copy of the License at
8#
9# http://www.apache.org/licenses/LICENSE-2.0
10#
11# Unless required by applicable law or agreed to in writing,
12# software distributed under the License is distributed on an
13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14# KIND, either express or implied. See the License for the
15# specific language governing permissions and limitations
16# under the License.
17from __future__ import annotations
19from typing import TYPE_CHECKING
21import structlog
23from airflow.exceptions import DeserializationError
24from airflow.models.serialized_dag import SerializedDagModel
25from airflow.timetables.simple import PartitionedAssetTimetable
27if TYPE_CHECKING: 27 ↛ 28line 27 didn't jump to line 28 because the condition on line 27 was never true
28 from sqlalchemy.orm import Session
31log = structlog.get_logger(logger_name=__name__)
34def _extract_partitioned_timetable(serdag: SerializedDagModel) -> PartitionedAssetTimetable | None:
35 """Return the ``PartitionedAssetTimetable`` carried by *serdag*, or ``None``."""
36 try:
37 timetable = serdag.dag.timetable
38 except (DeserializationError, TypeError):
39 # ``DeserializationError`` covers structural serialization failures so
40 # a corrupted serialized Dag silently degrades to non-partitioned
41 # rather than 500-ing the read-only UI page. ``TypeError`` covers the
42 # eager validation in ``RollupMapper.__init__`` (raised during
43 # timetable deserialization when an upstream mapper / window pair is
44 # incompatible), so a misconfigured rollup Dag also degrades to
45 # non-partitioned here rather than 500-ing.
46 # ``KeyError`` / ``AttributeError`` / ``ImportError`` are intentionally
47 # not caught: refactor bugs that rename an attribute on ``serdag.dag``
48 # or break an import path must surface to the caller rather than
49 # silently downgrading the route to non-rollup.
50 log.warning("Failed to deserialize timetable for Dag", dag_id=serdag.dag_id, exc_info=True)
51 return None
52 if not timetable.partitioned:
53 return None
54 if TYPE_CHECKING:
55 assert isinstance(timetable, PartitionedAssetTimetable)
56 return timetable
59def load_partitioned_timetable(dag_id: str, session: Session) -> PartitionedAssetTimetable | None:
60 """
61 Return the PartitionedAssetTimetable for *dag_id*, or None if absent or not partitioned.
63 Callers gate this behind ``DagModel.has_rollup_mappers``, which is only
64 populated for ``PartitionedAssetTimetable``. The ``TYPE_CHECKING`` assert
65 narrows the type for mypy without a runtime ``isinstance`` cost.
66 """
67 serdag = SerializedDagModel.get(dag_id=dag_id, session=session)
68 if serdag is None:
69 return None
70 return _extract_partitioned_timetable(serdag)
73def load_partitioned_timetables(
74 dag_ids: list[str], session: Session
75) -> dict[str, PartitionedAssetTimetable | None]:
76 """
77 Batch-load PartitionedAssetTimetables for *dag_ids* in a single query.
79 Routes that already gate per-Dag on ``DagModel.has_rollup_mappers`` should
80 use this when iterating over many Dags so ``SerializedDagModel`` is hit
81 once instead of once per Dag. Returns a dict keyed by ``dag_id``; entries
82 whose timetable failed to deserialize or is not partitioned are ``None``.
83 """
84 if not dag_ids:
85 return {}
86 return {
87 serdag.dag_id: _extract_partitioned_timetable(serdag)
88 for serdag in SerializedDagModel.get_latest_serialized_dags(dag_ids=dag_ids, session=session)
89 }