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]

Reply via email to