ferruzzi commented on code in PR #68917:
URL: https://github.com/apache/airflow/pull/68917#discussion_r3798272845
##########
airflow-core/src/airflow/serialization/definitions/dag.py:
##########
@@ -741,52 +744,112 @@ def _process_dagrun_deadline_alerts(
if not deadline_alert:
continue
- deserialized_deadline_alert = decode_deadline_alert(
- {
- Encoding.TYPE: DAT.DEADLINE_ALERT,
- Encoding.VAR: {
- DeadlineAlertFields.REFERENCE:
deadline_alert.reference,
- DeadlineAlertFields.INTERVAL: deadline_alert.interval,
- DeadlineAlertFields.CALLBACK:
deadline_alert.callback_def,
- },
- }
- )
-
- interval = deserialized_deadline_alert.interval
+ # Deadline creation is best-effort. A failure here must not
prevent the DagRun
+ # itself from being created. Use a plain try/except rather than
+ # ``session.begin_nested()`` since ``create_dagrun`` runs under
+ # ``prohibit_commit`` and releasing a SAVEPOINT would trip that
guard.
+ try:
+ deserialized_deadline_alert = decode_deadline_alert(
+ {
+ Encoding.TYPE: DAT.DEADLINE_ALERT,
+ Encoding.VAR: {
+ DeadlineAlertFields.REFERENCE:
deadline_alert.reference,
+ DeadlineAlertFields.INTERVAL:
deadline_alert.interval,
+ DeadlineAlertFields.CALLBACK:
deadline_alert.callback_def,
+ },
+ }
+ )
- if isinstance(interval, VariableInterval):
- interval = interval.resolve()
+ interval = deserialized_deadline_alert.interval
- if isinstance(deserialized_deadline_alert.reference,
SerializedReferenceModels.TYPES.DAGRUN):
- deadline_time =
deserialized_deadline_alert.reference.evaluate_with(
- session=session,
- interval=interval,
- # TODO : Pretty sure we can drop these last two; verify
after testing is complete
- dag_id=self.dag_id,
- run_id=orm_dagrun.run_id,
+ # Resolve the DagRun's team once, so a team-scoped
VariableInterval is looked up
+ # against the right team (not the global scope) and the stats
tag is consistent.
+ team_name = (
+ DagModel.get_team_name(self.dag_id, session=session)
+ if airflow_conf.getboolean("core", "multi_team")
+ else None
)
- if deadline_time is not None:
- session.add(
- Deadline(
- deadline_time=deadline_time,
- callback=deserialized_deadline_alert.callback,
- dagrun_id=orm_dagrun.id,
- deadline_alert_id=deadline_alert.id,
- dag_id=orm_dagrun.dag_id,
- bundle_name=orm_dagrun.dag_model.bundle_name,
- )
- )
- team_name = (
- DagModel.get_team_name(self.dag_id, session=session)
- if airflow_conf.getboolean("core", "multi_team")
- else None
- )
- stats.incr(
- "deadline_alerts.deadline_created",
- tags=prune_dict({"dag_id": self.dag_id, "team_name":
team_name}),
+ if isinstance(interval, VariableInterval):
+ interval = self._resolve_variable_interval(interval,
team_name=team_name, session=session)
+
+ if isinstance(deserialized_deadline_alert.reference,
SerializedReferenceModels.TYPES.DAGRUN):
+ deadline_time =
deserialized_deadline_alert.reference.evaluate_with(
+ session=session,
+ interval=interval,
+ # TODO : Pretty sure we can drop these last two;
verify after testing is complete
+ dag_id=self.dag_id,
+ run_id=orm_dagrun.run_id,
)
+ if deadline_time is not None:
+ session.add(
+ Deadline(
+ deadline_time=deadline_time,
+ callback=deserialized_deadline_alert.callback,
+ dagrun_id=orm_dagrun.id,
+ deadline_alert_id=deadline_alert.id,
+ dag_id=orm_dagrun.dag_id,
+ bundle_name=orm_dagrun.dag_model.bundle_name,
+ )
+ )
+ stats.incr(
+ "deadline_alerts.deadline_created",
+ tags=prune_dict({"dag_id": self.dag_id,
"team_name": team_name}),
Review Comment:
Nit: You have another metric `deadline_creation_failed` down on L808 which
isn't getting team_name. Maybe build a `metrics_tags = prune_dict({"dag_id":
self.dag_id, "team_name": team_name}` up above and reuse it for both, or add
`team_name` to the other one directly?
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1390,16 +1388,20 @@ def
test_dag_run_dag_versions_with_null_created_dag_version(self, dag_maker, ses
VariableInterval("my_key"),
],
)
- @mock.patch.object(Variable, "get")
@mock.patch.object(Deadline, "prune_deadlines")
- def test_dagrun_success_deadline(self, _, mock_get, interval, session,
deadline_test_dag):
+ def test_dagrun_success_deadline(self, _, interval, session,
deadline_test_dag):
def on_success_callable(context):
assert context["dag_run"].dag_id == "test_dag"
- future_date = datetime.datetime.now() + datetime.timedelta(days=365)
+ future_date = datetime.datetime(2037, 1, 1,
tzinfo=datetime.timezone.utc)
- # First value used during resolution
- mock_get.return_value = "5"
+ if isinstance(interval, VariableInterval):
+ # Seed via the metastore model, not the SDK Variable (whose set()
routes through
+ # SUPERVISOR_COMMS), so the row lands in the variable table on
this session.
+ from airflow.models.variable import Variable as VariableModel
Review Comment:
Looks like this (and other) local imports can be top-level??
##########
airflow-core/src/airflow/serialization/definitions/dag.py:
##########
@@ -741,52 +743,100 @@ def _process_dagrun_deadline_alerts(
if not deadline_alert:
continue
- deserialized_deadline_alert = decode_deadline_alert(
- {
- Encoding.TYPE: DAT.DEADLINE_ALERT,
- Encoding.VAR: {
- DeadlineAlertFields.REFERENCE:
deadline_alert.reference,
- DeadlineAlertFields.INTERVAL: deadline_alert.interval,
- DeadlineAlertFields.CALLBACK:
deadline_alert.callback_def,
- },
- }
- )
-
- interval = deserialized_deadline_alert.interval
+ # Deadline creation is best-effort. A failure here must not
prevent the DagRun
+ # itself from being created. Use a plain try/except rather than
+ # ``session.begin_nested()`` since ``create_dagrun`` runs under
+ # ``prohibit_commit`` and releasing a SAVEPOINT would trip that
guard.
+ try:
+ deserialized_deadline_alert = decode_deadline_alert(
+ {
+ Encoding.TYPE: DAT.DEADLINE_ALERT,
+ Encoding.VAR: {
+ DeadlineAlertFields.REFERENCE:
deadline_alert.reference,
+ DeadlineAlertFields.INTERVAL:
deadline_alert.interval,
+ DeadlineAlertFields.CALLBACK:
deadline_alert.callback_def,
+ },
+ }
+ )
- if isinstance(interval, VariableInterval):
- interval = interval.resolve()
+ interval = deserialized_deadline_alert.interval
- if isinstance(deserialized_deadline_alert.reference,
SerializedReferenceModels.TYPES.DAGRUN):
- deadline_time =
deserialized_deadline_alert.reference.evaluate_with(
- session=session,
- interval=interval,
- # TODO : Pretty sure we can drop these last two; verify
after testing is complete
- dag_id=self.dag_id,
- run_id=orm_dagrun.run_id,
+ # Resolve the DagRun's team once, so a team-scoped
VariableInterval is looked up
+ # against the right team (not the global scope) and the stats
tag is consistent.
+ team_name = (
+ DagModel.get_team_name(self.dag_id, session=session)
+ if airflow_conf.getboolean("core", "multi_team")
+ else None
)
- if deadline_time is not None:
- session.add(
- Deadline(
- deadline_time=deadline_time,
- callback=deserialized_deadline_alert.callback,
- dagrun_id=orm_dagrun.id,
- deadline_alert_id=deadline_alert.id,
- dag_id=orm_dagrun.dag_id,
- bundle_name=orm_dagrun.dag_model.bundle_name,
- )
- )
- team_name = (
- DagModel.get_team_name(self.dag_id, session=session)
- if airflow_conf.getboolean("core", "multi_team")
- else None
- )
- stats.incr(
- "deadline_alerts.deadline_created",
- tags=prune_dict({"dag_id": self.dag_id, "team_name":
team_name}),
+ if isinstance(interval, VariableInterval):
+ interval = self._resolve_variable_interval(interval,
team_name=team_name, session=session)
+
+ if isinstance(deserialized_deadline_alert.reference,
SerializedReferenceModels.TYPES.DAGRUN):
+ deadline_time =
deserialized_deadline_alert.reference.evaluate_with(
+ session=session,
+ interval=interval,
+ # TODO : Pretty sure we can drop these last two;
verify after testing is complete
+ dag_id=self.dag_id,
+ run_id=orm_dagrun.run_id,
)
+ if deadline_time is not None:
+ session.add(
+ Deadline(
+ deadline_time=deadline_time,
+ callback=deserialized_deadline_alert.callback,
+ dagrun_id=orm_dagrun.id,
+ deadline_alert_id=deadline_alert.id,
+ dag_id=orm_dagrun.dag_id,
+ bundle_name=orm_dagrun.dag_model.bundle_name,
+ )
+ )
+ stats.incr(
+ "deadline_alerts.deadline_created",
+ tags=prune_dict({"dag_id": self.dag_id,
"team_name": team_name}),
+ )
+ except Exception:
+ log.exception(
+ "Failed to create deadline for alert %s on DagRun %s
(dag_id=%s); "
+ "skipping this deadline, the DagRun is unaffected",
+ getattr(deadline_alert, "id", "<unknown>"),
+ orm_dagrun.run_id,
+ self.dag_id,
+ )
+ stats.incr("deadline_alerts.deadline_creation_failed",
tags={"dag_id": self.dag_id})
+
+ @staticmethod
+ def _resolve_variable_interval(
+ interval: VariableInterval, *, team_name: str | None, session: Session
+ ) -> datetime.timedelta:
+ """
+ Resolve a ``VariableInterval`` to a concrete ``timedelta`` at DagRun
creation.
+
+ The Variable is resolved using the standard secrets lookup order. The
scheduler
+ session is passed to the metastore backend to avoid creating a new
session
+ during DagRun creation.
+
+ :param interval: The ``VariableInterval`` to resolve.
+ :param team_name: Team owning the DagRun, forwarded to scope the
Variable lookup.
+ :param session: Scheduler session used for metadata database lookups.
+ :return: The resolved ``timedelta``.
+ :raises ValueError: If the Variable cannot be resolved or converted to
a valid ``timedelta``.
+ """
+ for backend in ensure_secrets_loaded():
Review Comment:
Sounds good. At this point it looks like it may be smarter to plumb
`session` into `Variable.get_variable_from_secrets` since we've basically
duplicated it here, but maybe you can just open an Issue to clean that up
instead of holding this PR up longer over that.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]