kaxil commented on code in PR #74222:
URL: https://github.com/apache/airflow/pull/74222#discussion_r4196073265
##########
airflow-core/src/airflow/api_fastapi/common/db/dag_runs.py:
##########
@@ -96,29 +95,17 @@ def attach_dag_versions_to_runs(dag_runs: Sequence[DagRun],
*, session: Session)
run_key_values = [(dr.dag_id, dr.run_id) for dr in runs_needing_versions]
- ti_sub = (
+ rows = session.execute(
Review Comment:
Same cause as the Gantt query: this used to union `TaskInstanceHistory`, and
the hook now limits it to the working set. When an earlier try of a run ran on
an older Dag version, that version drops out of `dag_versions` in the dag run
list and grid responses. The non-prefetched `DagRun.dag_versions` path still
picks it up through `historical_task_instances`, so the two paths now disagree.
`include_all_attempts=True` here would fix it.
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py:
##########
@@ -61,7 +60,7 @@ def get_gantt_data(
session: SessionDep,
) -> GanttResponse:
"""Get all task instance tries for Gantt chart."""
- # Exclude mapped tasks (use grid summaries) and UP_FOR_RETRY (already in
history)
+ # Pending retries retain timing for backoff; only the archived attempt
belongs on the chart.
Review Comment:
This select now goes through the `do_orm_execute` hook, so it only returns
working-set rows. After a retry the chart shows only the latest try. During
backoff the only live row is the UP_FOR_RETRY successor, which the clause below
drops, so the task disappears from the chart entirely. 3.3.2 unioned
`TaskInstanceHistory` here. Adding
`.execution_options(include_all_attempts=True)` should be enough: `archive()`
already moves unfinished states to FAILED, and the existing UP_FOR_RETRY clause
still drops the pending successor. `test_gantt.py` has no retried-task case,
which would have caught this. Fine as a follow-up.
##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_state_store.py:
##########
@@ -47,7 +47,7 @@
def _get_task_scope_for_ti(task_instance_id: UUID, session: Session) ->
TaskScope:
ti = session.get(TI, task_instance_id)
- if ti is None:
+ if ti is None or ti.working_set is not True:
Review Comment:
A GET from an archived attempt gets 404 here, and the SDK maps any 404 to
`TASK_STORE_NOT_FOUND`, so `task_state_store.get()` hands back the default.
`ResumableJobMixin` reads that as no saved job id. If the attempt was archived
while its worker was still running (heartbeat timeout through `handle_failure`,
or adopt-or-reset), it submits the job again. At 2026-10-30 the `set` after
that gets 410, so the new job id is never stored. Before this change the read
went to the shared (dag_id, run_id, task_id, map_index) scope and returned the
stored id. 3.3 also returned 404 here because the old UUID was gone. `get_xcom`
already returns 410 for an archived attempt when
`IdentifyArchivedTaskStateUpdates.is_applied`. Doing that here too would make
the SDK raise. The `Airflow-API-Version` header in
`test_archived_attempt_returns_404` has no effect because the `client` fixture
stubs out `require_auth`. The real-auth setup from
`test_archived_attempt_read_response_by_version` would cover both
versions.
##########
airflow-core/src/airflow/migrations/versions/0142_3_4_0_unify_task_attempt_ownership.py:
##########
@@ -0,0 +1,484 @@
+#
+# 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.
+
+"""
+Unify task attempt ownership without rewriting legacy XCom data.
Review Comment:
This drops `task_instance_history` and `hitl_detail_history` and renames
`xcom` to `xcom_v1` and `rendered_task_instance_fields` to `rtif_v1`, but the
PR has no newsfragment. Anyone with SQL, BI dashboards or `db clean --tables`
scripts against those names only finds out at upgrade time. A `significant`
newsfragment listing the old and new names would cover it.
##########
airflow-core/src/airflow/models/renderedtifields.py:
##########
@@ -335,15 +353,83 @@ def _do_delete_old_records(
run_ids_to_keep: list[str] | ScalarSelect[str],
session: Session,
) -> None:
- # This query might deadlock occasionally and it should be retried if
fails (see decorator)
- stmt = (
+ from airflow.models.taskinstance import TaskInstance
+
+ legacy = LegacyRenderedTaskInstanceFields.__table__
+ session.execute(
+ delete(legacy).where(
+ legacy.c.dag_id == dag_id,
+ legacy.c.task_id == task_id,
+ legacy.c.run_id.not_in(run_ids_to_keep),
+ )
+ )
+ session.execute(
delete(cls)
Review Comment:
This prune runs on every rendered-fields write
(`num_dag_runs_to_retain_rendered_fields` defaults to 30). `rtif_v2` has no
dag_id/task_id/run_id columns, so the task scope only comes from the EXISTS
into `task_instance`, and that matches every attempt of every run of this task,
not only the retained window. The old delete was a range on the RTIF primary
key. Since 0142 is unreleased, bounding the subquery to the runs that just left
the window (or putting the coordinates on `rtif_v2`) would keep the per-write
cost flat as history grows. `synchronize_session` also went from `False` to
`"fetch"`, which adds a pre-select on MySQL with nothing in the session to sync.
##########
airflow-core/src/airflow/models/xcom.py:
##########
@@ -86,35 +99,123 @@ class XComModel(TaskInstanceDependencies):
# separately, and enforce uniqueness with DagRun.id instead.
Index("idx_xcom_key", key),
Index("idx_xcom_task_instance", dag_id, task_id, run_id, map_index),
- CheckConstraint(mapped_length >= 0, name="mapped_length_not_negative"),
+ CheckConstraint(mapped_length >= 0,
name=conv("ck_xcom_mapped_length_not_negative")),
PrimaryKeyConstraint("dag_run_id", "task_id", "map_index", "key",
name="xcom_pkey"),
ForeignKeyConstraint(
[dag_id, task_id, run_id, map_index],
[
- "task_instance.dag_id",
- "task_instance.task_id",
- "task_instance.run_id",
- "task_instance.map_index",
+ "legacy_task_data_owner.dag_id",
+ "legacy_task_data_owner.task_id",
+ "legacy_task_data_owner.run_id",
+ "legacy_task_data_owner.map_index",
],
name="xcom_task_instance_fkey",
ondelete="CASCADE",
),
)
- dag_run = relationship(
- "DagRun",
- primaryjoin="XComModel.dag_run_id == foreign(DagRun.id)",
- uselist=False,
- lazy="joined",
- passive_deletes="all",
- )
- logical_date = association_proxy("dag_run", "logical_date")
- task = relationship(
- "TaskInstance",
- viewonly=True,
- lazy="raise",
+class XComModelV2(Base):
+ """XCom values stored per task attempt, keyed by the attempt UUID."""
+
+ __tablename__ = "xcom_v2"
+
+ id: Mapped[UUID] = mapped_column(Uuid(), primary_key=True,
default=uuid6.uuid7)
+ task_instance_id: Mapped[UUID] = mapped_column(Uuid(), nullable=False)
+ key: Mapped[str] = mapped_column(String(512, **COLLATION_ARGS),
nullable=False)
+ value: Mapped[Any] = mapped_column(JSON().with_variant(postgresql.JSONB,
"postgresql"), nullable=True)
+ timestamp: Mapped[datetime] = mapped_column(UtcDateTime,
default=timezone.utcnow, nullable=False)
+ dag_result: Mapped[bool | None] = mapped_column(Boolean, nullable=True,
default=False)
+ mapped_length: Mapped[int | None] = mapped_column(Integer, nullable=True)
+
+ __table_args__ = (
+ PrimaryKeyConstraint("id", name="xcom_v2_pkey"),
+ UniqueConstraint("task_instance_id", "key", name="xcom_v2_ti_key_uq"),
Review Comment:
`xcom_v1` keeps `idx_xcom_key`, but `xcom_v2` has no index on `key` alone.
The global XCom list (`/dags/~/dagRuns/~/taskInstances/~/xcomEntries` with a
key or key-prefix filter, which the Browse > XComs page sends) used that index,
and every row written after the upgrade lands here. Adding the same index to
`xcom_v2` in 0142 would keep that search on an index.
##########
airflow-core/src/airflow/models/taskinstance.py:
##########
@@ -639,7 +701,9 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload):
duration: Mapped[float | None] = mapped_column(Float, nullable=True)
state: Mapped[str | None] = mapped_column(String(20), nullable=True)
try_number: Mapped[int] = mapped_column(Integer, default=0)
- max_tries: Mapped[int] = mapped_column(Integer, server_default="-1")
+ max_tries: Mapped[int] = mapped_column(Integer, server_default="-1",
nullable=False)
+ working_set: Mapped[bool | None] = mapped_column(Boolean, default=True,
server_default=true())
+ archived_reason: Mapped[str | None] = mapped_column(String(50))
Review Comment:
`prepare_db_for_next_try` is the only caller of `archive()` and always
passes `reason="retry"`, but it also runs for a user clear, restart completion,
the orphan reset in the scheduler and the restore of a removed task. So every
archived row says "retry" (or "legacy" from the migration), and nothing reads
the column yet. Could `prepare_db_for_next_try` take the reason (as a
`Literal`) so the value means something before AIP-111 or the UI starts relying
on it?
--
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]