amoghrajesh commented on code in PR #71211:
URL: https://github.com/apache/airflow/pull/71211#discussion_r3765733861


##########
providers/amazon/src/airflow/providers/amazon/aws/operators/glue.py:
##########
@@ -384,6 +364,147 @@ def on_kill(self):
             if not response["SuccessfulSubmissions"]:
                 self.log.error("Failed to stop AWS Glue Job: %s. Run Id: %s", 
self.job_name, self._job_run_id)
 
+    def _set_job_run_id(self, context: Context, job_run_id: str) -> None:
+        """Record the run id and surface its console link; guarded so one 
attempt logs it once."""
+        if self._job_run_id == job_run_id:
+            return
+        self._job_run_id = job_run_id
+        GlueJobRunDetailsLink.persist(
+            context=context,
+            operator=self,
+            region_name=self.hook.conn_region_name,
+            aws_partition=self.hook.conn_partition,
+            job_name=urllib.parse.quote(self.job_name, safe=""),
+            job_run_id=job_run_id,
+        )
+        self.log.info(
+            "You can monitor this Glue Job run at: %s",
+            GlueJobRunDetailsLink.format_str.format(
+                
aws_domain=GlueJobRunDetailsLink.get_aws_domain(self.hook.conn_partition),
+                region_name=self.hook.conn_region_name,
+                job_name=urllib.parse.quote(self.job_name, safe=""),
+                job_run_id=job_run_id,
+            ),
+        )
+
+    def _build_script_args(self, context: Context) -> dict:
+        script_args = dict(self.script_args)
+        if self.openlineage_inject_parent_job_info:
+            self.log.info("Injecting OpenLineage parent job information into 
Glue job arguments.")
+            script_args = 
inject_parent_job_information_into_glue_arguments(script_args, context)
+        if self.openlineage_inject_transport_info:
+            self.log.info("Injecting OpenLineage transport information into 
Glue job arguments.")
+            script_args = 
inject_transport_information_into_glue_arguments(script_args, context)
+        if self.durable:
+            script_args, _ = self._prepare_script_args_with_task_uuid(context, 
base_args=script_args)
+        return script_args
+
+    def _find_previous_job_run(self, context: Context, task_uuid: str) -> str 
| None:
+        """
+        Look for a Glue job run this task instance already started.
+
+        Checks XCom for a cached run id first, then falls back to a task-UUID 
scan. The XCom tier
+        only works on Airflow 2.11-3.2; Airflow 3.3+ clears task XComs before 
every non-deferral
+        attempt, so it always misses there and the scan runs every time.
+        """
+        ti = context["ti"]
+        previous_job_run_id = ti.xcom_pull(key="glue_job_run_id", 
task_ids=ti.task_id)
+        if previous_job_run_id:
+            try:
+                job_run = self.hook.conn.get_job_run(JobName=self.job_name, 
RunId=previous_job_run_id)
+                state = job_run.get("JobRun", {}).get("JobRunState")
+                self.log.info("Previous Glue job_run_id: %s, state: %s", 
previous_job_run_id, state)
+                if state in ("RUNNING", "STARTING"):
+                    return previous_job_run_id
+            except Exception:
+                self.log.warning("Failed to get previous Glue job run state", 
exc_info=True)
+        else:
+            try:
+                existing = self._find_job_run_id_by_task_uuid(task_uuid)
+                if existing:
+                    existing_job_run_id, existing_job_run_state = existing
+                    self.log.info(
+                        "Found Glue job_run_id by task UUID: %s, state: %s",
+                        existing_job_run_id,
+                        existing_job_run_state,
+                    )
+                    if existing_job_run_state in ("RUNNING", "STARTING"):
+                        ti.xcom_push(key="glue_job_run_id", 
value=existing_job_run_id)
+                        return existing_job_run_id
+            except Exception:
+                self.log.warning("Failed to find previous Glue job run by task 
UUID", exc_info=True)
+        return None
+
+    def _has_stored_external_id(self, context: Context) -> bool:
+        task_state_store = context.get("task_state_store")
+        return task_state_store is not None and 
task_state_store.get(self.external_id_key) is not None
+
+    def submit_job(self, context: Context) -> str:
+        """Start a Glue job run and return its run id, or reconnect to one 
this task already started."""
+        script_args = self._build_script_args(context)
+        # Scan only when there's nothing else to go on: first attempt, no 
store, or a store that
+        # never recorded this key. A stored id (even terminal) means the 
caller already decided.
+        if self.durable and context["ti"].try_number > 1 and not 
self._has_stored_external_id(context):
+            existing_job_run_id = self._find_previous_job_run(context, 
script_args[self.TASK_UUID_ARG])
+            if existing_job_run_id:
+                self._set_job_run_id(context, existing_job_run_id)
+                return existing_job_run_id
+        self.log.info(
+            "Initializing AWS Glue Job: %s. Wait for completion: %s",
+            self.job_name,
+            self.wait_for_completion,
+        )
+        # A prior get_job_status call may have set this to a stale, terminal 
run id while checking
+        # whether to reconnect. Clear it so on_kill has nothing to act on if 
initialize_job raises.
+        self._job_run_id = None
+        glue_job_run = self.hook.initialize_job(script_args, 
self.run_job_kwargs)
+        # Set before polling so on_kill can stop the run even if the worker 
dies immediately after.
+        self._set_job_run_id(context, glue_job_run["JobRunId"])
+        # Downstream tasks read this key directly; it's also what feeds the 
XCom tier above.
+        context["ti"].xcom_push(key="glue_job_run_id", 
value=glue_job_run["JobRunId"])
+        return glue_job_run["JobRunId"]
+
+    def get_job_status(self, external_id: JsonValue, context: Context) -> str:
+        """Query the raw job run state; a run id Glue no longer knows about 
degrades to NOT_FOUND."""
+        job_run_id = cast("str", external_id)
+        # This is the first place a reconnecting attempt learns the run id.
+        self._set_job_run_id(context, job_run_id)
+        try:
+            return self.hook.get_job_state(self.job_name, job_run_id)
+        except ClientError as e:
+            if e.response["Error"]["Code"] == "EntityNotFoundException":
+                return NOT_FOUND_STATE
+            raise
+
+    def is_job_active(self, status: str) -> bool:
+        return status not in (*JOB_RUN_TERMINAL_STATES, NOT_FOUND_STATE)
+
+    def is_job_succeeded(self, status: str) -> bool:
+        return status in JOB_RUN_SUCCESS_STATES
+
+    def poll_until_complete(self, external_id: JsonValue, context: Context) -> 
None:
+        job_run_id = cast("str", external_id)
+        self._set_job_run_id(context, job_run_id)
+        if not self.wait_for_completion:
+            self.log.info("AWS Glue Job: %s. Run Id: %s", self.job_name, 
job_run_id)
+            return
+        glue_job_run = self.hook.job_completion(
+            self.job_name, job_run_id, self.verbose, self.sleep_before_return
+        )
+        state = glue_job_run["JobRunState"]
+        self.log.info("AWS Glue Job: %s status: %s. Run Id: %s", 
self.job_name, state, job_run_id)
+        if state not in JOB_RUN_SUCCESS_STATES:

Review Comment:
   Moved the STOPPED as failure statement out of the "Durable execution" 
section entirely and stated plainly that `durable=False` does not change it, 
since the raise fires unconditionally. 
   
   While tracing this I found the deferrable verbose path (a separate 
hand-rolled loop in the trigger) still treats STOPPED as success, disagreeing 
with both the waiter and our fix - filed #71490 for that, out of scope here.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to