Nataneljpwd commented on code in PR #71949:
URL: https://github.com/apache/airflow/pull/71949#discussion_r4148990327
##########
providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py:
##########
@@ -457,6 +458,17 @@ def _resolve_connection(self) -> dict[str, Any]:
conn_data["keytab"] =
self._create_keytab_path_from_base64_keytab(
base64_keytab, conn_data["principal"]
)
+ # Construct the Standalone Restendpoint
+ if (
+ conn.conn_type == "spark"
+ and conn_data["master"].startswith("spark://")
+ and conn_data["deploy_mode"] == "cluster"
+ and "," not in conn_data["master"] # only consider single
master, non-HA for now
Review Comment:
can you implement it to use HA? should not be too hard, should just be a
small change where the curl is retried, maybe in a loop or a helper command
##########
providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py:
##########
@@ -320,41 +352,81 @@ def
test_resolve_spark_submit_env_vars_use_krb5ccache_missing_KRB5CCNAME_env(sel
):
hook._build_spark_submit_command(self._spark_job_file)
- def test_build_track_driver_status_command(self):
+ @pytest.mark.parametrize(
+ ("conn_id", "expected_command"),
+ [
+ (
+ "spark_standalone_cluster",
+ [
+ "/usr/bin/curl",
+ "--max-time",
+ "30",
+
"http://spark-standalone-master:6066/v1/submissions/status/driver-1",
+ ],
+ ),
+ (
+ "spark_yarn_cluster",
+ [
+ "spark-submit",
+ "--master",
+ "yarn://yarn-master",
+ "--status",
+ "driver-1",
+ ],
+ ),
+ (
+ "spark_standalone_cluster_rpc_endpoint",
+ [
+ "/usr/bin/curl",
+ "--max-time",
+ "30",
+
"http://spark-standalone-master-rpc-endpoint:6066/v1/submissions/status/driver-1",
+ ],
+ ),
+ (
+ "spark_standalone_cluster_ha",
+ [
+ "spark-submit",
+ "--master",
+ "spark://m1:6066,m2:6066",
+ "--status",
+ "driver-1",
+ ],
+ ),
+ (
+ "spark_standalone_cluster_ipv6",
+ [
+ "/usr/bin/curl",
+ "--max-time",
+ "30",
+ "http://[2001:db8::1]:6066/v1/submissions/status/driver-1",
+ ],
+ ),
+ (
+ "spark_standalone_cluster_ipv6_ha",
+ [
+ "spark-submit",
+ "--master",
+ "spark://[2001:db8::1]:6066,[1993:db8::1]:6066",
+ "--status",
+ "driver-1",
+ ],
+ ),
+ ],
+ )
+ def test_build_track_driver_status_command(self, conn_id,
expected_command):
# note this function is only relevant for spark setup matching below
condition
# 'spark://' in self._connection['master'] and
self._connection['deploy_mode'] == 'cluster'
# Given
- hook_spark_standalone_cluster =
SparkSubmitHook(conn_id="spark_standalone_cluster")
- hook_spark_standalone_cluster._driver_id = "driver-20171128111416-0001"
- hook_spark_yarn_cluster = SparkSubmitHook(conn_id="spark_yarn_cluster")
- hook_spark_yarn_cluster._driver_id = "driver-20171128111417-0001"
+ hook = SparkSubmitHook(conn_id=conn_id)
+ hook._driver_id = "driver-1"
Review Comment:
I like this change done to the test, looks great and way better, thank you!
##########
providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py:
##########
@@ -487,6 +560,7 @@ def test_resolve_connection_yarn_default_connection(self):
"keytab": None,
"rest_scheme": "http",
"rest_port": 6066,
+ "rest_endpoint": None,
Review Comment:
some of the connections you changed look very similar, maybe it is worth
creating some kind of fixture for it?
##########
providers/apache/spark/tests/unit/apache/spark/hooks/test_spark_submit.py:
##########
@@ -747,10 +829,63 @@ def
test_resolve_connection_spark_standalone_cluster_connection(self):
"keytab": None,
"rest_scheme": "http",
"rest_port": 6066,
+ "rest_endpoint": "http://spark-standalone-master:6066",
}
assert connection == expected_spark_connection
assert cmd[0] == "spark-submit"
+ @pytest.mark.parametrize(
+ ("master", "rest_scheme", "rest_port", "expected"),
+ [
+ (
+ "spark://spark-standalone-master-rpc-endpoint:7078",
+ "http",
+ 6067,
+ "http://spark-standalone-master-rpc-endpoint:6067",
+ ),
+ (
+ "spark://spark-standalone-master-rpc-endpoint:7077",
+ "https",
+ 7443,
+ "https://spark-standalone-master-rpc-endpoint:7443",
+ ),
+ ],
+ )
+ def
test_resolve_connection_spark_standalone_cluster_connection_rpc_endpoint(
+ self,
+ create_connection_without_db,
+ master,
+ rest_scheme,
+ rest_port,
+ expected,
+ ):
+ create_connection_without_db(
+ Connection(
+ conn_id="spark_standalone_cluster_rpc_endpoint_parametrized",
+ conn_type="spark",
+ host=master,
+ extra={
+ "deploy-mode": "cluster",
+ "rest-scheme": rest_scheme,
+ "rest-port": rest_port,
Review Comment:
would be nice to set the base rest endpoint here as well
##########
providers/apache/spark/src/airflow/providers/apache/spark/hooks/spark_submit.py:
##########
@@ -457,6 +458,17 @@ def _resolve_connection(self) -> dict[str, Any]:
conn_data["keytab"] =
self._create_keytab_path_from_base64_keytab(
base64_keytab, conn_data["principal"]
)
+ # Construct the Standalone Restendpoint
+ if (
+ conn.conn_type == "spark"
+ and conn_data["master"].startswith("spark://")
+ and conn_data["deploy_mode"] == "cluster"
+ and "," not in conn_data["master"] # only consider single
master, non-HA for now
Review Comment:
also, what if I have a custom rest_endpoint? what if it's not the same as
the connection or master? as if some api gateway routing to a different place
per submission and status? this is the reason I liked the direction more than
the competing PR, yet HA ans setting a custom rest_endpoint are what's required
IMO
--
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]