From a0bf1dec3b62d2153773908dfe1987f916434547 Mon Sep 17 00:00:00 2001 From: Jed Cunningham Date: Thu, 13 Mar 2025 00:32:43 -0600 Subject: [PATCH] Remove unused DAG reparse classmethods in asset manager This functionality was removed in #44866, but the methods weren't removed. The reparse functionality is being refactored, so removing them is easier than repairing them. --- airflow/assets/manager.py | 31 ------------------------------- 1 file changed, 31 deletions(-) diff --git a/airflow/assets/manager.py b/airflow/assets/manager.py index cecd76b07c3f4..f71e23d794a01 100644 --- a/airflow/assets/manager.py +++ b/airflow/assets/manager.py @@ -35,7 +35,6 @@ DagScheduleAssetReference, DagScheduleAssetUriReference, ) -from airflow.models.dagbag import DagPriorityParsingRequest from airflow.stats import Stats from airflow.utils.log.logging_mixin import LoggingMixin @@ -260,36 +259,6 @@ def _postgres_queue_dagruns(cls, asset_id: int, dags_to_queue: set[DagModel], se stmt = insert(AssetDagRunQueue).values(asset_id=asset_id).on_conflict_do_nothing() session.execute(stmt, values) - @classmethod - def _send_dag_priority_parsing_request(cls, file_locs: Iterable[str], session: Session) -> None: - if session.bind.dialect.name == "postgresql": - return cls._postgres_send_dag_priority_parsing_request(file_locs, session) - return cls._slow_path_send_dag_priority_parsing_request(file_locs, session) - - @classmethod - def _slow_path_send_dag_priority_parsing_request(cls, file_locs: Iterable[str], session: Session) -> None: - def _send_dag_priority_parsing_request_if_needed(fileloc: str) -> str | None: - # Don't error whole transaction when a single DagPriorityParsingRequest item conflicts. - # https://docs.sqlalchemy.org/en/14/orm/session_transaction.html#using-savepoint - req = DagPriorityParsingRequest(fileloc=fileloc) - try: - with session.begin_nested(): - session.merge(req) - except exc.IntegrityError: - cls.logger().debug("Skipping request %s, already present", req, exc_info=True) - return None - return req.fileloc - - for fileloc in file_locs: - _send_dag_priority_parsing_request_if_needed(fileloc) - - @classmethod - def _postgres_send_dag_priority_parsing_request(cls, file_locs: Iterable[str], session: Session) -> None: - from sqlalchemy.dialects.postgresql import insert - - stmt = insert(DagPriorityParsingRequest).on_conflict_do_nothing() - session.execute(stmt, [{"fileloc": fileloc} for fileloc in file_locs]) - def resolve_asset_manager() -> AssetManager: """Retrieve the asset manager."""