Skip to content

Implement on_killed support for tasks during DAG run timeout - #42005

Closed
MRLab12 wants to merge 7 commits into
apache:mainfrom
MRLab12:on-kill-option-for-dag-runs-fixed-branch
Closed

MRLab12 wants to merge 7 commits into
apache:mainfrom
MRLab12:on-kill-option-for-dag-runs-fixed-branch

Conversation

@MRLab12

@MRLab12 MRLab12 commented Sep 4, 2024 •

Copy link
Copy Markdown
Contributor

  • This is a PR to fix Test workflow issues in the previous PR.

This PR addresses issue #41036 by adding support for the on_killed callback on tasks that are still running when a DAG run reaches its timeout.

Key changes:

  1. Added a new configuration option call_on_kill_on_dagrun_timeout to control whether tasks should be killed when a DAG run times out (default: True).
  2. Updated logic in _schedule_dag_run to call task on_kill if call_on_kill_on_dagrun_timeout is enabled.

closes: #41036

Testing:

  • Added unit test for the new timeout handling logic
  • Verified that the new configuration option works as expected

@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch from a4127a2 to 5499321 Compare September 4, 2024 16:53
@MRLab12 MRLab12 changed the title On kill option for dag runs fixed branch Implement on_killed support for tasks during DAG run timeout Sep 4, 2024
@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch 3 times, most recently from 601e5bb to fc0db1a Compare September 5, 2024 14:20
@MRLab12

MRLab12 commented Sep 5, 2024

Copy link
Copy Markdown
Contributor Author

After the last test fix I did there, I'm thinking what would the correct default behavior for this feature should be? Currently it's default True, but if we want to normally just skip the TI, then the default should be False.

@MRLab12
MRLab12 marked this pull request as ready for review September 5, 2024 14:51
@MRLab12

MRLab12 commented Sep 17, 2024

Copy link
Copy Markdown
Contributor Author

@RNHTTR Checking in here, have you gotten a chance to look at these changes?

Comment thread airflow/jobs/scheduler_job_runner.py Outdated

@RNHTTR RNHTTR Oct 3, 2024 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should this have a specific exception handling for AttributeError in the case where the operator doesn't have an on_kill method defined? In this case, I think we'd want to self.log.info("on_kill method does not exist for task %s: %s", task_instance, e).

What do you think?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Or maybe something like...

if hasattr(task, 'on_kill'):
    try:
        ...
else:
    self.log.warn("Task %s was configured to execute its on_kill method after a DAG run timeout, but it does not have the method 'on_kill' method defined", task_instance)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If there can be a case of an operator without an on_kill, then definitely we should check for that. I think using hasattr to check for it before getting the task and using on_kill is better to prevent the AttributeError. I will make this change.

@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch 3 times, most recently from 9d373b5 to 316888c Compare October 5, 2024 16:31
@MRLab12

MRLab12 commented Oct 5, 2024

Copy link
Copy Markdown
Contributor Author

I had a conflict and made a mistake when pushing the commits again. After getting the commits pushed, the build is failing with the error: Please ask maintainer to assign the 'legacy api' label to the PR in order to continue . Is this correct?

@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch from 9b14cac to 8250b6b Compare October 6, 2024 21:49
Comment thread airflow/models/baseoperator.py Outdated

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is there any context/previous discussion about this?

The reason I ask is that I'm wary of adding yet more options of ways things can run

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@RNHTTR can help clarify, but basically this change would allow tasks to run on_kill when a dagrun reaches timeout. The use case given in the issue is that "Some users would like for externally running workloads (e.g. snowflake, emr, bigquery, etc etc) to stop executing when a DAG run times out." . With this flag the user will have more control on what behavior they want when there's a timeout.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This all stems from the fairly controversial idea of setting running tasks to skipped if a DAG run reaches its timeout.

There are a couple paths forward:

  1. Change the behavior after a dagrun_timeout is reached to mark running TIs as failed.
  2. Proceed with this PR
  3. Encourage users to use an on_skipped_callback for external systems that need to be shutdown when a task moves to the skipped state after a DAG run timeout.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ashb @RNHTTR Checking in here. Do we want to move forward with this PR?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

IMO this would be a good change. It's somewhat counterintuitive that tasks will be marked as skipped when a Dagrun times out, and enabling on_kill on Dagrun timeout is intuitive.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ashb Any thoughts on what's the best path forward?

Also, I mentioned in a comment here that there has been a significant change to the DAG properties. If we want to use this new call_on_kill_on_dagrun_timeout, where should it go now?

@gopidesupavan gopidesupavan added the legacy api Whether legacy API changes should be allowed in PR label Oct 11, 2024
@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch 2 times, most recently from 654931a to b8148b1 Compare October 14, 2024 13:43
Comment thread airflow/jobs/scheduler_job_runner.py Outdated
@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch from 1a20120 to 06fd3b4 Compare October 23, 2024 17:57
@MRLab12
MRLab12 force-pushed the on-kill-option-for-dag-runs-fixed-branch branch from 06fd3b4 to 5e4df97 Compare October 30, 2024 17:18
@MRLab12

MRLab12 commented Nov 6, 2024 •

Copy link
Copy Markdown
Contributor Author

Looks like there has been an extensive change that was merged recently here that changed how the DAG properties are defined. I'm going through the changes in the PR, but would the changes from here go in Airflow's core rather than in the Task SDK?

Trying to fix the merge conflict and a bit confused in where the properties go now.

@RNHTTR

RNHTTR commented Dec 23, 2024

Copy link
Copy Markdown
Contributor

@MRLab12 I think they definitely need to go into the Task SDK, because that's what users will interface with. Because the models inherit from their associated SDK objects (example), I don't think you need to duplicate them there unless you'll be explicitly writing code at the model level, which I don't think is necessary.

@github-actions

Copy link
Copy Markdown
Contributor

This pull request has been automatically marked as stale because it has not had recent activity. It will be closed in 5 days if no further activity occurs. Thank you for your contributions.

@github-actions github-actions Bot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Feb 11, 2025
@github-actions github-actions Bot closed this Feb 16, 2025
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:API Airflow's REST/HTTP API area:scheduler area:serialization legacy api Whether legacy API changes should be allowed in PR stale Stale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Tasks that are running when a DAG run reaches its timeout should optionally support on_killed

4 participants