kaxil commented on code in PR #72100:
URL: https://github.com/apache/airflow/pull/72100#discussion_r4121935345


##########
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 to reconnect to, but this is not 
guaranteed. The state store still matters for the small set of triggers that 
don't implement ``on_kill`` (for example ``GlueJobCompleteTrigger`` and 
``LivyTrigger``): for those, keep task state (``keep_task_state``) when 
clearing so the next attempt reconnects to the job still running instead of 
submitting a duplicate.

Review Comment:
   Partly my fault, this one: my round-2 comment named `GlueJobCompleteTrigger` 
and `LivyTrigger` as the cases that need the box ticked, and that was wrong, 
because neither operator's deferrable path writes a job id to the state store. 
`LivyOperator` only goes through `execute_resumable` when `not self.deferrable 
and self._polling_interval > 0` (`livy.py:220-223`); the deferred branch posts 
the batch and defers with nothing stored. `GlueJobOperator` says so in its 
deferred branch (`glue.py:358-359`) and reattaches through the task-UUID scan 
in `submit_job`, gated on `durable`, whatever the state store holds. So a user 
who follows this sentence and ticks the box still gets a second Livy batch. 
Suggest dropping `keep_task_state` from this advice: for deferred Glue the 
lever is `durable=True`, and for deferred Livy it is cancelling the batch 
before clearing. The same pair is repeated in `resumable-tasks.rst:187-192`.



##########
airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py:
##########
@@ -59,6 +59,43 @@
 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()``; a 
custom
+    ``[workers] state_store_backend`` payload is left orphaned there with no 
reclaim path other than
+    its own lifecycle/TTL policy.
+
+    :param event: what prompted the discard, used as the log event name.
+    """
+    backend = _get_db_backend()

Review Comment:
   Is the metastore-only discard intended for `[state_store] backend` too, not 
just `[workers] state_store_backend`? The docstring covers the second, but the 
first is a different knob: `resolve_state_backend()` reads 
`conf.getimport("state_store", "backend")`, and every Execution API state route 
(`execution_api/routes/task_state_store.py:69-111`) plus the clear-on-success 
path (`execution_api/routes/task_instances.py:544`) go through 
`get_state_backend()`. With a custom `[state_store] backend`, the worker reads 
and writes there, while this discard deletes rows from `task_state_store` that 
the backend never wrote. The clear reports success and the next attempt resumes 
from the old checkpoint, which is the case this PR exists to prevent. Either 
discarding through `get_state_backend()`, or rejecting `keep_task_state=false` 
when a non-metastore backend is configured, would close it.



##########
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 ``KubernetesPodOperator`` with ``on_kill_action="keep_pod"``.

Review Comment:
   The `keep_pod` example doesn't hold for `KubernetesPodOperator` with the 
default `reattach_on_restart=True`. When the state store comes back empty, 
`get_or_create_pod` falls back to `find_pod` (`pod.py:745-750`), whose label 
selector leaves out the try number (`:1557`), and on a kill with `keep_pod` 
`cleanup` returns early (`:1371`), so the pod never gets the `already_checked` 
label. The next attempt reattaches to the old pod rather than starting a second 
one. `resumable-tasks.rst:180` carries the same example; the sync 
`GlueJobOperator` with `stop_job_run_on_kill=False` next to it is accurate and 
could make the point on its own.



##########
airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx:
##########
@@ -86,25 +87,29 @@ 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");
   const past = selectedOptions.includes("past");
   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);
 
   useEffect(() => {
     if (openDialog) {
       setNote(taskInstance?.note ?? null);
+      setKeepTaskState(keepTaskStateDefault);
     }
-  }, [openDialog, taskInstance?.note]);
+  }, [openDialog, taskInstance?.note, keepTaskStateDefault]);

Review Comment:
   Two small things on this effect. Because it also depends on 
`taskInstance?.note`, a refetch that changes the note while the dialog is open 
resets `keepTaskState` to the default, so a user who ticked the box and then 
confirms discards the state after all. It's rare (someone else has to edit the 
note meanwhile), but the loss is silent and irreversible; a separate effect 
keyed on `openDialog` alone for `keepTaskState` avoids it. Separately, no 
dialog test covers this wiring: the only UI test that mentions keep-task-state 
is `Settings.test.tsx`, so the round-5 group-dialog regression would have 
passed CI. One test that stores `CLEAR_KEEP_TASK_STATE_KEY=true`, opens a 
dialog, and asserts the checkbox is checked and the request carries 
`keep_task_state: true` would guard all three dialogs' shared pattern.



##########
airflow-core/docs/core-concepts/resumable-tasks.rst:
##########
@@ -145,26 +145,55 @@ 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, including via ``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, and marking a task as failed or 
success (which clears
+downstream tasks as a side effect), 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
+  ``KubernetesPodOperator`` with ``on_kill_action="keep_pod"`` opts into this, 
while
+  ``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

Review Comment:
   This reads backwards: the Triggerer cancels the trigger and then *runs* its 
`on_kill`, it doesn't cancel the `on_kill`. `cancel_triggers` calls 
`task.cancel(_USER_ACTION_CANCEL_MSG)` (`triggerer_job_runner.py:1456`), and 
the run loop catches that message and awaits `trigger.on_kill()` under 
`_ON_CANCEL_TIMEOUT` (`:1708-1711`). `task-state-store.rst:288` already 
describes it correctly. Suggest: "Instead, the Triggerer cancels the orphaned 
trigger and runs the trigger's ``on_kill``, bounded by ``[triggerer] 
on_kill_timeout``."



##########
airflow-core/docs/core-concepts/resumable-tasks.rst:
##########
@@ -145,26 +145,55 @@ 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, including via ``airflowctl 
dags clear``, which

Review Comment:
   The exceptions here miss the core CLI: `airflow tasks clear` and `airflow 
dags clear` reach `clear_task_instances()` without going through the REST 
route, and the only discard call is `routes/public/task_instances.py:995`, so 
both still keep task state. The newsfragment and `dag-run.rst:261-263` say so, 
but this is the page the others link to, and `dag-run.rst` still documents 
`airflow tasks clear`. Someone who fixes a bug and re-runs from the core CLI 
resumes from the stale checkpoint, the case described just above. Worth adding 
them to this sentence, and qualifying the one-line cross-references in 
`task-state-store.rst:307-309` and `task-and-asset-state-store.rst:63-65` as 
"through the REST API, UI or airflowctl".



-- 
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]

Reply via email to