Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz
, '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

Add support in AWS Batch Operator for multinode jobs - #28321

Closed
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs
Closed

Add support in AWS Batch Operator for multinode jobs#28321
camilleanne wants to merge 16 commits into
apache:mainfrom
camilleanne:ct/aws-batch-operator-for-multinode-jobs

Conversation

@camilleanne

@camilleannecamilleanne commented Dec 12, 2022

Copy link
Copy Markdown
  • Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
  • Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.

closes: #25522

I had a hard time running tests locally, so I'm opening as a draft PR initially although I don't anticipate any changes beyond syncing with main, but I'd like to confirm test success before making available for review.


^ 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 Dec 12, 2022
@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 (flake8, 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.
    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

Comment on lines +158 to +159
self.container_overrides = overrides
self.node_overrides = node_overrides

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.

Just a question. Did you check are this arguments mutually exclusive?
AWS API doesn't mention it however everything might possible because even new eksPropertiesOverride not marked as exclusive.

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.

yes they are, specifying both returns an input validation error.
BTW @Taragolis , could you approve running the workflow to help get this PR ready-to-review ? 🙏

@camilleanne
camilleanne marked this pull request as ready for review December 16, 2022 01:26
@camilleanne

Copy link
Copy Markdown
Author

Ok I got a handle on all tests finally :) ready for review now.

job_id,
)
log_configuration = (
job_node_range_properties[0].get("container", {}).get("logConfiguration", {})

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.

Is it possible to have zero element in the array ? i.e. should we add a check on len == 0 and a user-friendly error message ?

Comment on lines +452 to +453
self.log.warning(
"AWS Batch job (%s) is neither a container nor multinode job. Log info not found."

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.

Maybe this could be an error log, considering the user-provided input is invalid for this kind of request ?

The other warning logs in this method are mostly informative (there are several node groups, which is important info for the user to know, but doesn't require any action), this one I think requires user action, and thus more attention.

Comment on lines +349 to +357
},
}
},
}
],
},
}
]
}

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.

beautiful 😄

Comment threadtests/providers/amazon/aws/operators/test_batch.py Outdated
@@ -108,12 +110,13 @@ class BatchOperator(BaseOperator):
"job_queue",
"overrides",

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class, so:

Suggested change
"overrides",
"container_overrides",

Thinking about it, I wonder if this is a breaking change... I don't know too well how this works 😬

@TaragolisTaragolisDec 16, 2022

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.

@camilleanne@vandonr-amz I recommend deprecate old parameter but keep it for a while, so users would have a time for change their code.
It is much easier achieve in subclass of BaseOperator (this PR case), because all arguments are keyword.

Some examples

Use old attribute value only if new attribute not set

http_conn_id=kwargs.pop("http_conn_id", None)
ifhttp_conn_id:
warnings.warn(
"Parameter `http_conn_id` is deprecated. Please use `slack_webhook_conn_id` instead.",
DeprecationWarning,
stacklevel=2,
)
ifslack_webhook_conn_id:
raiseAirflowException("You cannot provide both `slack_webhook_conn_id` and `http_conn_id`.")
slack_webhook_conn_id=http_conn_id

Use old attribute value only if it equal new attribute value or not new attribute not set

ifmax_tries:
warnings.warn(
f"Parameter `{self.__class__.__name__}.max_tries` is deprecated and will be removed "
"in a future release. Please use method `max_polling_attempts` instead.",
DeprecationWarning,
stacklevel=2,
)
ifmax_polling_attemptsandmax_polling_attempts!=max_tries:
raiseException("max_polling_attempts must be the same value as max_tries")
else:
self.max_polling_attempts=max_tries

@vandonr-amzvandonr-amzDec 16, 2022

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.

given the fact that it was just a rename for readability, I wonder if it wouldn't be simpler to just keep the old name ?
but then it could be confusing to users...

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.

I thought it fine if rename this argument because one day we might also add eks_properties_overrides, so override might confuse users more rather than change attribute name.
And deprecation warning give time to change arguments in end users DAG code.

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.

TIL the names here need to match not the parameter name in the ctor but the name of the attribute in the class

To translate that into ELI5: if you have self.foo = bar in your operator, the template field would be named "foo" to match "self.foo", not "bar". :P

@o-nikolas

Copy link
Copy Markdown
Contributor

Hey @camilleanne,

Any plans to pick this one up again?

@github-actions

Copy link
Copy Markdown
Contributor

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

@github-actionsgithub-actionsBot added the stale Stale PRs per the .github/workflows/stale.yml policy file label Mar 13, 2023
@o-nikolas

Copy link
Copy Markdown
Contributor

@vandonr-amz is taking on this work on in #29522

Closing this PR with that context.

dimberman pushed a commit that referenced this pull request Apr 12, 2023
picking up #28321 after it's been somewhat abandoned by the original author.
Addressed my own comment about empty array, and it should be good to go I think.
Initial description from @camilleanne:
Adds support for AWS Batch multinode jobs by allowing a node_overrides json object to be passed through to the boto3 submit_job method.
Adds support for multinode jobs by properly parsing the output of describe_jobs (which is different for container vs multinode) to extract the log stream name.
closes: #25522
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providersprovider:amazonAWS/Amazon - related issuesstaleStale PRs per the .github/workflows/stale.yml policy file

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support AWS Batch multinode job types

5 participants

@camilleanne@o-nikolas@ferruzzi@Taragolis@vandonr-amz