Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

AwsGlueJobOperator: add run_job_kwargs to Glue job run #16796

Merged
merged 2 commits into from
Oct 7, 2021
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 7 additions & 2 deletions airflow/providers/amazon/aws/hooks/glue.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,18 +95,23 @@ def get_iam_execution_role(self) -> Dict:
self.log.error("Failed to create aws glue job, error: %s", general_error)
raise

def initialize_job(self, script_arguments: Optional[dict] = None) -> Dict[str, str]:
def initialize_job(
self,
script_arguments: Optional[dict] = None,
run_kwargs: Optional[dict] = None,
) -> Dict[str, str]:
"""
Initializes connection with AWS Glue
to run job
:return:
"""
glue_client = self.get_conn()
script_arguments = script_arguments or {}
run_kwargs = run_kwargs or {}

try:
job_name = self.get_or_create_glue_job()
job_run = glue_client.start_job_run(JobName=job_name, Arguments=script_arguments)
job_run = glue_client.start_job_run(JobName=job_name, Arguments=script_arguments, **run_kwargs)
return job_run
except Exception as general_error:
self.log.error("Failed to run aws glue job, error: %s", general_error)
Expand Down
6 changes: 5 additions & 1 deletion airflow/providers/amazon/aws/operators/glue.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@ class AwsGlueJobOperator(BaseOperator):
:type iam_role_name: Optional[str]
:param create_job_kwargs: Extra arguments for Glue Job Creation
:type create_job_kwargs: Optional[dict]
:param run_job_kwargs: Extra arguments for Glue Job Run
:type run_job_kwargs: Optional[dict]
"""

template_fields = ('script_args',)
Expand All @@ -77,6 +79,7 @@ def __init__(
s3_bucket: Optional[str] = None,
iam_role_name: Optional[str] = None,
create_job_kwargs: Optional[dict] = None,
run_job_kwargs: Optional[dict] = None,
**kwargs,
):
super().__init__(**kwargs)
Expand All @@ -94,6 +97,7 @@ def __init__(
self.s3_protocol = "s3://"
self.s3_artifacts_prefix = 'artifacts/glue-scripts/'
self.create_job_kwargs = create_job_kwargs
self.run_job_kwargs = run_job_kwargs or {}

def execute(self, context):
"""
Expand Down Expand Up @@ -124,7 +128,7 @@ def execute(self, context):
create_job_kwargs=self.create_job_kwargs,
)
self.log.info("Initializing AWS Glue Job: %s", self.job_name)
glue_job_run = glue_job.initialize_job(self.script_args)
glue_job_run = glue_job.initialize_job(self.script_args, self.run_job_kwargs)
glue_job_run = glue_job.job_completion(self.job_name, glue_job_run['JobRunId'])
self.log.info(
"AWS Glue Job: %s status: %s. Run Id: %s",
Expand Down
3 changes: 2 additions & 1 deletion tests/providers/amazon/aws/hooks/test_glue.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ def test_get_or_create_glue_job(self, mock_get_conn, mock_get_iam_execution_role
def test_initialize_job(self, mock_get_conn, mock_get_or_create_glue_job, mock_get_job_state):
some_data_path = "s3://glue-datasets/examples/medicare/SampleData.csv"
some_script_arguments = {"--s3_input_data_path": some_data_path}
some_run_kwargs = {"NumberOfWorkers": 5}
some_script = "s3:/glue-examples/glue-scripts/sample_aws_glue_job.py"
some_s3_bucket = "my-includes"

Expand All @@ -96,7 +97,7 @@ def test_initialize_job(self, mock_get_conn, mock_get_or_create_glue_job, mock_g
s3_bucket=some_s3_bucket,
region_name=self.some_aws_region,
)
glue_job_run = glue_job_hook.initialize_job(some_script_arguments)
glue_job_run = glue_job_hook.initialize_job(some_script_arguments, some_run_kwargs)
glue_job_run_state = glue_job_hook.get_job_state(glue_job_run['JobName'], glue_job_run['JobRunId'])
assert glue_job_run_state == mock_job_run_state, 'Mocks but be equal'

Expand Down
2 changes: 1 addition & 1 deletion tests/providers/amazon/aws/operators/test_glue.py
Original file line number Diff line number Diff line change
Expand Up @@ -58,5 +58,5 @@ def test_execute_without_failure(
mock_initialize_job.return_value = {'JobRunState': 'RUNNING', 'JobRunId': '11111'}
mock_get_job_state.return_value = 'SUCCEEDED'
glue.execute(None)
mock_initialize_job.assert_called_once_with({})
mock_initialize_job.assert_called_once_with({}, {})
assert glue.job_name == 'my_test_job'