anmolxlight commented on code in PR #71939:
URL: https://github.com/apache/airflow/pull/71939#discussion_r4149401966
##########
providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py:
##########
@@ -676,22 +694,31 @@ def _build_track_driver_status_command(self) -> list[str]:
"""
curl_max_wait_time = 30
spark_host = self._connection["master"]
- if spark_host.endswith(":6066"):
- spark_host = spark_host.replace("spark://", "http://")
- connection_cmd = [
- "/usr/bin/curl",
- "--max-time",
- str(curl_max_wait_time),
- f"{spark_host}/v1/submissions/status/{self._driver_id}",
- ]
- self.log.info(connection_cmd)
-
- # The driver id so we can poll for its status
+ # spark:// indicates Spark standalone cluster mode
+ if "spark://" in spark_host:
if not self._driver_id:
raise AirflowException(
"Invalid status: attempted to poll driver status but no
driver id is known. Giving up."
)
+ urls = [
+ f"{base}/v1/submissions/status/{self._driver_id}"
+ for base in self._get_standalone_rest_base_urls()
+ ]
+ if len(urls) == 1:
+ connection_cmd = [
+ "/usr/bin/curl",
+ "--max-time",
+ str(curl_max_wait_time),
+ urls[0],
+ ]
+ else:
+ # HA: try each master in order (mirrors
_StandaloneSparkSubmitBackend.get_job_status)
+ curl_cmds = " || ".join(
+ f"/usr/bin/curl --fail --max-time {curl_max_wait_time}
{shlex.quote(u)}" for u in urls
Review Comment:
Fixed in 96d8b89e3f5.
**What I checked first:** `git show upstream/main:.../spark_submit.py | grep
curl` shows `"/usr/bin/curl"` is pre-existing on main (line 682), not
introduced here. But it was worth fixing anyway, because this hook already
resolves `spark-submit` from `PATH` on every other path (via `spark_binary`),
so hardcoding curl was inconsistent.
**What changed:** `curl` is now bare, resolved from `PATH` by `sh -c`. No
`which`/`command -v` lookup needed: the shell already does that lookup as part
of running the command, and pre-resolving it in Python would bake a host path
back into the command. So this works on conda, Homebrew, and slim images where
curl is not in `/usr/bin`.
**Proof it is not hardcoded:**
`test_standalone_curl_command_resolves_curl_from_path` puts a `curl` shim
script in `tmp_path` on `PATH` and runs the emitted command. It asserts the
shim ran (`SHIM_RAN` in stdout). If the code used `/usr/bin/curl` the real curl
would run instead and the assert would fail.
```python
shim = tmp_path / "curl"
shim.write_text("#!/bin/sh\necho SHIM_RAN \"$@\"\n")
shim.chmod(0o755)
cmd = SparkSubmitHook._build_standalone_curl_command(hook, [url])
proc = subprocess.run(cmd, capture_output=True, text=True, check=False,
env={"PATH": f"{tmp_path}:{os.environ['PATH']}"})
assert "SHIM_RAN" in proc.stdout, "curl was not resolved from PATH"
```
##########
providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py:
##########
@@ -676,22 +694,31 @@ def _build_track_driver_status_command(self) -> list[str]:
"""
curl_max_wait_time = 30
spark_host = self._connection["master"]
- if spark_host.endswith(":6066"):
- spark_host = spark_host.replace("spark://", "http://")
- connection_cmd = [
- "/usr/bin/curl",
- "--max-time",
- str(curl_max_wait_time),
- f"{spark_host}/v1/submissions/status/{self._driver_id}",
- ]
- self.log.info(connection_cmd)
-
- # The driver id so we can poll for its status
+ # spark:// indicates Spark standalone cluster mode
+ if "spark://" in spark_host:
if not self._driver_id:
raise AirflowException(
"Invalid status: attempted to poll driver status but no
driver id is known. Giving up."
)
+ urls = [
+ f"{base}/v1/submissions/status/{self._driver_id}"
+ for base in self._get_standalone_rest_base_urls()
+ ]
+ if len(urls) == 1:
+ connection_cmd = [
+ "/usr/bin/curl",
+ "--max-time",
+ str(curl_max_wait_time),
+ urls[0],
+ ]
+ else:
+ # HA: try each master in order (mirrors
_StandaloneSparkSubmitBackend.get_job_status)
+ curl_cmds = " || ".join(
+ f"/usr/bin/curl --fail --max-time {curl_max_wait_time}
{shlex.quote(u)}" for u in urls
Review Comment:
You were right, and it was worse than a log-noise issue. Fixed in
96d8b89e3f5.
**You were right that stdout/stderr could ruin it.** I proved the mechanism
by actually running the command the hook emits, against a local HTTP server
serving Spark's real pretty-printed REST payload, with
`stderr=subprocess.STDOUT` exactly as `_start_driver_status_tracking` does it
(spark_submit.py:1185-1191).
curl writes its **progress meter to stderr**. Merged into stdout, its
carriage returns splice into the JSON body and `_process_spark_status_log`
parses a corrupted value. Before the fix:
```
'FINISHED-}' != 'FINISHED'
```
That is not a terminal state, so `while self._driver_status not in
["FINISHED", "UNKNOWN", ...]` never exits — the poll loop hangs until the
10-miss limit. So the answer to "does it work with `||`" is: it silently
produced a status that never terminates.
**Fix:** `--silent` (kills the meter) plus `--show-error` (keeps `curl: (7)
Failed to connect ...` for failover diagnostics).
Two things I checked so the fix is not just cosmetic:
- **Error text is safe.** `--show-error` still emits `curl: (7)` lines onto
the merged stream. Those contain neither `submissionId` nor `driverState`, so
`_process_spark_status_log` skips them. The pre-existing
`test_process_spark_driver_status_log_bad_response` already pinned that for the
single-host case;
`test_standalone_curl_command_runs_and_falls_through_to_second_master` now pins
it with a real connection error on the merged stream, asserting `_driver_status
== "FINISHED"`.
- **Exit codes still propagate.** `sh -c 'curl ... || curl ...'` returns the
last curl's status (measured rc=7 when all masters are refused), so
`missed_job_status_reports` still increments and the existing retry accounting
is unchanged.
**New test runs the real command instead of comparing a string:**
```python
proc = subprocess.run(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
text=True, check=False)
assert proc.returncode == 0, proc.stdout
assert served.is_set(), "second master was never reached"
assert "% Total" not in proc.stdout # meter must not reach stdout
hook._process_spark_status_log(iter(proc.stdout.splitlines()))
assert hook._driver_status == "FINISHED" # fails without --silent
```
That last assert is what caught the bug: it failed with `'FINISHED-}'`
before I added `--silent`.
Full hook suite: 136 passed.
##########
providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py:
##########
@@ -1292,19 +1319,39 @@ def _poll_k8s_driver_via_api(self) -> str | None:
def _build_spark_driver_kill_command(self) -> list[str]:
"""
- Construct the spark-submit command to kill a driver.
+ Construct the command to kill a driver.
:return: full command to kill a driver
"""
- # Assume that spark-submit is present in the path to the executing user
- connection_cmd = [self._connection["spark_binary"]]
+ # spark:// indicates Spark standalone cluster mode
+ if "spark://" in self._connection["master"]:
+ # spark-submit --kill derives its REST URL from the master URL
(binary RPC
+ # port), which cannot serve REST requests — use the REST API
directly.
+ urls = [
+ f"{base}/v1/submissions/kill/{self._driver_id}"
+ for base in self._get_standalone_rest_base_urls()
+ ]
+ if len(urls) == 1:
+ connection_cmd = [
+ "/usr/bin/curl",
+ "-X",
+ "DELETE",
+ urls[0],
+ ]
+ else:
+ # HA: try each master in order
+ curl_cmds = " || ".join(f"/usr/bin/curl --fail -X DELETE
{shlex.quote(u)}" for u in urls)
+ connection_cmd = ["sh", "-c", curl_cmds]
Review Comment:
Agreed, it was duplicated. Fixed in 96d8b89e3f5.
**Extracted `_build_standalone_curl_command(urls, curl_args=None)`**, used
by both `_build_track_driver_status_command` and
`_build_spark_driver_kill_command`. The kill path just passes `curl_args=["-X",
"DELETE"]`.
**And you were right about the `len` branch.** It's gone. `" || ".join(...)`
on a one-element list yields exactly one command, no trailing `||`, so no `if
len(urls) == 1` is needed. That is what made the single-master and HA cases
collapse into one code path:
```python
def _build_standalone_curl_command(self, urls, curl_args=None):
curl_max_wait_time = 30
args = ["--silent", "--show-error", "--fail", "--max-time",
str(curl_max_wait_time), *(curl_args or [])]
curl_cmds = " || ".join(f"curl {' '.join(args)} {shlex.quote(u)}" for u
in urls)
return ["sh", "-c", curl_cmds]
```
Note `sh -c` is now used for the single-master case too (it was a plain argv
list before). That is intentional and load-bearing: it is what lets the shell
resolve `curl` from `PATH` instead of hardcoding `/usr/bin/curl`, which was
your earlier comment. Both concerns are satisfied by the same line.
Also dropped `_get_standalone_rest_base_url()` (the singular back-compat
wrapper) — nothing referenced it, including the tests, so it was dead code.
The three `curl` expectations were previously weak substring/`in` checks
that would have passed on a malformed command. They are now exact-equality
asserts on the whole argv, so a wrong URL, a missing `--fail`, or a stray `||`
all fail.
--
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]