This is an automated email from the ASF dual-hosted git repository.
amoghrajesh pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 5e785acfb3b Discard the task state values for a task instance when it
is cleared (#72100)
5e785acfb3b is described below
commit 5e785acfb3b206ed2a24dcb063298c7fd47fe465
Author: Amogh Desai <[email protected]>
AuthorDate: Wed Sep 30 15:24:24 2026 +0530
Discard the task state values for a task instance when it is cleared
(#72100)
---
.../task-and-asset-state-store.rst | 4 +
airflow-core/docs/core-concepts/dag-run.rst | 11 +-
.../docs/core-concepts/resumable-tasks.rst | 75 ++++++++----
.../docs/core-concepts/task-state-store.rst | 6 +-
airflow-core/newsfragments/72100.significant.rst | 59 ++++++++++
.../core_api/datamodels/task_instances.py | 5 +
.../core_api/openapi/v2-rest-api-generated.yaml | 7 ++
.../core_api/routes/public/task_instances.py | 8 ++
.../core_api/services/public/task_instances.py | 41 +++++++
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 6 +
.../airflow/ui/openapi-gen/requests/types.gen.ts | 4 +
.../airflow/ui/public/i18n/locales/en/common.json | 4 +
.../airflow/ui/public/i18n/locales/en/dags.json | 1 +
.../TaskInstance/ClearGroupTaskInstanceDialog.tsx | 40 +++++--
.../TaskInstance/ClearTaskInstanceDialog.test.tsx | 127 +++++++++++++++++++++
.../Clear/TaskInstance/ClearTaskInstanceDialog.tsx | 23 +++-
.../src/airflow/ui/src/constants/localStorage.ts | 1 +
.../airflow/ui/src/hooks/useUserSettings.test.tsx | 21 ++++
.../src/airflow/ui/src/hooks/useUserSettings.ts | 4 +
.../ui/src/pages/Settings/Settings.test.tsx | 8 ++
.../src/airflow/ui/src/pages/Settings/Settings.tsx | 13 +++
.../TaskInstances/BulkClearTaskInstancesButton.tsx | 34 ++++--
.../ui/src/queries/useBulkClearTaskInstances.ts | 2 +
.../core_api/routes/public/test_task_instances.py | 117 +++++++++++++++++++
.../src/airflowctl/api/datamodels/generated.py | 7 ++
airflow-ctl/src/airflowctl/api/operations.py | 2 +-
.../tests/airflow_ctl/api/test_operations.py | 19 ++-
27 files changed, 605 insertions(+), 44 deletions(-)
diff --git
a/airflow-core/docs/administration-and-deployment/task-and-asset-state-store.rst
b/airflow-core/docs/administration-and-deployment/task-and-asset-state-store.rst
index ba26194d4ec..d4a68591294 100644
---
a/airflow-core/docs/administration-and-deployment/task-and-asset-state-store.rst
+++
b/airflow-core/docs/administration-and-deployment/task-and-asset-state-store.rst
@@ -60,6 +60,10 @@ Number of days after which task state store rows expire.
When a key is written w
``clear_on_success``
~~~~~~~~~~~~~~~~~~~~
+Retention and ``clear_on_success`` are not the only ways task state store
entries get removed:
+clearing a task instance through the REST API, the UI, or ``airflowctl`` also
discards them by
+default unless ``keep_task_state`` is set. See
:doc:`/core-concepts/resumable-tasks`.
+
When ``True``, all task state store keys for a task instance are automatically
deleted when that task instance moves to the ``success`` state. Defaults to
``False``, which preserves task state store entries after success for
observability (e.g. the submitted job ID or the last row count is still
readable from the UI or REST API after the run completes).
.. important::
diff --git a/airflow-core/docs/core-concepts/dag-run.rst
b/airflow-core/docs/core-concepts/dag-run.rst
index dad6eb51147..9db308ddd3c 100644
--- a/airflow-core/docs/core-concepts/dag-run.rst
+++ b/airflow-core/docs/core-concepts/dag-run.rst
@@ -250,8 +250,17 @@ There are multiple options you can select to re-run -
* **Downstream** - The downstream tasks in the current Dag
* **Recursive** - All the tasks in the child Dags and parent Dags
* **Failed** - Only the failed tasks in the Dag's most recent run
+* **Keep task state and resume** - Preserve the task's ``task_state_store``
entries instead of
+ discarding them, so the next attempt resumes from a checkpoint or reconnects
to an external job
+ instead of starting over. See :doc:`resumable-tasks` for details.
-You can also clear the task through CLI using the command:
+Clearing an individual task instance discards its task state store entries by
default, so the next
+attempt starts from the beginning instead of resuming a checkpoint or
reconnecting to an external
+job. See :doc:`resumable-tasks` for when to keep it instead.
+
+You can also clear the task through CLI using ``airflowctl tasks clear``,
which discards task state
+the same way and accepts ``--keep-task-state`` to opt back in. The legacy,
deprecated
+``airflow tasks clear`` command below always keeps task state:
.. code-block:: bash
diff --git a/airflow-core/docs/core-concepts/resumable-tasks.rst
b/airflow-core/docs/core-concepts/resumable-tasks.rst
index 2fca68b978a..69d267a955b 100644
--- a/airflow-core/docs/core-concepts/resumable-tasks.rst
+++ b/airflow-core/docs/core-concepts/resumable-tasks.rst
@@ -145,26 +145,61 @@ existing job on retry instead of submitting a new one.
For more details and a working example, see
:class:`~airflow.sdk.ResumableJobMixin`.
-**Clearing a task is treated the same as a retry**
-
-Clearing a task instance does not delete its ``task_state_store`` rows -- they
are only removed
-when the ``dag_run`` itself is deleted, or by :ref:`airflow state-store clean
-<task-and-asset-state-store-cleanup>`. For a checkpointed task this is usually
what you want:
-clearing resumes from the last checkpoint rather than starting over.
-
-For an operator with durable execution, it means clearing a task whose
external job already
-succeeded reads that stored result back and returns immediately, without
resubmitting the job. If
-you want clearing to always resubmit regardless of a prior success, set
-``[state_store] clear_on_success = True``, which deletes a task's state store
rows automatically
-when it moves to ``SUCCESS`` (see
:doc:`/administration-and-deployment/task-and-asset-state-store`).
-
-This does not guarantee the external job is still there to reconnect to,
though. Clearing a task
-that is actively running (``deferrable=False``) stops the worker process,
which runs the
-operator's ``on_kill``. Most operators with durable execution cancel the
external job there by
-default, so the next attempt finds it already stopped instead of still running
-- an operator that
-leaves the job running by default on kill is the exception, check its own
docs. Deferred tasks
-(``deferrable=True``) don't have this problem: there is no actively polling
worker process for the
-clear to interrupt.
+**Retries resume, clearing starts over**
+
+A retry keeps the task's ``task_state_store`` entries, which is what makes
crash recovery work: the
+next attempt reads the checkpoint or the external job id written by the
attempt before it.
+
+Clearing a task discards them. Clearing means "run this again", and a
checkpoint records how far a
+task got, not what it got there with. If you fixed the code or the upstream
data and cleared the
+task, resuming would leave the work done before the fix in place and silently
mix it with the
+corrected work. So by default a cleared task starts from the beginning.
+
+To resume from the checkpoint instead, set ``keep_task_state`` when clearing,
or tick the
+corresponding box in the clear dialog. That is the right choice when nothing
about the inputs or the
+code changed and you only want the task to carry on where it stopped.
+
+This applies to clearing individual task instances through the REST API, the
UI, or
+``airflowctl dags clear`` (which clears every task instance in the matched Dag
run(s) through this
+same endpoint). Clearing an entire Dag run through the Clear Run dialog/API,
marking a task as
+failed or success (which clears downstream tasks as a side effect), and the
core CLI's
+``airflow tasks clear`` / ``airflow dags clear`` (which call
``clear_task_instances()`` directly and
+never reach the discard endpoint) all still keep task state unconditionally
today; see
+`#72929 <https://github.com/apache/airflow/issues/72929>`_.
+
+**Clearing a task that submitted an external job**
+
+For an operator with durable execution the stored value is an external job id,
so discarding it has
+a different consequence: the next attempt submits a new job rather than
reconnecting to the existing
+one.
+
+Whether that matters depends on what happened to the job:
+
+* Most operators cancel the external job in ``on_kill``, so clearing a
*running* task stops the job
+ and there is nothing left to reconnect to. Submitting a fresh one is the
only option anyway.
+* An operator configured, or defaulting, to leave the job running on kill
keeps it alive, so a
+ fresh submission runs alongside it. Check the operator's own docs — for
example
+ ``GlueJobOperator`` defaults ``stop_job_run_on_kill`` to ``False`` and so
leaves the job running
+ unless you turn it on.
+* Clearing a *failed* task never runs ``on_kill`` at all, so an external job
that outlived the
+ worker is still running.
+* A deferred task has no worker process to run ``on_kill`` on. Instead, the
Triggerer cancels the
+ orphaned trigger and runs the trigger's ``on_kill``, bounded by
``[triggerer] on_kill_timeout``.
+ Most triggers cancel the external job there too.
+ ``GlueJobCompleteTrigger`` and ``LivyTrigger`` don't implement ``on_kill``,
but neither writes a
+ job id to the state store when deferred either, so ``keep_task_state`` will
not help for them. For
+ a deferred ``GlueJobOperator``, ``durable`` defaults to ``True``, so it
reattaches to a still-running
+ Glue run by default and clearing it does not start over; stop the Glue run
before clearing, or set
+ ``durable=False``, to force a fresh one. For a deferred ``LivyOperator``,
cancel the batch yourself
+ before clearing.
+
+In the cases where the job is left running (the second bullet or a failed
task), pass
+``keep_task_state`` so the next attempt reconnects to the job already in
flight instead of paying
+for a second one.
+
+Note that ``[state_store] clear_on_success`` is a separate control: it
discards a task's entries as
+soon as it reaches ``SUCCESS``, so nothing is left for a later clear to find
either way (see
+:doc:`/administration-and-deployment/task-and-asset-state-store`).
.. _concepts-resumable-tasks-async:
diff --git a/airflow-core/docs/core-concepts/task-state-store.rst
b/airflow-core/docs/core-concepts/task-state-store.rst
index 46bd2c88f4c..c972fd7861c 100644
--- a/airflow-core/docs/core-concepts/task-state-store.rst
+++ b/airflow-core/docs/core-concepts/task-state-store.rst
@@ -285,7 +285,7 @@ If the worker process crashes, the task instance is
retried. Task store data wri
Deferrable tasks
~~~~~~~~~~~~~~~~
-Once a task defers, the Triggerer handles continuity across poke cycles. Use
task state store in deferrable tasks only when you need to survive an
operator-initiated clear, not for normal poke continuity.
+Once a task defers, the Triggerer handles continuity across poke cycles.
Clearing a deferred task does not synchronously cancel its trigger: the
Triggerer only notices the trigger is orphaned on its next iteration, then
cancels it via the trigger's ``on_kill``, bounded by ``[triggerer]
on_kill_timeout``. A new attempt can therefore start before that cancellation
finishes. Most triggers implement ``on_kill`` to cancel the external job there,
so the next attempt usually finds nothing left [...]
Mapped tasks
@@ -304,6 +304,10 @@ To wipe state across all map indices of a task, use the
:doc:`Core API </adminis
Automatic cleanup (``clear_on_success``)
----------------------------------------
+Task state store entries are also removed when a task instance is cleared
through the REST API, the
+UI, or ``airflowctl``: clearing discards them by default so the next attempt
starts over, unless
+``keep_task_state`` is set. See :doc:`resumable-tasks` for that behaviour.
+
When ``[state_store] clear_on_success = True``, all task state store keys for
a task instance are automatically deleted when the task moves to the
``success`` state. This is useful for reducing storage when post-success
observability is not needed.
.. note::
diff --git a/airflow-core/newsfragments/72100.significant.rst
b/airflow-core/newsfragments/72100.significant.rst
new file mode 100644
index 00000000000..e0c94bce297
--- /dev/null
+++ b/airflow-core/newsfragments/72100.significant.rst
@@ -0,0 +1,59 @@
+Clearing a task now discards its task state store entries by default
+
+Clearing a task instance discards its ``task_state_store`` entries, so the
next attempt starts from
+the beginning instead of resuming from a checkpoint or reconnecting to an
external job recorded by
+the attempt that was cleared.
+
+Retries are unaffected. They keep task state exactly as before, which is what
crash recovery relies
+on. Only a deliberate clear discards.
+
+This only applies to clearing individual task instances (the task-instance
clear endpoint / dialog,
+and ``airflowctl dags clear``, which clears every task instance in the matched
Dag run(s) through
+the same endpoint). Clearing an entire Dag run through the "Clear Run"
dialog/API, and marking a
+task as failed or success (which clears downstream tasks as a side effect),
still keep task state
+unconditionally today; extending discard-by-default to those paths is tracked
in
+`#72929 <https://github.com/apache/airflow/issues/72929>`_.
+
+**Why**
+
+Clearing means "run this again". A checkpoint records how far a task got, not
what it got there
+with, so resuming after the code or the upstream data changed left work done
before the fix in place
+and silently mixed it with the corrected work. Clearing a task whose external
job had already
+succeeded was worse: the operator read the stored result back and returned in
seconds having run
+nothing.
+
+**Keeping the old behaviour**
+
+Pass ``keep_task_state=True`` to the clear task instances endpoint, or tick
"keep task state" in the
+clear dialog. Use it when nothing about the inputs or the code changed and the
task should carry on
+where it stopped, or when an external job is still running and you want the
next attempt to
+reconnect rather than submit a duplicate. To make "keep task state" the
default so you don't have to
+tick it on every clear, turn it on under Settings > Clearing > "Keep task
state on clear"; that
+default is per-browser and does not affect anyone else on the same deployment.
+
+Operators with durable execution are worth particular attention. Clearing a
*failed* task never runs
+``on_kill``, so an external job that outlived its worker is still running, and
discarding the stored
+id means submitting a second one. The same applies to operators configured to
leave their job alive
+on kill, such as ``GlueJobOperator``, which defaults ``stop_job_run_on_kill``
to ``False``.
+
+This also widens the permission surface of clearing: a task-instance clear
(gated by the
+task-instance edit permission) now deletes ``task_state_store`` rows for that
task instance, the same
+rows the dedicated task-state-store endpoint gates behind the delete
permission.
+
+On the CLI, ``airflowctl tasks clear`` discards task state the same way and
has a
+``--keep-task-state`` flag to opt back in. ``airflowctl dags clear`` discards
task state too, since
+it clears every task instance in the matched Dag run(s) through the same
endpoint, but it has no
+``--keep-task-state`` equivalent yet. ``airflow dags clear`` and ``airflow
tasks clear``
+(airflow-core; only ``tasks clear`` is deprecated, in favor of ``airflowctl
tasks clear``) are
+unaffected: they call ``clear_task_instances()`` directly and never touch task
state.
+
+* Types of change
+
+ * [ ] Dag changes
+ * [ ] Config changes
+ * [x] API changes
+ * [x] CLI changes
+ * [x] Behaviour changes
+ * [ ] Plugin changes
+ * [ ] Dependency changes
+ * [ ] Code interface changes
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
index a4d65063b60..1e7564afe48 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
@@ -253,6 +253,11 @@ class ClearTaskInstancesBody(StrictBaseModel):
"and finally ``False`` (the historical default for clear/rerun).",
)
prevent_running_task: bool = False
+ keep_task_state: bool = Field(
+ default=False,
+ description="Keep the task state store entries of the cleared task
instances so the next "
+ "attempt resumes from them. By default they are discarded, so the task
starts over.",
+ )
note: Annotated[str, StringConstraints(max_length=1000)] | None = None
@model_validator(mode="before")
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index 1ca26d0f9c2..60841257d6e 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -13028,6 +13028,13 @@ components:
type: boolean
title: Prevent Running Task
default: false
+ keep_task_state:
+ type: boolean
+ title: Keep Task State
+ description: Keep the task state store entries of the cleared task
instances
+ so the next attempt resumes from them. By default they are
discarded,
+ so the task starts over.
+ default: false
note:
anyOf:
- type: string
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
index 803f5edc0d8..65f8b57ec8e 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
@@ -107,6 +107,7 @@ from airflow.api_fastapi.core_api.openapi.exceptions import
create_openapi_http_
from airflow.api_fastapi.core_api.security import GetUserDep,
ReadableTIFilterDep, requires_access_dag
from airflow.api_fastapi.core_api.services.public.task_instances import (
BulkTaskInstanceService,
+ _discard_task_state_store,
_get_task_group_task_ids,
_get_task_group_task_instances,
_patch_task_group_state,
@@ -992,6 +993,13 @@ def post_clear_task_instances(
except AirflowClearRunningTaskException as e:
raise HTTPException(status.HTTP_409_CONFLICT, str(e)) from e
+ # After the clear has succeeded, so a failed clear cannot take the
task state with it.
+ # This is the only clear path that discards task state today; Dag-run
clear and
+ # mark-as-failed/success (which clear downstream tasks) still keep it
unconditionally.
+ # It is tracked through https://github.com/apache/airflow/issues/72929
+ if not body.keep_task_state:
+ _discard_task_state_store(task_instances, session,
event="Discarded task state on clear")
+
if body.note is not None:
_patch_task_instance_note(
task_instance_body=body,
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
index 476ce7bb7f2..1c9a92768c0 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
@@ -59,6 +59,47 @@ from airflow.utils.state import TaskInstanceState
log = structlog.get_logger(__name__)
+def _discard_task_state_store(tis: Sequence[TI], session: Session, *, event:
str) -> None:
+ """
+ Discard the task state store entries of each task instance.
+
+ A failure is logged and re-raised so the request fails and the session
rolls back, rather than
+ reporting success while some entries survive undiscarded.
+
+ This only drops the metadata DB reference row via ``_get_db_backend()``;
it does not go through
+ ``get_state_backend()``. A custom ``[workers] state_store_backend``
payload is left orphaned
+ with no reclaim path other than its own lifecycle/TTL policy, and a custom
``[state_store]
+ backend`` is not touched at all: the worker still reads and writes there,
so a clear reports
+ success while the actual state survives and a later attempt can resume
from it. Closing this
+ gap needs a server-side path to the configured state backend and is
tracked for a future
+ change; today, this discard is only exact for the default metastore
backend.
+
+ :param event: what prompted the discard, used as the log event name.
+ """
+ backend = _get_db_backend()
+ for ti in tis:
+ scope = TaskScope(
+ dag_id=ti.dag_id,
+ run_id=ti.run_id,
+ task_id=ti.task_id,
+ map_index=ti.map_index if ti.map_index is not None else -1,
+ )
+ try:
+ backend.clear(scope=scope, session=session)
+ except Exception:
+ log.warning(
+ "Failed to discard task state",
+ discard_event=event,
+ dag_id=ti.dag_id,
+ run_id=ti.run_id,
+ task_id=ti.task_id,
+ map_index=ti.map_index,
+ exc_info=True,
+ )
+ raise
+ log.info(event, task_instance_count=len(tis))
+
+
def _clear_task_state_store_on_success(tis: Sequence[TI], session: Session) ->
None:
"""Clear task state store rows for each TI if clear_on_success is
enabled."""
if not conf.getboolean("state_store", "clear_on_success", fallback=False):
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 8ec5c8c5116..729d8a5ad80 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -2446,6 +2446,12 @@ export const $ClearTaskInstancesBody = {
title: 'Prevent Running Task',
default: false
},
+ keep_task_state: {
+ type: 'boolean',
+ title: 'Keep Task State',
+ description: 'Keep the task state store entries of the cleared
task instances so the next attempt resumes from them. By default they are
discarded, so the task starts over.',
+ default: false
+ },
note: {
anyOf: [
{
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 189c9f182f7..4017f33daee 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -770,6 +770,10 @@ export type ClearTaskInstancesBody = {
*/
run_on_latest_version?: boolean | null;
prevent_running_task?: boolean;
+ /**
+ * Keep the task state store entries of the cleared task instances so the
next attempt resumes from them. By default they are discarded, so the task
starts over.
+ */
+ keep_task_state?: boolean;
note?: string | null;
};
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
index 7a781ef45ed..d3408928095 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
@@ -271,6 +271,10 @@
"selected": "Selected",
"settings": {
"clearing": {
+ "keepTaskState": {
+ "helper": "Keep task state store entries by default when clearing task
instances, instead of discarding them.",
+ "label": "Keep task state on clear"
+ },
"preventRunningTask": {
"helper": "Skip tasks that are currently running when clearing task
instances.",
"label": "Prevent clearing running tasks"
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
index d77b200dbb6..704ecc41118 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
@@ -79,6 +79,7 @@
"downstream": "Downstream",
"existingTasks": "Clear existing tasks",
"future": "Future",
+ "keepTaskState": "Keep task state and resume",
"onlyFailed": "Clear only failed tasks",
"past": "Past",
"preventRunningTasks": "Prevent rerun if task is running",
diff --git
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
index 0658d9d789f..7f6c6fad8ee 100644
---
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
+++
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
@@ -16,7 +16,7 @@
* specific language governing permissions and limitations
* under the License.
*/
-import { useState } from "react";
+import { useEffect, useState } from "react";
import { Button, Flex } from "@chakra-ui/react";
import { useTranslation } from "react-i18next";
@@ -31,7 +31,7 @@ import { Checkbox, Modal, SegmentedControl } from
"src/system-components";
import { ActionAccordion } from "src/components/ActionAccordion";
import { useRerunWithLatestVersion } from
"src/components/Clear/useRerunWithLatestVersion";
-import { useClearTaskInstanceDefaultOptions } from "src/hooks/useUserSettings";
+import { useClearKeepTaskStateDefault, useClearTaskInstanceDefaultOptions }
from "src/hooks/useUserSettings";
import { useClearTaskInstances } from "src/queries/useClearTaskInstances";
import { useClearTaskInstancesDryRun } from
"src/queries/useClearTaskInstancesDryRun";
import { isStatePending, useAutoRefresh } from "src/utils";
@@ -49,13 +49,8 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
const { dagId = "", runId = "" } = useParams();
const groupId = taskInstance.task_id;
- const { isPending, mutate } = useClearTaskInstances({
- dagId,
- dagRunId: runId,
- onSuccessConfirm: onClose,
- });
-
const [clearTaskInstanceDefaultOptions] =
useClearTaskInstanceDefaultOptions();
+ const [keepTaskStateDefault] = useClearKeepTaskStateDefault();
const [selectedOptions, setSelectedOptions] =
useState<Array<string>>(clearTaskInstanceDefaultOptions);
const onlyFailed = selectedOptions.includes("onlyFailed");
@@ -63,8 +58,28 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
const future = selectedOptions.includes("future");
const upstream = selectedOptions.includes("upstream");
const downstream = selectedOptions.includes("downstream");
+ const [keepTaskState, setKeepTaskState] = useState(keepTaskStateDefault);
const [note, setNote] = useState<string | null>(null);
+ const onCloseDialog = () => {
+ setNote(null);
+ setKeepTaskState(keepTaskStateDefault);
+ onClose();
+ };
+
+ useEffect(() => {
+ if (open) {
+ setNote(null);
+ setKeepTaskState(keepTaskStateDefault);
+ }
+ }, [open, keepTaskStateDefault]);
+
+ const { isPending, mutate } = useClearTaskInstances({
+ dagId,
+ dagRunId: runId,
+ onSuccessConfirm: onCloseDialog,
+ });
+
const { data: dagDetails } = useDagServiceGetDagDetails({
dagId,
});
@@ -137,6 +152,7 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
include_past: past,
include_upstream: upstream,
...(note === null ? {} : { note }),
+ ...(keepTaskState ? { keep_task_state: true } : {}),
only_failed: onlyFailed,
run_on_latest_version: runOnLatestVersion,
task_group_id: groupId,
@@ -146,6 +162,12 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
>
<CgRedo /> {translate("modal.confirm")}
</Button>
+ <Checkbox
+ checked={keepTaskState}
+ onCheckedChange={(event) =>
setKeepTaskState(Boolean(event.checked))}
+ >
+ {translate("dags:runAndTaskActions.options.keepTaskState")}
+ </Checkbox>
{shouldShowRunOnLatestOption ? (
<Checkbox
checked={runOnLatestVersionForced || runOnLatestVersion}
@@ -163,7 +185,7 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
</>
}
lazyMount
- onOpenChange={onClose}
+ onOpenChange={onCloseDialog}
open={open}
title={
<>
diff --git
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.test.tsx
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.test.tsx
new file mode 100644
index 00000000000..def2c5e5469
--- /dev/null
+++
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.test.tsx
@@ -0,0 +1,127 @@
+/*!
+ * 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.
+ */
+import "@testing-library/jest-dom";
+import { fireEvent, render, screen, waitFor } from "@testing-library/react";
+import { http, HttpResponse } from "msw";
+import { setupServer, type SetupServer } from "msw/node";
+import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from
"vitest";
+
+import type { TaskInstanceResponse } from "openapi/requests/types.gen";
+
+import { CLEAR_KEEP_TASK_STATE_KEY } from "src/constants/localStorage";
+import { handlers } from "src/mocks/handlers";
+import { Wrapper } from "src/utils/Wrapper";
+
+import ClearTaskInstanceDialog from "./ClearTaskInstanceDialog";
+
+const DAG_ID = "test_dag";
+const DAG_RUN_ID = "run_1";
+const TASK_ID = "task_1";
+
+const taskInstance: TaskInstanceResponse = {
+ dag_display_name: "Test DAG",
+ dag_id: DAG_ID,
+ dag_run_id: DAG_RUN_ID,
+ dag_version: null,
+ duration: null,
+ end_date: null,
+ executor: null,
+ executor_config: "{}",
+ hostname: null,
+ id: "test_task_instance",
+ logical_date: "2025-01-01T00:00:00Z",
+ map_index: -1,
+ max_tries: 0,
+ note: null,
+ operator: "EmptyOperator",
+ operator_name: "EmptyOperator",
+ pid: null,
+ pool: "default_pool",
+ pool_slots: 1,
+ priority_weight: null,
+ queue: null,
+ queued_when: null,
+ rendered_fields: undefined,
+ rendered_map_index: null,
+ run_after: "2025-01-01T00:00:00Z",
+ scheduled_when: null,
+ start_date: null,
+ state: "success",
+ task_display_name: "task_1",
+ task_id: TASK_ID,
+ trigger: null,
+ triggerer_job: null,
+ try_number: 1,
+ unixname: null,
+};
+
+const affectedTasks = { task_instances: [taskInstance], total_entries: 1 };
+
+const mutateMock = vi.fn();
+
+vi.mock("src/queries/useClearTaskInstances", () => ({
+ useClearTaskInstances: () => ({ isPending: false, mutate: mutateMock }),
+}));
+
+// Mocked directly (rather than via MSW) because it backs both the
affected-tasks
+// list in the main dialog and the confirmation dialog's running-task gate; a
+// same-tick response for both keeps the gate from blocking on a real fetch.
+vi.mock("src/queries/useClearTaskInstancesDryRun", () => ({
+ useClearTaskInstancesDryRun: () => ({ data: affectedTasks, isFetching:
false, isPending: false }),
+}));
+
+let server: SetupServer;
+
+beforeAll(() => {
+ server = setupServer(
+ ...handlers,
+ http.get(`/api/v2/dags/${DAG_ID}/details`, () => HttpResponse.json({})),
+ http.get(`/api/v2/dags/${DAG_ID}/dagRuns/${DAG_RUN_ID}`, () =>
HttpResponse.json({ dag_versions: [] })),
+ );
+ server.listen({ onUnhandledRequest: "bypass" });
+});
+afterEach(() => {
+ server.resetHandlers();
+ localStorage.clear();
+});
+afterAll(() => server.close());
+
+describe("ClearTaskInstanceDialog", () => {
+ it("seeds the keep-task-state checkbox from the stored default and sends it
on confirm", async () => {
+ localStorage.setItem(CLEAR_KEEP_TASK_STATE_KEY, JSON.stringify(true));
+
+ render(<ClearTaskInstanceDialog onClose={vi.fn()} open
taskInstance={taskInstance} />, {
+ wrapper: Wrapper,
+ });
+
+ const keepTaskStateCheckbox = screen.getByRole("checkbox", { name:
/keepTaskState/iu });
+
+ expect(keepTaskStateCheckbox).toBeChecked();
+
+ const confirmButton = await screen.findByRole("button", { name:
/modal\.confirm/iu });
+
+ fireEvent.click(confirmButton);
+
+ await waitFor(() => expect(mutateMock).toHaveBeenCalled());
+
+ const [{ requestBody }] = mutateMock.mock.calls[0] as [{ requestBody: {
keep_task_state?: boolean } }];
+
+ expect(requestBody.keep_task_state).toBe(true);
+ });
+});
diff --git
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
index c1d18f46a86..8d04075adb9 100644
---
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
+++
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
@@ -33,6 +33,7 @@ import { useRerunWithLatestVersion } from
"src/components/Clear/useRerunWithLate
import Time from "src/components/Time";
import {
+ useClearKeepTaskStateDefault,
useClearPreventRunningTaskDefault,
useClearTaskInstanceDefaultOptions,
} from "src/hooks/useUserSettings";
@@ -86,6 +87,7 @@ const ClearTaskInstanceDialog = (props: Props) => {
const [clearTaskInstanceDefaultOptions] =
useClearTaskInstanceDefaultOptions();
const [preventRunningTaskDefault] = useClearPreventRunningTaskDefault();
+ const [keepTaskStateDefault] = useClearKeepTaskStateDefault();
const [selectedOptions, setSelectedOptions] =
useState<Array<string>>(clearTaskInstanceDefaultOptions);
const onlyFailed = selectedOptions.includes("onlyFailed");
@@ -93,6 +95,7 @@ const ClearTaskInstanceDialog = (props: Props) => {
const future = selectedOptions.includes("future");
const upstream = selectedOptions.includes("upstream");
const downstream = selectedOptions.includes("downstream");
+ const [keepTaskState, setKeepTaskState] = useState(keepTaskStateDefault);
const [preventRunningTask, setPreventRunningTask] =
useState(preventRunningTaskDefault);
const [note, setNote] = useState<string | null>(taskInstance?.note ?? null);
@@ -103,8 +106,18 @@ const ClearTaskInstanceDialog = (props: Props) => {
}
}, [openDialog, taskInstance?.note]);
+ // Separate from the note effect above: this must only reset on open, not on
every
+ // note refetch, or a background note change while the dialog is open
silently
+ // discards the user's checked box.
+ useEffect(() => {
+ if (openDialog) {
+ setKeepTaskState(keepTaskStateDefault);
+ }
+ }, [openDialog, keepTaskStateDefault]);
+
const onCloseDialog = () => {
setNote(taskInstance?.note ?? null);
+ setKeepTaskState(keepTaskStateDefault);
closeDialog();
};
@@ -211,10 +224,16 @@ const ClearTaskInstanceDialog = (props: Props) => {
>
<CgRedo /> {translate("modal.confirm")}
</Button>
+ <Checkbox
+ checked={keepTaskState}
+ onCheckedChange={(event) =>
setKeepTaskState(Boolean(event.checked))}
+ style={{ marginRight: "auto" }}
+ >
+ {translate("dags:runAndTaskActions.options.keepTaskState")}
+ </Checkbox>
<Checkbox
checked={preventRunningTask}
onCheckedChange={(event) =>
setPreventRunningTask(Boolean(event.checked))}
- style={{ marginRight: "auto" }}
>
{translate("dags:runAndTaskActions.options.preventRunningTasks")}
</Checkbox>
@@ -342,6 +361,7 @@ const ClearTaskInstanceDialog = (props: Props) => {
only_failed: onlyFailed,
run_on_latest_version: runOnLatestVersion,
task_ids: taskIds,
+ ...(keepTaskState ? { keep_task_state: true } : {}),
...(preventRunningTask ? { prevent_running_task: true } :
{}),
},
});
@@ -364,6 +384,7 @@ const ClearTaskInstanceDialog = (props: Props) => {
only_failed: onlyFailed,
run_on_latest_version: runOnLatestVersion,
task_ids: allMapped ? [taskId] : [[taskId, mapIndex as
number]],
+ ...(keepTaskState ? { keep_task_state: true } : {}),
...(preventRunningTask ? { prevent_running_task: true } : {}),
},
});
diff --git a/airflow-core/src/airflow/ui/src/constants/localStorage.ts
b/airflow-core/src/airflow/ui/src/constants/localStorage.ts
index 1152dd06f5b..43b814aa0d3 100644
--- a/airflow-core/src/airflow/ui/src/constants/localStorage.ts
+++ b/airflow-core/src/airflow/ui/src/constants/localStorage.ts
@@ -35,6 +35,7 @@ export const DEFAULT_TASK_GROUPS_EXPANDED_KEY =
"default_task_groups_expanded";
export const CLEAR_RUN_DEFAULT_OPTIONS_KEY = "clear_run_default_options";
export const CLEAR_TASK_INSTANCE_DEFAULT_OPTIONS_KEY =
"clear_task_instance_default_options";
export const CLEAR_PREVENT_RUNNING_TASK_KEY = "clear_prevent_running_task";
+export const CLEAR_KEEP_TASK_STATE_KEY = "clear_keep_task_state";
export const MARK_TASK_INSTANCE_DEFAULT_OPTIONS_KEY =
"mark_task_instance_default_options";
export const DEFAULT_TASK_INSTANCE_TAB_KEY = "default_task_instance_tab";
export const DEFAULT_LANDING_PAGE_KEY = "default_landing_page";
diff --git a/airflow-core/src/airflow/ui/src/hooks/useUserSettings.test.tsx
b/airflow-core/src/airflow/ui/src/hooks/useUserSettings.test.tsx
index ab9143f1bf4..9987b1ac52d 100644
--- a/airflow-core/src/airflow/ui/src/hooks/useUserSettings.test.tsx
+++ b/airflow-core/src/airflow/ui/src/hooks/useUserSettings.test.tsx
@@ -20,6 +20,7 @@ import { act, renderHook } from "@testing-library/react";
import { afterEach, describe, expect, it } from "vitest";
import {
+ CLEAR_KEEP_TASK_STATE_KEY,
CLEAR_PREVENT_RUNNING_TASK_KEY,
CLEAR_RUN_DEFAULT_OPTIONS_KEY,
CLEAR_TASK_INSTANCE_DEFAULT_OPTIONS_KEY,
@@ -30,6 +31,7 @@ import {
} from "src/constants/localStorage";
import {
+ useClearKeepTaskStateDefault,
useClearPreventRunningTaskDefault,
useClearRunDefaultOptions,
useClearTaskInstanceDefaultOptions,
@@ -131,6 +133,25 @@ describe("useClearPreventRunningTaskDefault", () => {
});
});
+describe("useClearKeepTaskStateDefault", () => {
+ it("defaults to false", () => {
+ const { result } = renderHook(() => useClearKeepTaskStateDefault());
+
+ expect(result.current[0]).toBe(false);
+ });
+
+ it("persists a new value", () => {
+ const { result } = renderHook(() => useClearKeepTaskStateDefault());
+
+ act(() => {
+ result.current[1](true);
+ });
+
+ expect(result.current[0]).toBe(true);
+ expect(JSON.parse(localStorage.getItem(CLEAR_KEEP_TASK_STATE_KEY) ??
"false")).toBe(true);
+ });
+});
+
describe("useDefaultTaskInstanceTab", () => {
it("defaults to logs when nothing is stored", () => {
const { result } = renderHook(() => useDefaultTaskInstanceTab());
diff --git a/airflow-core/src/airflow/ui/src/hooks/useUserSettings.ts
b/airflow-core/src/airflow/ui/src/hooks/useUserSettings.ts
index cfa8877ab95..e37ab971b33 100644
--- a/airflow-core/src/airflow/ui/src/hooks/useUserSettings.ts
+++ b/airflow-core/src/airflow/ui/src/hooks/useUserSettings.ts
@@ -21,6 +21,7 @@ import { useLocalStorage } from "usehooks-ts";
import type { Direction } from "src/components/Graph/DirectionDropdown";
import {
+ CLEAR_KEEP_TASK_STATE_KEY,
CLEAR_PREVENT_RUNNING_TASK_KEY,
CLEAR_RUN_DEFAULT_OPTIONS_KEY,
CLEAR_TASK_INSTANCE_DEFAULT_OPTIONS_KEY,
@@ -60,6 +61,9 @@ export const useClearTaskInstanceDefaultOptions = () =>
export const useClearPreventRunningTaskDefault = () =>
useLocalStorage<boolean>(CLEAR_PREVENT_RUNNING_TASK_KEY, true);
+/** Default state of the "keep task state" checkbox when clearing task
instances. */
+export const useClearKeepTaskStateDefault = () =>
useLocalStorage<boolean>(CLEAR_KEEP_TASK_STATE_KEY, false);
+
/** Default selection for the "Mark as" task-instance dialog toggle (past /
future / … ). */
export const useMarkTaskInstanceDefaultOptions = () =>
useLocalStorage<Array<string>>(MARK_TASK_INSTANCE_DEFAULT_OPTIONS_KEY, []);
diff --git a/airflow-core/src/airflow/ui/src/pages/Settings/Settings.test.tsx
b/airflow-core/src/airflow/ui/src/pages/Settings/Settings.test.tsx
index 670aefb81ae..8520bc31622 100644
--- a/airflow-core/src/airflow/ui/src/pages/Settings/Settings.test.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Settings/Settings.test.tsx
@@ -23,6 +23,7 @@ import { initReactI18next } from "react-i18next";
import { afterEach, beforeAll, describe, expect, it } from "vitest";
import {
+ CLEAR_KEEP_TASK_STATE_KEY,
CLEAR_PREVENT_RUNNING_TASK_KEY,
DEFAULT_GRAPH_DIRECTION_KEY,
DEFAULT_TASK_GROUPS_EXPANDED_KEY,
@@ -45,6 +46,7 @@ beforeAll(async () => {
common: {
settings: {
clearing: {
+ keepTaskState: { helper: "helper", label: "Keep task state on
clear" },
preventRunningTask: { helper: "helper", label: "Prevent clearing
running tasks" },
runSelection: { helper: "helper", label: "Default run clear
selection" },
taskSelection: { helper: "helper", label: "Default task clear
selection" },
@@ -128,6 +130,7 @@ describe("Settings page", () => {
"default-graph-direction",
"default-task-instance-tab",
"clear-prevent-running-task",
+ "clear-keep-task-state",
]) {
expect(screen.getByTestId(testId)).toBeInTheDocument();
}
@@ -139,11 +142,15 @@ describe("Settings page", () => {
// The prevent-running switch defaults to on.
expect(screen.getByTestId("clear-prevent-running-task")).toHaveAttribute("data-state",
"checked");
+
+ // The keep-task-state switch defaults to off, matching discard-by-default.
+
expect(screen.getByTestId("clear-keep-task-state")).toHaveAttribute("data-state",
"unchecked");
});
it("reflects stored values in the controls", () => {
localStorage.setItem(DEFAULT_GRAPH_DIRECTION_KEY, JSON.stringify("DOWN"));
localStorage.setItem(CLEAR_PREVENT_RUNNING_TASK_KEY,
JSON.stringify(false));
+ localStorage.setItem(CLEAR_KEEP_TASK_STATE_KEY, JSON.stringify(true));
localStorage.setItem(DEFAULT_TASK_INSTANCE_TAB_KEY,
JSON.stringify("details"));
localStorage.setItem(DEFAULT_LANDING_PAGE_KEY,
JSON.stringify("dashboard"));
@@ -154,6 +161,7 @@ describe("Settings page", () => {
expect(screen.getByTestId("default-graph-direction")).toHaveTextContent("DOWN-LABEL");
expect(screen.getByTestId("default-task-instance-tab")).toHaveTextContent("DETAILS-TAB");
expect(screen.getByTestId("clear-prevent-running-task")).toHaveAttribute("data-state",
"unchecked");
+
expect(screen.getByTestId("clear-keep-task-state")).toHaveAttribute("data-state",
"checked");
});
});
diff --git a/airflow-core/src/airflow/ui/src/pages/Settings/Settings.tsx
b/airflow-core/src/airflow/ui/src/pages/Settings/Settings.tsx
index fcce8e29872..17a7623eff1 100644
--- a/airflow-core/src/airflow/ui/src/pages/Settings/Settings.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Settings/Settings.tsx
@@ -27,6 +27,7 @@ import type { Direction } from
"src/components/Graph/DirectionDropdown";
import type { DefaultTaskInstanceTab } from "src/constants/tab";
import {
+ useClearKeepTaskStateDefault,
useClearPreventRunningTaskDefault,
useClearRunDefaultOptions,
useClearTaskInstanceDefaultOptions,
@@ -175,6 +176,7 @@ export const Settings = () => {
const [clearRunOptions, setClearRunOptions] = useClearRunDefaultOptions();
const [clearTaskOptions, setClearTaskOptions] =
useClearTaskInstanceDefaultOptions();
const [preventRunningTask, setPreventRunningTask] =
useClearPreventRunningTaskDefault();
+ const [keepTaskState, setKeepTaskState] = useClearKeepTaskStateDefault();
const [markTaskOptions, setMarkTaskOptions] =
useMarkTaskInstanceDefaultOptions();
const [defaultTaskInstanceTab, setDefaultTaskInstanceTab] =
useDefaultTaskInstanceTab();
const [defaultLandingPage, setDefaultLandingPage] = useDefaultLandingPage();
@@ -285,6 +287,17 @@ export const Settings = () => {
helper={translate("settings.clearing.preventRunningTask.helper")}
label={translate("settings.clearing.preventRunningTask.label")}
/>
+ <SettingRow
+ control={
+ <Switch
+ checked={keepTaskState}
+ data-testid="clear-keep-task-state"
+ onCheckedChange={(event) => setKeepTaskState(event.checked)}
+ />
+ }
+ helper={translate("settings.clearing.keepTaskState.helper")}
+ label={translate("settings.clearing.keepTaskState.label")}
+ />
</Section>
<Section title={translate("settings.marking.title")}>
<ToggleSetting
diff --git
a/airflow-core/src/airflow/ui/src/pages/TaskInstances/BulkClearTaskInstancesButton.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstances/BulkClearTaskInstancesButton.tsx
index c28c0597492..70657038df9 100644
---
a/airflow-core/src/airflow/ui/src/pages/TaskInstances/BulkClearTaskInstancesButton.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/TaskInstances/BulkClearTaskInstancesButton.tsx
@@ -29,6 +29,7 @@ import { Checkbox, Modal, SegmentedControl } from
"src/system-components";
import { ActionAccordion } from "src/components/ActionAccordion";
import { ErrorAlert } from "src/components/ErrorAlert";
+import { useClearKeepTaskStateDefault } from "src/hooks/useUserSettings";
import { useBulkClearDryRun } from "src/queries/useBulkClearDryRun";
import { useBulkClearTaskInstances } from
"src/queries/useBulkClearTaskInstances";
@@ -40,19 +41,23 @@ type Props = {
const BulkClearTaskInstancesButton = ({ clearSelections, selectedTaskInstances
}: Props) => {
const { t: translate } = useTranslation();
const { onClose, onOpen, open } = useDisclosure();
+ const [keepTaskStateDefault] = useClearKeepTaskStateDefault();
const [selectedOptions, setSelectedOptions] =
useState<Array<string>>(["downstream"]);
const [note, setNote] = useState<string | null>(null);
+ const [keepTaskState, setKeepTaskState] = useState(keepTaskStateDefault);
const [preventRunningTask, setPreventRunningTask] = useState(true);
- const { bulkClear, error, isPending } = useBulkClearTaskInstances({
- clearSelections,
- onSuccessConfirm: onClose,
- });
const handleClose = () => {
setNote(null);
+ setKeepTaskState(keepTaskStateDefault);
onClose();
};
+ const { bulkClear, error, isPending } = useBulkClearTaskInstances({
+ clearSelections,
+ onSuccessConfirm: handleClose,
+ });
+
const past = selectedOptions.includes("past");
const future = selectedOptions.includes("future");
const upstream = selectedOptions.includes("upstream");
@@ -117,12 +122,20 @@ const BulkClearTaskInstancesButton = ({ clearSelections,
selectedTaskInstances }
<ActionAccordion affectedTasks={affectedTasks} groupByRunId
note={note} setNote={setNote} />
<ErrorAlert error={error} />
<Flex alignItems="center" justifyContent="space-between" mt={3}>
- <Checkbox
- checked={preventRunningTask}
- onCheckedChange={(event) =>
setPreventRunningTask(Boolean(event.checked))}
- >
- {translate("dags:runAndTaskActions.options.preventRunningTasks")}
- </Checkbox>
+ <Flex alignItems="center" gap={4}>
+ <Checkbox
+ checked={preventRunningTask}
+ onCheckedChange={(event) =>
setPreventRunningTask(Boolean(event.checked))}
+ >
+ {translate("dags:runAndTaskActions.options.preventRunningTasks")}
+ </Checkbox>
+ <Checkbox
+ checked={keepTaskState}
+ onCheckedChange={(event) =>
setKeepTaskState(Boolean(event.checked))}
+ >
+ {translate("dags:runAndTaskActions.options.keepTaskState")}
+ </Checkbox>
+ </Flex>
<Button
disabled={affectedTasks.total_entries === 0}
loading={isPending || isFetching}
@@ -133,6 +146,7 @@ const BulkClearTaskInstancesButton = ({ clearSelections,
selectedTaskInstances }
includeOnlyFailed: onlyFailed,
includePast: past,
includeUpstream: upstream,
+ keepTaskState,
note,
preventRunningTask,
});
diff --git
a/airflow-core/src/airflow/ui/src/queries/useBulkClearTaskInstances.ts
b/airflow-core/src/airflow/ui/src/queries/useBulkClearTaskInstances.ts
index 8f2274b7b00..30bd319256e 100644
--- a/airflow-core/src/airflow/ui/src/queries/useBulkClearTaskInstances.ts
+++ b/airflow-core/src/airflow/ui/src/queries/useBulkClearTaskInstances.ts
@@ -44,6 +44,7 @@ export type BulkClearOptions = {
includeOnlyFailed: boolean;
includePast: boolean;
includeUpstream: boolean;
+ keepTaskState: boolean;
note: string | null;
preventRunningTask: boolean;
};
@@ -96,6 +97,7 @@ export const useBulkClearTaskInstances = ({ clearSelections,
onSuccessConfirm }:
include_upstream: options.includeUpstream,
note: options.note,
only_failed: options.includeOnlyFailed,
+ ...(options.keepTaskState ? { keep_task_state: true } : {}),
...(options.preventRunningTask ? { prevent_running_task: true }
: {}),
task_ids: tis.map((ti) =>
ti.map_index >= 0 ? ([ti.task_id, ti.map_index] as [string,
number]) : ti.task_id,
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
index 8d93e289a44..44d8a65f2f7 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
@@ -4346,6 +4346,123 @@ class
TestPostClearTaskInstances(TestTaskInstanceEndpoint):
ti_id = response_data["task_instances"][0]["id"]
_check_task_instance_note(session, ti_id, {"content":
"placeholder-note", "user_id": None})
+ def _seed_task_state(self, session, dag_id, task_id=None, map_index=None):
+ """Store one task state key for the TI matching the given
task_id/map_index."""
+ stmt = select(TaskInstance).where(TaskInstance.dag_id == dag_id)
+ if task_id is not None:
+ stmt = stmt.where(TaskInstance.task_id == task_id)
+ if map_index is not None:
+ stmt = stmt.where(TaskInstance.map_index == map_index)
+ ti = session.scalars(stmt).one()
+ MetastoreBackend().set(
+ TaskScope(dag_id=ti.dag_id, run_id=ti.run_id, task_id=ti.task_id,
map_index=ti.map_index),
+ "job_id",
+ "app_1234",
+ session=session,
+ )
+ session.commit()
+
+ def _task_state_rows(self, session, dag_id, task_id=None, map_index=None):
+ stmt = select(TaskStateStoreModel).where(TaskStateStoreModel.dag_id ==
dag_id)
+ if task_id is not None:
+ stmt = stmt.where(TaskStateStoreModel.task_id == task_id)
+ if map_index is not None:
+ stmt = stmt.where(TaskStateStoreModel.map_index == map_index)
+ return session.scalars(stmt).all()
+
+ @pytest.mark.db_test
+ @pytest.mark.parametrize(
+ ("payload_extra", "expect_kept"),
+ [
+ pytest.param({}, False, id="default-discards"),
+ pytest.param({"keep_task_state": True}, True, id="keep-preserves"),
+ pytest.param({"keep_task_state": False}, False,
id="explicit-false-discards"),
+ ],
+ )
+ def test_clear_task_state_store(self, test_client, session, payload_extra,
expect_kept):
+ """Clearing one mapped index discards only that index's task state,
unless kept."""
+ dag_id = "example_task_mapping_second_order"
+ self.create_task_instances(
+ session,
+ dag_id=dag_id,
+ task_instances=[
+ {"logical_date": DEFAULT_DATETIME_1, "state": State.FAILED},
+ {
+ "logical_date": DEFAULT_DATETIME_1 + dt.timedelta(days=1),
+ "state": State.FAILED,
+ "map_indexes": (0, 1),
+ },
+ ],
+ update_extras=False,
+ )
+ self._seed_task_state(session, dag_id, "times_2", 0)
+ self._seed_task_state(session, dag_id, "times_2", 1)
+ assert self._task_state_rows(session, dag_id, "times_2", 0)
+ assert self._task_state_rows(session, dag_id, "times_2", 1)
+
+ response = test_client.post(
+ f"/dags/{dag_id}/clearTaskInstances",
+ json={
+ "dry_run": False,
+ "reset_dag_runs": False,
+ "only_failed": True,
+ "task_ids": [["times_2", 0]],
+ **payload_extra,
+ },
+ )
+ assert response.status_code == 200
+
+ session.expire_all()
+ assert bool(self._task_state_rows(session, dag_id, "times_2", 0)) is
expect_kept
+ # The untargeted map index is never touched, regardless of
keep_task_state.
+ assert self._task_state_rows(session, dag_id, "times_2", 1)
+
+ @pytest.mark.db_test
+ def test_clear_dry_run_does_not_discard_task_state(self, test_client,
session):
+ """A dry run previews the clear and must not touch task state."""
+ dag_id = "example_python_operator"
+ self.create_task_instances(
+ session,
+ dag_id=dag_id,
+ task_instances=[{"logical_date": DEFAULT_DATETIME_1, "state":
State.FAILED}],
+ update_extras=False,
+ )
+ self._seed_task_state(session, dag_id)
+
+ response = test_client.post(
+ f"/dags/{dag_id}/clearTaskInstances",
+ json={"dry_run": True, "reset_dag_runs": False, "only_failed":
True},
+ )
+ assert response.status_code == 200
+ assert response.json()["total_entries"] == 1
+
+ session.expire_all()
+ assert self._task_state_rows(session, dag_id)
+
+ @pytest.mark.db_test
+ @mock.patch.object(MetastoreBackend, "clear",
side_effect=RuntimeError("boom"))
+ def test_clear_task_state_store_discard_failure_fails_request(self,
mock_clear, test_client, session):
+ """A backend.clear() failure fails the whole request instead of
reporting a partial clear as success."""
+ dag_id = "example_python_operator"
+ self.create_task_instances(
+ session,
+ dag_id=dag_id,
+ task_instances=[{"logical_date": DEFAULT_DATETIME_1, "state":
State.FAILED}],
+ update_extras=False,
+ )
+ self._seed_task_state(session, dag_id)
+
+ with pytest.raises(RuntimeError, match="boom"):
+ test_client.post(
+ f"/dags/{dag_id}/clearTaskInstances",
+ json={"dry_run": False, "reset_dag_runs": False,
"only_failed": True},
+ )
+
+ session.expire_all()
+ assert self._task_state_rows(session, dag_id)
+ ti = session.scalars(select(TaskInstance).where(TaskInstance.dag_id ==
dag_id)).one()
+ assert ti.state == State.FAILED
+
@pytest.mark.parametrize(
("task_group_id", "expected_task_ids"),
[
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 4331cdb6299..b06357521f6 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -304,6 +304,13 @@ class ClearTaskInstancesBody(BaseModel):
),
] = None
prevent_running_task: Annotated[bool | None, Field(title="Prevent Running
Task")] = False
+ keep_task_state: Annotated[
+ bool | None,
+ Field(
+ description="Keep the task state store entries of the cleared task
instances so the next attempt resumes from them. By default they are discarded,
so the task starts over.",
+ title="Keep Task State",
+ ),
+ ] = False
note: Annotated[Note | None, Field(title="Note")] = None
diff --git a/airflow-ctl/src/airflowctl/api/operations.py
b/airflow-ctl/src/airflowctl/api/operations.py
index c5f0dc392f0..126dd42eb0d 100644
--- a/airflow-ctl/src/airflowctl/api/operations.py
+++ b/airflow-ctl/src/airflowctl/api/operations.py
@@ -768,7 +768,7 @@ class TasksOperations(BaseOperations):
"""Clear task instances of a Dag; with dry_run (the default) only
previews the affected task instances."""
self.response = self.client.post(
f"dags/{dag_id}/clearTaskInstances",
- json=clear_task_instances.model_dump(mode="json",
exclude_none=True),
+ json=clear_task_instances.model_dump(mode="json",
exclude_defaults=True),
)
return
TaskInstanceCollectionResponse.model_validate_json(self.response.content)
diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
index 24ed44e2620..3a7933eb7aa 100644
--- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
+++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
@@ -2001,7 +2001,7 @@ class TestTasksOperations:
)
def test_clear(self):
- expected_body = self.clear_task_instances.model_dump(mode="json",
exclude_none=True)
+ expected_body = self.clear_task_instances.model_dump(mode="json",
exclude_defaults=True)
def handle_request(request: httpx.Request) -> httpx.Response:
assert request.url.path ==
f"/api/v2/dags/{self.dag_id}/clearTaskInstances"
@@ -2015,6 +2015,23 @@ class TestTasksOperations:
response = client.tasks.clear(self.dag_id, self.clear_task_instances)
assert response == self.task_instance_collection_response
+ def test_clear_omits_default_valued_fields(self):
+ """A payload built entirely from field defaults must not send
keep_task_state.
+
+ This covers cases like a server with 3.3.2 and earlier doesn't have
this field and rejects unknown keys,
+ so sending it unconditionally would break every clear against those
servers.
+ """
+
+ def handle_request(request: httpx.Request) -> httpx.Response:
+ body = json.loads(request.content.decode())
+ assert "keep_task_state" not in body
+ return httpx.Response(
+ 200,
json=json.loads(self.task_instance_collection_response.model_dump_json())
+ )
+
+ client = make_api_client(transport=httpx.MockTransport(handle_request))
+ client.tasks.clear(self.dag_id, ClearTaskInstancesBody())
+
class TestVariablesOperations:
key = "key"