potiuk commented on code in PR #70653:
URL: https://github.com/apache/airflow/pull/70653#discussion_r4064418347
##########
task-sdk/src/airflow/sdk/bases/resumablejobmixin.py:
##########
@@ -149,42 +149,50 @@ def execute_resumable(self, context: Context) -> Any:
if external_id:
stats.incr("resumable_job.reconnect_attempt",
tags=stats_tags)
- status = self.get_job_status(external_id, context)
-
- span.set_attribute("resumable.external_id",
str(external_id))
- span.set_attribute("resumable.prior_status", status)
-
- if self.is_job_active(status):
- # Job is still running, skip submission and reconnect
to it.
- span.set_attribute("resumable.decision", "reconnect")
- stats.incr("resumable_job.reconnect_success",
tags=stats_tags)
- self.log.info(
- "Reconnecting to existing job",
- external_id_key=self.external_id_key,
- external_id=external_id,
- status=status,
- )
- reconnect_to = external_id
- elif self.is_job_succeeded(status):
- # Job already finished successfully, skip polling and
return result directly.
- span.set_attribute("resumable.decision",
"already_succeeded")
- stats.incr("resumable_job.already_succeeded",
tags=stats_tags)
- self.log.info(
- "Job already completed successfully, skipping
resubmission",
- external_id_key=self.external_id_key,
- external_id=external_id,
- )
- already_succeeded_id = external_id
- else:
- # Job is in a terminal failed state, fall through and
submit a new job.
- span.set_attribute("resumable.decision",
"terminal_resubmit")
- stats.incr("resumable_job.terminal_resubmit",
tags=stats_tags)
- self.log.warning(
- "Prior job in terminal state, resubmitting fresh",
- external_id_key=self.external_id_key,
- external_id=external_id,
- status=status,
- )
+ try:
+ status = self.get_job_status(external_id, context)
+
+ span.set_attribute("resumable.external_id",
str(external_id))
+ span.set_attribute("resumable.prior_status", status)
+
+ if self.is_job_active(status):
+ # Job is still running, skip submission and
reconnect to it.
+ span.set_attribute("resumable.decision",
"reconnect")
+ stats.incr("resumable_job.reconnect_success",
tags=stats_tags)
+ self.log.info(
+ "Reconnecting to existing job",
+ external_id_key=self.external_id_key,
+ external_id=external_id,
+ status=status,
+ )
+ reconnect_to = external_id
+ elif self.is_job_succeeded(status):
+ # Job already finished successfully, skip polling
and return result directly.
+ span.set_attribute("resumable.decision",
"already_succeeded")
+ stats.incr("resumable_job.already_succeeded",
tags=stats_tags)
+ self.log.info(
+ "Job already completed successfully, skipping
resubmission",
+ external_id_key=self.external_id_key,
+ external_id=external_id,
+ )
+ already_succeeded_id = external_id
+ else:
+ # Job is in a terminal failed state, fall through
and submit a new job.
+ span.set_attribute("resumable.decision",
"terminal_resubmit")
+ stats.incr("resumable_job.terminal_resubmit",
tags=stats_tags)
+ self.log.warning(
+ "Prior job in terminal state, resubmitting
fresh",
+ external_id_key=self.external_id_key,
+ external_id=external_id,
+ status=status,
+ )
+ except Exception:
Review Comment:
The `try` starts at line 152 and wraps the branch bodies, not just the
decision — so an outcome counter and this one can both fire for a single
attempt.
Each branch increments and *then* logs: `:161` → `:162`, `:172` → `:173`,
`:182` → `:183`. If one of those logging calls raises, the outcome has already
been counted and this `except` adds `reconnect_failure` on top, giving
`reconnect_attempt` 1 and outcomes 2. That's the reconciliation your E2E table
demonstrates, and it's the thing this PR exists to establish. Line 193 has the
same problem in the trace: it overwrites a `resumable.decision` that was
already set to `reconnect` / `already_succeeded` / `terminal_resubmit`.
Worth saying why I don't think this is academic: `self.log` is structlog,
and remote log processors are injected into the global chain — #66633 is
currently open about exactly that path. A logging call failing because a remote
sink is unreachable is the same class of event this metric is meant to catch.
Narrowing the `try` to the fetch and classification, then branching on the
result outside it, makes the exclusivity structural:
```python
try:
status = self.get_job_status(external_id, context)
if self.is_job_active(status):
decision = "reconnect"
elif self.is_job_succeeded(status):
decision = "already_succeeded"
else:
decision = "terminal_resubmit"
except Exception:
span.set_attribute("resumable.decision", "reconnect_failure")
stats.incr("resumable_job.reconnect_failure", tags=stats_tags)
raise
```
##########
task-sdk/tests/task_sdk/bases/test_resumablejobmixin.py:
##########
@@ -296,6 +296,26 @@ def test_terminal_resubmit_fires_when_job_failed(self):
assert "resumable_job.reconnect_success" not in called_names
assert "resumable_job.fresh_submit" not in called_names
+ @pytest.mark.parametrize("failing_method", ["get_job_status",
"is_job_active"])
Review Comment:
Both parameters here raise *before* any outcome counter fires, so they're
the two positions where double-counting can't happen — the test passes for the
same reason the bug hides.
Two additions worth having: `"is_job_succeeded"` (the third classification
call, at `resumablejobmixin.py:169`) is one more token in this list; and a case
where the *log* call raises after an outcome was counted would pin the point in
my comment on line 189 — it fails today, and passes once the `try` is narrowed.
--
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]