Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Increase waiter_max_attempts default value in EcsRunTaskOperator - #33712

Merged
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main
Aug 30, 2023
Merged

Increase waiter_max_attempts default value in EcsRunTaskOperator#33712
eladkal merged 16 commits into
apache:mainfrom
mjsqu:main

Conversation

@mjsqu

@mjsqumjsqu commented Aug 24, 2023

Copy link
Copy Markdown
Contributor

closes: #33711
related: #33016, #30586

Prior to v8.0.0 the default waiter_max_attempts parameter in EcsRunTaskOperator was set to sys.maxsize. The value is a large integer, which should always exceed the execution_timeout value for a given task. I believe this is why the line setting the value to sys.maxsize is accompanied by a comment:

waiter.config.max_attempts=sys.maxsize# timeout is managed by airflow

Changes in v8.0.0 introduced a configurable value defaulting to 100 for this parameter, with a delay defaulting to 6 seconds (the changes were introduced to resolve#30586). The line above remains, but I think it is being overridden by the call to waiter.wait:

waiter.wait(
cluster=self.cluster,
tasks=[self.arn],
WaiterConfig={
"Delay": self.waiter_delay,
"MaxAttempts": self.waiter_max_attempts, # !!! Overriding previously set max_attempts?
},
)

The defaults for these previously unavailable parameters have been set to:

waiter_delay: int=6,
waiter_max_attempts: int=100,

I believe the introduction of a default value of 100 has changed the behaviour of the operator - now when not specified, all tasks detach from the running ECS Task at 10 minutes, reporting failed status back to Airflow.

The fix replaces the default value with sys.maxsize - which should revert the waiter functionality back to prior working versions.

^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in newsfragments.

@boring-cyborgboring-cyborgBot added area:providers provider:amazon AWS/Amazon - related issues labels Aug 24, 2023
@boring-cyborg

Copy link
Copy Markdown

Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contribution Guide (https://github.com/apache/airflow/blob/main/CONTRIBUTING.rst)
Here are some useful points:

  • Pay attention to the quality of your code (ruff, mypy and type annotations). Our pre-commits will help you with that.
  • In case of a new feature add useful documentation (in docstrings or in docs/ directory). Adding a new operator? Check this short guide Consider adding an example DAG that shows how users should use it.
  • Consider using Breeze environment for testing locally, it's a heavy docker but it ships with a working Airflow and a lot of integrations.
  • Be patient and persistent. It might take some time to get a review or get the final approval from Committers.
  • Please follow ASF Code of Conduct for all communication including (but not limited to) comments on Pull Requests, Mailing list and Slack.
  • Be sure to read the Airflow Coding style.
  • Always keep your Pull Requests rebased, otherwise your build might fail due to changes not related to your commits.
    Apache Airflow is a community-driven project and together we are making it better 🚀.
    In case of doubts contact the developers at:
    Mailing List: dev@airflow.apache.org
    Slack: https://s.apache.org/airflow-slack

@mjsqu

mjsqu commented Aug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

While I wait for my sleep 800 command to run having manually set waiter_max_attempts to sys.maxsize...

Here's the documentation for the ECS TasksStopped waiter: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/ecs/waiter/TasksStopped.html

If not set, it has defaults of

WaiterConfig={
'Delay': 6,
'MaxAttempts': 100
}

too. I think this was the motivation for setting these as the defaults, however the prior functionality would override the MaxAttempts to a really high number and let Airflow's execution_timeout parameter handle long-running tasks. The execution_timeout method works well, when the Airflow timeout is hit, the ECS task is killed.

wait_for_completion: bool = True,
waiter_delay: int = 6,
waiter_max_attempts: int = 100,
waiter_max_attempts: int = sys.maxsize, # Unless specified timeout is managed by airflow

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.

There are two things here:

  1. Changing the default value to infinite could be a breaking change (cc: @eladkal )
  2. If it's okay to change it, the best way is to make it None and replace the loop in the Trigger with while True when it's None.

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

@mjsqumjsquAug 25, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks for the review. I agree that adding sys.maxsize must have been a workaround and is a bit non-standard so perhaps this PR might be able to resolve that by using some other method.

I had considered if the max_attempts and delay parameters could be calculated at runtime by taking into consideration the execution_timeout value, however this would require a bit more of a rewrite of the Operator.

@eladkaleladkalAug 25, 2023

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.

Agree it might be seen as breaking but looking at the problem description I think the previous change was accidentally breaking (i.e. changing from effectively infinite to finite) and this one is simply a bugfix restoring it.

Then I can accept this as non breaking change but a bug fix. Yet requesting to add a note explaining it in the top of the CHANGELOG. I will bake it to the proper release version during release process.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Changelog note added 456ee50

Unsure of the best way to refer back to the original PR that I am partially reverting, I trust you'll properly format the Changelog message when it comes to release time. TIA.

Changelog
---------

* ``Fix revert default waiter max attempts to sys.maxsize in 'EcsRunTaskOperator' (partial revert of #30586) (#33712)``

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.

Lets modify this to a more explanatory text.
Dont mention PR numbers. The commit title that we add during release cycle will do that and refer to this PR for more details.

What we are seeking in the note is just a simple explnatipn of the situation "A bug intoduced in provider version... Caused.... In this version we are fixing it by returning value to ..."
Something like this.

@eladkaleladkal left a comment

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.

LGTM

needs rebase and resolve conflicts

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

LGTM

needs rebase and resolve conflicts

Resolved conflict in CHANGELOG.rst - 56782d5

@eladkal

Copy link
Copy Markdown
Contributor

Tests fail :(

@mjsqu

mjsqu commented Aug 27, 2023

Copy link
Copy Markdown
ContributorAuthor

Tests fail :(

Tests were failing due to the size of sys.maxsize. The waiter_max_attempts value is used in a timedelta calculation here (https://github.com/mjsqu/airflow/blob/b9453e9674e69203ca8658aa15f7e5c24e88f4db/airflow/providers/amazon/aws/operators/ecs.py#L564)

>>>timedelta(seconds=sys.maxsize*delay+60)
Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>OverflowError: PythoninttoolargetoconverttoCint

To resolve this, I have set the default waiter_max_attempts value high enough such that the waiter would run for 1M years. This allows the timedelta calculation to succeed, while providing a value that is likely longer than any ECS Tasks will ever run for.

max_attempts=1000000*365*24*60*10>>>timedelta(seconds=max_attempts*delay+60)
datetime.timedelta(days=365000000, seconds=60)

Alternatively, if possible, it would be good to calculate a default_max_attempts at runtime using execution_timeout, I'm just unsure how.

I'm open to suggestions about other values to use instead of 1 million years, there are lower values that would be equally sensible.

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Latest test failures:

cluster="c", tasks=["arn"], WaiterConfig={"Delay": 6, "MaxAttempts": 1000000*365*24*60*10}
)
>assert1000000*365*24*60*10==client_mock.get_waiter.return_value.config.max_attemptsEAssertionError: assert ((((1000000*365) *24) *60) *10) ==<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>E+where<MagicMockname='client.get_waiter().config.max_attempts'id='139852189801392'>=<MagicMockname='client.get_waiter().config'id='139852190309584'>.max_attemptsE+where<MagicMockname='client.get_waiter().config'id='139852190309584'>=<MagicMockname='client.get_waiter()'id='139852190269200'>.configE+where<MagicMockname='client.get_waiter()'id='139852190269200'>=<MagicMockname='client.get_waiter'id='139852190252720'>.return_valueE+where<MagicMockname='client.get_waiter'id='139852190252720'>=<MagicMockname='client'id='139852190228688'>.get_waiter

Struggling to understand how the assertion fails, perhaps there's a precision issue such that client_mock.get_waiter.return_value.config.max_attempts does not return the exact value it was provided with?

@mjsqu

Copy link
Copy Markdown
ContributorAuthor

Anyone following this PR/Issue, we have set the waiter_max_attempts value to sys.maxsize in the interim for affected EcsRunTaskOperator tasks, but you could choose a sufficiently large value to force the waiter to run for longer than execution_timeout:

waiter_max_attempts * waiter_delay > execution_timeout

e.g. for an execution_timeout of 1 hour (3600 seconds), set:

waiter_max_attempts=650, # Greater than 600waiter_delay=6,
# Total waiter time = 3,900 seconds == 65 minutes

return

waiter = self.client.get_waiter("tasks_stopped")
waiter.config.max_attempts = sys.maxsize # timeout is managed by airflow

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.

The failed check was not valid after removing this line.

@mjsqumjsqu changed the title Replace waiter_max_attempts default with sys.maxsizeIncrease EcsRunOperator waiter_max_attempts default valueAug 27, 2023

@vincbeckvincbeck left a comment

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.

Thanks for fixing that bug!

Comment threadairflow/providers/amazon/aws/operators/ecs.py Outdated
@eladkaleladkal changed the title Increase EcsRunOperator waiter_max_attempts default valueIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkaleladkal changed the title Increase waiter_max_attempts default value in EcsRunTaskOperatorIncrease waiter_max_attempts default value in EcsRunTaskOperatorAug 30, 2023
@eladkal
eladkal merged commit ea44ed9 into apache:mainAug 30, 2023
@boring-cyborg

Copy link
Copy Markdown

Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

EcsRunTaskOperator waiter default waiter_max_attempts too low - all Airflow tasks detach from ECS tasks at 10 minutes

6 participants

@mjsqu@eladkal@potiuk@hussein-awala@o-nikolas@vincbeck