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
4 changes: 3 additions & 1 deletion airflow-core/docs/migrations-ref.rst
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,9 @@ Here's the list of all the Database Migrations that are executed via when you ru
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| Revision ID | Revises ID | Airflow Version | Description |
+=========================+==================+===================+==============================================================+
| ``c4e7a1f9b2d0`` (head) | ``436dc127462c`` | ``3.4.0`` | Add index on asset.uri. |
| ``9f3c2d4a1b7e`` (head) | ``c4e7a1f9b2d0`` | ``3.4.0`` | Add start_from_trigger to trigger. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``c4e7a1f9b2d0`` | ``436dc127462c`` | ``3.4.0`` | Add index on asset.uri. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``436dc127462c`` | ``5a5d3253e946`` | ``3.4.0`` | Drop span_status column. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
Expand Down
1 change: 1 addition & 0 deletions airflow-core/newsfragments/69842.bugfix.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Keep the triggerer running when a deferred task is missing from its serialized Dag.
1 change: 1 addition & 0 deletions airflow-core/newsfragments/69843.improvement.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Avoid loading serialized Dags for ordinary deferred task triggers by persisting whether Dag context is required.
Original file line number Diff line number Diff line change
Expand Up @@ -688,6 +688,7 @@ def _create_ti_state_update_query_and_update_state(
kwargs={},
queue=ti_patch_payload.queue,
team_name=get_team_name_for_ti(task_instance_id, session),
start_from_trigger=False,
)
trigger_row.encrypted_kwargs = trigger_kwargs
session.add(trigger_row)
Expand Down
1 change: 1 addition & 0 deletions airflow-core/src/airflow/dag_processing/collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -1186,6 +1186,7 @@ def add_asset_trigger_references(
classpath=triggers[trigger_hash]["classpath"],
kwargs=triggers[trigger_hash]["kwargs"],
team_name=team_name,
start_from_trigger=False,
)
for trigger_hash in all_trigger_hashes
if trigger_hash not in orm_triggers
Expand Down
64 changes: 47 additions & 17 deletions airflow-core/src/airflow/jobs/triggerer_job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
from airflow._shared.observability.metrics import stats
from airflow._shared.timezones import timezone
from airflow.configuration import conf
from airflow.exceptions import TaskNotFound
from airflow.executors import workloads
from airflow.executors.workloads.task import TaskInstanceDTO
from airflow.jobs.base_job_runner import BaseJobRunner
Expand Down Expand Up @@ -133,6 +134,7 @@
from airflow.utils.helpers import log_filename_template_renderer, prune_dict
from airflow.utils.log.logging_mixin import LoggingMixin
from airflow.utils.session import create_session, provide_session
from airflow.utils.state import TaskInstanceState

if TYPE_CHECKING:
from opentelemetry.util._decorator import _AgnosticContextManager
Expand Down Expand Up @@ -873,36 +875,64 @@ def _create_workload(
ti=ser_ti, # type: ignore
)

workload = workloads.RunTrigger(
id=trigger.id,
classpath=trigger.classpath,
encrypted_kwargs=trigger.encrypted_kwargs,
ti=ser_ti,
timeout_after=trigger.task_instance.trigger_timeout,
)
if trigger.start_from_trigger is False:
return workload

serialized_dag_model = dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id,
session=session,
)

if serialized_dag_model:
if serialized_dag_model is None:
log.warning(
"Serialized Dag for context-required trigger was not found; skipping trigger",
trigger_id=trigger.id,
ti_id=trigger.task_instance.id,
dag_id=trigger.task_instance.dag_id,
task_id=trigger.task_instance.task_id,
dag_version_id=trigger.task_instance.dag_version_id,
)
return None

try:
task = serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
except TaskNotFound:
trigger.task_instance.state = TaskInstanceState.REMOVED
trigger.task_instance.trigger_id = None
log.warning(
"Task for deferred trigger was not found in serialized Dag; removing TaskInstance",
trigger_id=trigger.id,
ti_id=trigger.task_instance.id,
dag_id=trigger.task_instance.dag_id,
task_id=trigger.task_instance.task_id,
dag_version_id=trigger.task_instance.dag_version_id,
)
return None

if trigger.start_from_trigger is None:
trigger.start_from_trigger = task.start_from_trigger

if trigger.start_from_trigger is False:
return workload

log.info("Start from trigger enabled for task %s", task.task_id)
dag_run = trigger.task_instance.get_dagrun(session=session)

# When a TaskInstance of a Trigger contains a task with start_from_trigger enabled,
# it means we need to load the SerializedDagModel so we can build a RuntimeTaskInstance later on which
# will allow us to build a context on which we will render the templated fields.
if task.start_from_trigger:
log.info("Start from trigger enabled for task %s", task.task_id)
dag_run = trigger.task_instance.get_dagrun(session=session)

return workloads.RunTrigger(
id=trigger.id,
classpath=trigger.classpath,
encrypted_kwargs=trigger.encrypted_kwargs,
ti=ser_ti,
timeout_after=trigger.task_instance.trigger_timeout,
dag_data=serialized_dag_model.data,
dag_run_data=dag_run.dag_run_data.model_dump(exclude_unset=True),
)
return workloads.RunTrigger(
id=trigger.id,
classpath=trigger.classpath,
encrypted_kwargs=trigger.encrypted_kwargs,
ti=ser_ti,
timeout_after=trigger.task_instance.trigger_timeout,
dag_data=serialized_dag_model.data,
dag_run_data=dag_run.dag_run_data.model_dump(exclude_unset=True),
)

def fetch_trigger_details(self, trigger_ids: set[int], *, session: Session) -> dict[int, Trigger]:
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
"""
Add start_from_trigger to trigger.

Revision ID: 9f3c2d4a1b7e
Revises: c4e7a1f9b2d0
Create Date: 2026-07-13 00:00:00.000000
"""

from __future__ import annotations

import sqlalchemy as sa
from alembic import op

from airflow.migrations.utils import disable_sqlite_fkeys

# revision identifiers, used by Alembic.
revision = "9f3c2d4a1b7e"
down_revision = "c4e7a1f9b2d0"
branch_labels = None
depends_on = None
airflow_version = "3.4.0"


def upgrade():
"""Add start_from_trigger to trigger."""
with disable_sqlite_fkeys(op):
with op.batch_alter_table("trigger", schema=None) as batch_op:
batch_op.add_column(sa.Column("start_from_trigger", sa.Boolean(), nullable=True))


def downgrade():
"""Remove start_from_trigger from trigger."""
with disable_sqlite_fkeys(op):
with op.batch_alter_table("trigger", schema=None) as batch_op:
batch_op.drop_column("start_from_trigger")
1 change: 1 addition & 0 deletions airflow-core/src/airflow/models/taskinstance.py
Original file line number Diff line number Diff line change
Expand Up @@ -1784,6 +1784,7 @@ def defer_task(self, *, session: Session = NEW_SESSION) -> bool:
classpath=start_trigger_args.trigger_cls,
kwargs=trigger_kwargs,
team_name=team_name,
start_from_trigger=True,
)

# First, make the trigger entry
Expand Down
7 changes: 5 additions & 2 deletions airflow-core/src/airflow/models/trigger.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
from traceback import format_exception
from typing import TYPE_CHECKING, Any

from sqlalchemy import ForeignKey, Integer, String, Text, delete, func, or_, select, update
from sqlalchemy import Boolean, ForeignKey, Integer, String, Text, delete, func, or_, select, update
from sqlalchemy.ext.associationproxy import association_proxy
from sqlalchemy.orm import Mapped, Session, mapped_column, relationship, selectinload
from sqlalchemy.sql.functions import coalesce
Expand Down Expand Up @@ -99,6 +99,7 @@ class Trigger(Base):
created_date: Mapped[datetime.datetime] = mapped_column(UtcDateTime, nullable=False)
triggerer_id: Mapped[int | None] = mapped_column(Integer, nullable=True)
queue: Mapped[str | None] = mapped_column(String(256), nullable=True)
start_from_trigger: Mapped[bool | None] = mapped_column(Boolean, nullable=True)

# Denormalized from dag_bundle_team to keep the triggerer's ~1s polling queries join-free,
# especially since it's eventually consistent and trigger rows are ephemeral.
Expand Down Expand Up @@ -132,13 +133,15 @@ def __init__(
created_date: datetime.datetime | None = None,
queue: str | None = None,
team_name: str | None = None,
start_from_trigger: bool | None = False,
) -> None:
super().__init__()
self.classpath = classpath
self.encrypted_kwargs = self.encrypt_kwargs(kwargs)
self.created_date = created_date or timezone.utcnow()
self.queue = queue
self.team_name = team_name
self.start_from_trigger = start_from_trigger

@property
def kwargs(self) -> dict[str, Any]:
Expand Down Expand Up @@ -200,7 +203,7 @@ def rotate_fernet_key(self):
def from_object(cls, trigger: BaseTrigger) -> Trigger:
"""Alternative constructor that creates a trigger row based directly off of a Trigger object."""
classpath, kwargs = trigger.serialize()
return cls(classpath=classpath, kwargs=kwargs)
return cls(classpath=classpath, kwargs=kwargs, start_from_trigger=False)

@classmethod
@provide_session
Expand Down
2 changes: 1 addition & 1 deletion airflow-core/src/airflow/utils/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,7 @@ class MappedClassProtocol(Protocol):
"3.1.8": "509b94a1042d",
"3.2.0": "1d6611b6ab7c",
"3.3.0": "d2f4e1b3c5a7",
"3.4.0": "c4e7a1f9b2d0",
"3.4.0": "9f3c2d4a1b7e",
}

# Prefix used to identify tables holding data moved during migration.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1597,6 +1597,7 @@ def test_ti_update_state_to_deferred(
"key": "value",
"moment": datetime(2024, 12, 18, 00, 00, 1, tzinfo=timezone.utc),
}
assert t[0].start_from_trigger is False
if queues_enabled:
assert t[0].queue == "default"
else:
Expand Down
1 change: 1 addition & 0 deletions airflow-core/tests/unit/dag_processing/test_collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -394,6 +394,7 @@ def test_add_asset_trigger_references_populates_team_name(
triggers = session.scalars(select(Trigger)).all()
assert len(triggers) == 1
assert triggers[0].team_name == expected
assert triggers[0].start_from_trigger is False

@pytest.mark.usefixtures("testing_dag_bundle")
def test_add_asset_trigger_references_hash_consistency(self, dag_maker, session):
Expand Down
Loading