From 698ec1011cce6109fdfdb1f3ac9cb3c55d9c8e2f Mon Sep 17 00:00:00 2001 From: Jacob Klegar Date: Thu, 14 Jan 2021 17:53:54 -0500 Subject: [PATCH 1/2] allow https url for spark ingestion jar Signed-off-by: Jacob Klegar --- sdk/python/feast/pyspark/launchers/aws/emr_utils.py | 10 ++++++---- sdk/python/feast/pyspark/launchers/k8s/k8s.py | 2 +- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/sdk/python/feast/pyspark/launchers/aws/emr_utils.py b/sdk/python/feast/pyspark/launchers/aws/emr_utils.py index 49806b4a083..f50742877c8 100644 --- a/sdk/python/feast/pyspark/launchers/aws/emr_utils.py +++ b/sdk/python/feast/pyspark/launchers/aws/emr_utils.py @@ -80,12 +80,14 @@ def _random_string(length) -> str: return "".join(random.choice(string.ascii_lowercase) for _ in range(length)) -def _upload_jar(jar_s3_prefix: str, local_path: str) -> str: - with open(local_path, "rb") as f: - uri = urlparse(os.path.join(jar_s3_prefix, os.path.basename(local_path))) +def _upload_jar(jar_s3_prefix: str, jar_path: str) -> str: + if jar_path.startswith("https://"): + return jar_path + with open(jar_path, "rb") as f: + uri = urlparse(os.path.join(jar_s3_prefix, os.path.basename(jar_path))) return urlunparse( get_staging_client(uri.scheme).upload_fileobj( - f, local_path, remote_uri=uri, + f, jar_path, remote_uri=uri, ) ) diff --git a/sdk/python/feast/pyspark/launchers/k8s/k8s.py b/sdk/python/feast/pyspark/launchers/k8s/k8s.py index 71bbd44b0cd..8fef5c09785 100644 --- a/sdk/python/feast/pyspark/launchers/k8s/k8s.py +++ b/sdk/python/feast/pyspark/launchers/k8s/k8s.py @@ -276,7 +276,7 @@ def historical_feature_retrieval( return cast(RetrievalJob, self._job_from_job_info(job_info)) def _upload_jar(self, jar_path: str) -> str: - if jar_path.startswith("s3://") or jar_path.startswith("s3a://"): + if jar_path.startswith("s3://") or jar_path.startswith("s3a://") or jar_path.startswith("https://"): return jar_path elif jar_path.startswith("file://"): local_jar_path = urlparse(jar_path).path From dead34c4223085a688eb27ba5034d81ec096a336 Mon Sep 17 00:00:00 2001 From: Jacob Klegar Date: Thu, 14 Jan 2021 18:04:12 -0500 Subject: [PATCH 2/2] lint Signed-off-by: Jacob Klegar --- sdk/python/feast/pyspark/launchers/aws/emr_utils.py | 4 +--- sdk/python/feast/pyspark/launchers/k8s/k8s.py | 6 +++++- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/sdk/python/feast/pyspark/launchers/aws/emr_utils.py b/sdk/python/feast/pyspark/launchers/aws/emr_utils.py index f50742877c8..1194ba72d33 100644 --- a/sdk/python/feast/pyspark/launchers/aws/emr_utils.py +++ b/sdk/python/feast/pyspark/launchers/aws/emr_utils.py @@ -86,9 +86,7 @@ def _upload_jar(jar_s3_prefix: str, jar_path: str) -> str: with open(jar_path, "rb") as f: uri = urlparse(os.path.join(jar_s3_prefix, os.path.basename(jar_path))) return urlunparse( - get_staging_client(uri.scheme).upload_fileobj( - f, jar_path, remote_uri=uri, - ) + get_staging_client(uri.scheme).upload_fileobj(f, jar_path, remote_uri=uri) ) diff --git a/sdk/python/feast/pyspark/launchers/k8s/k8s.py b/sdk/python/feast/pyspark/launchers/k8s/k8s.py index 8fef5c09785..5a98a184171 100644 --- a/sdk/python/feast/pyspark/launchers/k8s/k8s.py +++ b/sdk/python/feast/pyspark/launchers/k8s/k8s.py @@ -276,7 +276,11 @@ def historical_feature_retrieval( return cast(RetrievalJob, self._job_from_job_info(job_info)) def _upload_jar(self, jar_path: str) -> str: - if jar_path.startswith("s3://") or jar_path.startswith("s3a://") or jar_path.startswith("https://"): + if ( + jar_path.startswith("s3://") + or jar_path.startswith("s3a://") + or jar_path.startswith("https://") + ): return jar_path elif jar_path.startswith("file://"): local_jar_path = urlparse(jar_path).path