amoghrajesh commented on code in PR #72100: URL: https://github.com/apache/airflow/pull/72100#discussion_r4013165266
########## airflow-core/docs/core-concepts/resumable-tasks.rst: ########## @@ -145,26 +145,48 @@ 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. Clearing an entire Dag run, and marking a task Review Comment: Handled in [comments from kaxil](https://github.com/apache/airflow/pull/72100/commits/04276bf6fc7a80865e87ee319b12ca75fb575c1b) ########## airflow-core/newsfragments/72100.significant.rst: ########## @@ -0,0 +1,52 @@ +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. + +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"``. + +On the CLI, ``airflowctl dags clear`` discards task state the same way, 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 to opt back in. ``airflow dags clear`` and ``airflow tasks clear`` (airflow-core, Review Comment: Handled in [comments from kaxil](https://github.com/apache/airflow/pull/72100/commits/04276bf6fc7a80865e87ee319b12ca75fb575c1b) -- 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]
