From 294ad105fc0345220079324c15f1edd80f2fbcf8 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Thu, 16 Sep 2021 18:20:08 -0700 Subject: [PATCH 1/5] Upload docker image to ECR during feast apply Signed-off-by: Felix Wang --- sdk/python/feast/feature_store.py | 43 +++++++++++++++++++ .../infra/feature_servers/feature_server.py | 1 + sdk/python/feast/repo_operations.py | 4 ++ 3 files changed, 48 insertions(+) create mode 100644 sdk/python/feast/infra/feature_servers/feature_server.py diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 8d8344d8eb3..a25a47477b9 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -45,6 +45,7 @@ update_data_sources_with_inferred_event_timestamp_col, update_entities_with_inferred_types_from_feature_views, ) +from feast.infra.feature_servers.feature_server import FEATURE_SERVER_IMAGE_FOR_TYPE from feast.infra.provider import Provider, RetrievalJob, get_provider from feast.on_demand_feature_view import OnDemandFeatureView from feast.online_response import OnlineResponse, _infer_online_entity_rows @@ -1025,6 +1026,48 @@ def serve(self, port: int) -> None: feature_server.start_server(self, port) + @log_exceptions_and_usage + def upload_docker_image(self) -> None: + """Upload the docker image for the feature consumption server to the cloud.""" + + # TODO: add error checking and avoid hardcoding the region + repository_name = "feast-python-server-test" + feature_server_type = ( + self.config.feature_server.type if self.config.feature_server else None + ) + if feature_server_type == "aws_lambda": + import base64 + + import boto3 + import docker + from botocore.exceptions import ClientError + + docker_client = docker.from_env() + image_name = FEATURE_SERVER_IMAGE_FOR_TYPE[feature_server_type] + docker_client.images.pull(image_name) + + ecr_client = boto3.client("ecr", region_name="us-west-2") + try: + ecr_client.create_repository(repositoryName=repository_name) + except ClientError: + pass + auth_token = ecr_client.get_authorization_token()["authorizationData"][0][ + "authorizationToken" + ] + username, password = base64.b64decode(auth_token).decode("utf-8").split(":") + + sts_client = boto3.client("sts") + aws_account = sts_client.get_caller_identity()["Account"] + ecr_address = f"{aws_account}.dkr.ecr.us-west-2.amazonaws.com" + docker_client.login( + username=username, password=password, registry=ecr_address + ) + + # Pushing will likely take several minutes. + image = docker_client.images.get(image_name) + image.tag(f"{ecr_address}/{repository_name}:latest") + docker_client.api.push(f"{ecr_address}/{repository_name}:latest") + def _entity_row_to_key(row: GetOnlineFeaturesRequestV2.EntityRow) -> EntityKeyProto: names, values = zip(*row.fields.items()) diff --git a/sdk/python/feast/infra/feature_servers/feature_server.py b/sdk/python/feast/infra/feature_servers/feature_server.py new file mode 100644 index 00000000000..8ce7168acef --- /dev/null +++ b/sdk/python/feast/infra/feature_servers/feature_server.py @@ -0,0 +1 @@ +FEATURE_SERVER_IMAGE_FOR_TYPE = {"aws_lambda": "feastdockerbot/feast-python-server"} diff --git a/sdk/python/feast/repo_operations.py b/sdk/python/feast/repo_operations.py index a1ab7f89c24..38e831597f2 100644 --- a/sdk/python/feast/repo_operations.py +++ b/sdk/python/feast/repo_operations.py @@ -257,6 +257,10 @@ def apply_total(repo_config: RepoConfig, repo_path: Path, skip_source_validation # Commit the update to the registry only after successful infra update registry.commit() + repo_config = store.config + if repo_config.feature_server and repo_config.feature_server.enabled: + store.upload_docker_image() + def _tag_registry_entities_for_keep_delete( project: str, registry: Registry, repo: ParsedRepo From 2493f0bb786e4c84fc3e8d3b61480f77ebf31db9 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Mon, 20 Sep 2021 13:25:06 -0700 Subject: [PATCH 2/5] Change dockerhub reference and add error checking around imports Signed-off-by: Felix Wang --- sdk/python/feast/feature_store.py | 21 ++++++++++++++----- .../infra/feature_servers/feature_server.py | 2 +- 2 files changed, 17 insertions(+), 6 deletions(-) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index a25a47477b9..2107b5e3c00 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -1028,9 +1028,9 @@ def serve(self, port: int) -> None: @log_exceptions_and_usage def upload_docker_image(self) -> None: - """Upload the docker image for the feature consumption server to the cloud.""" + """Upload the docker image for the feature server to the cloud.""" - # TODO: add error checking and avoid hardcoding the region + # TODO(felixwang9817): add error checking, logging, and avoid hardcoding the region repository_name = "feast-python-server-test" feature_server_type = ( self.config.feature_server.type if self.config.feature_server else None @@ -1038,9 +1038,20 @@ def upload_docker_image(self) -> None: if feature_server_type == "aws_lambda": import base64 - import boto3 - import docker - from botocore.exceptions import ClientError + try: + import boto3 + from botocore.exceptions import ClientError + except ImportError as e: + from feast.errors import FeastExtrasDependencyImportError + + raise FeastExtrasDependencyImportError("aws", str(e)) + + try: + import docker + except ImportError as e: + from feast.errors import FeastExtrasDependencyImportError + + raise FeastExtrasDependencyImportError("docker", str(e)) docker_client = docker.from_env() image_name = FEATURE_SERVER_IMAGE_FOR_TYPE[feature_server_type] diff --git a/sdk/python/feast/infra/feature_servers/feature_server.py b/sdk/python/feast/infra/feature_servers/feature_server.py index 8ce7168acef..e5a0ad840a8 100644 --- a/sdk/python/feast/infra/feature_servers/feature_server.py +++ b/sdk/python/feast/infra/feature_servers/feature_server.py @@ -1 +1 @@ -FEATURE_SERVER_IMAGE_FOR_TYPE = {"aws_lambda": "feastdockerbot/feast-python-server"} +FEATURE_SERVER_IMAGE_FOR_TYPE = {"aws_lambda": "feastdev/feature-server"} From 5076e767a9c5aba9c00916cf65008db13f2a1146 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Mon, 20 Sep 2021 17:37:46 -0700 Subject: [PATCH 3/5] Move logic to update_infra and add logging and error checking Signed-off-by: Felix Wang --- sdk/python/feast/constants.py | 2 + sdk/python/feast/errors.py | 8 +++ sdk/python/feast/feature_store.py | 54 --------------- sdk/python/feast/infra/aws.py | 67 +++++++++++++++++-- .../infra/feature_servers/feature_server.py | 1 - .../feast/infra/passthrough_provider.py | 6 ++ sdk/python/feast/infra/provider.py | 10 +++ sdk/python/feast/repo_operations.py | 4 -- sdk/python/setup.py | 1 + 9 files changed, 90 insertions(+), 63 deletions(-) delete mode 100644 sdk/python/feast/infra/feature_servers/feature_server.py diff --git a/sdk/python/feast/constants.py b/sdk/python/feast/constants.py index 2ecd4f9307b..a53eb860422 100644 --- a/sdk/python/feast/constants.py +++ b/sdk/python/feast/constants.py @@ -16,3 +16,5 @@ # Maximum interval(secs) to wait between retries for retry function MAX_WAIT_INTERVAL: str = "60" + +AWS_LAMBDA_FEATURE_SERVER_IMAGE = "feastdev/feature-server" diff --git a/sdk/python/feast/errors.py b/sdk/python/feast/errors.py index 0d4fb929d95..b7f97627d62 100644 --- a/sdk/python/feast/errors.py +++ b/sdk/python/feast/errors.py @@ -210,6 +210,14 @@ def __init__( ) +class DockerDaemonNotRunning(Exception): + def __init__(self): + super().__init__( + "The Docker Python sdk cannot connect to the Docker daemon. Please make sure you have" + "the docker daemon installed, and that it is running." + ) + + class RegistryInferenceFailure(Exception): def __init__(self, repo_obj_type: str, specific_issue: str): super().__init__( diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 2107b5e3c00..8d8344d8eb3 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -45,7 +45,6 @@ update_data_sources_with_inferred_event_timestamp_col, update_entities_with_inferred_types_from_feature_views, ) -from feast.infra.feature_servers.feature_server import FEATURE_SERVER_IMAGE_FOR_TYPE from feast.infra.provider import Provider, RetrievalJob, get_provider from feast.on_demand_feature_view import OnDemandFeatureView from feast.online_response import OnlineResponse, _infer_online_entity_rows @@ -1026,59 +1025,6 @@ def serve(self, port: int) -> None: feature_server.start_server(self, port) - @log_exceptions_and_usage - def upload_docker_image(self) -> None: - """Upload the docker image for the feature server to the cloud.""" - - # TODO(felixwang9817): add error checking, logging, and avoid hardcoding the region - repository_name = "feast-python-server-test" - feature_server_type = ( - self.config.feature_server.type if self.config.feature_server else None - ) - if feature_server_type == "aws_lambda": - import base64 - - try: - import boto3 - from botocore.exceptions import ClientError - except ImportError as e: - from feast.errors import FeastExtrasDependencyImportError - - raise FeastExtrasDependencyImportError("aws", str(e)) - - try: - import docker - except ImportError as e: - from feast.errors import FeastExtrasDependencyImportError - - raise FeastExtrasDependencyImportError("docker", str(e)) - - docker_client = docker.from_env() - image_name = FEATURE_SERVER_IMAGE_FOR_TYPE[feature_server_type] - docker_client.images.pull(image_name) - - ecr_client = boto3.client("ecr", region_name="us-west-2") - try: - ecr_client.create_repository(repositoryName=repository_name) - except ClientError: - pass - auth_token = ecr_client.get_authorization_token()["authorizationData"][0][ - "authorizationToken" - ] - username, password = base64.b64decode(auth_token).decode("utf-8").split(":") - - sts_client = boto3.client("sts") - aws_account = sts_client.get_caller_identity()["Account"] - ecr_address = f"{aws_account}.dkr.ecr.us-west-2.amazonaws.com" - docker_client.login( - username=username, password=password, registry=ecr_address - ) - - # Pushing will likely take several minutes. - image = docker_client.images.get(image_name) - image.tag(f"{ecr_address}/{repository_name}:latest") - docker_client.api.push(f"{ecr_address}/{repository_name}:latest") - def _entity_row_to_key(row: GetOnlineFeaturesRequestV2.EntityRow) -> EntityKeyProto: names, values = zip(*row.fields.items()) diff --git a/sdk/python/feast/infra/aws.py b/sdk/python/feast/infra/aws.py index 0f4f2e07381..90a1e4a6b59 100644 --- a/sdk/python/feast/infra/aws.py +++ b/sdk/python/feast/infra/aws.py @@ -5,6 +5,10 @@ from tempfile import TemporaryFile from urllib.parse import urlparse +from colorama import Fore, Style + +import feast +from feast.constants import AWS_LAMBDA_FEATURE_SERVER_IMAGE from feast.errors import S3RegistryBucketForbiddenAccess, S3RegistryBucketNotExist from feast.infra.passthrough_provider import PassthroughProvider from feast.protos.feast.core.Registry_pb2 import Registry as RegistryProto @@ -13,11 +17,66 @@ class AwsProvider(PassthroughProvider): - """ - This class only exists for backwards compatibility. - """ + def _upload_docker_image(self) -> None: + import base64 + + try: + import boto3 + except ImportError as e: + from feast.errors import FeastExtrasDependencyImportError + + raise FeastExtrasDependencyImportError("aws", str(e)) + + try: + import docker + from docker.errors import APIError + except ImportError as e: + from feast.errors import FeastExtrasDependencyImportError - pass + raise FeastExtrasDependencyImportError("docker", str(e)) + + try: + docker_client = docker.from_env() + except APIError: + from feast.errors import DockerDaemonNotRunning + + raise DockerDaemonNotRunning() + + print( + f"Pulling remote image {Style.BRIGHT + Fore.GREEN}{AWS_LAMBDA_FEATURE_SERVER_IMAGE}{Style.RESET_ALL}:" + ) + docker_client.images.pull(AWS_LAMBDA_FEATURE_SERVER_IMAGE) + + repository_name = "feast-python-server-test" + ecr_client = boto3.client("ecr") + try: + print( + f"Creating remote ECR repository {Style.BRIGHT + Fore.GREEN}{repository_name}{Style.RESET_ALL}:" + ) + response = ecr_client.create_repository(repositoryName=repository_name) + repository_uri = response["repository"]["repositoryUri"] + except ecr_client.exceptions.RepositoryAlreadyExistsException: + response = ecr_client.describe_repositories( + repositoryNames=[repository_name] + ) + repository_uri = response["repositories"][0]["repositoryUri"] + + auth_token = ecr_client.get_authorization_token()["authorizationData"][0][ + "authorizationToken" + ] + username, password = base64.b64decode(auth_token).decode("utf-8").split(":") + + ecr_address = repository_uri.split("/")[0] + docker_client.login(username=username, password=password, registry=ecr_address) + + image = docker_client.images.get(AWS_LAMBDA_FEATURE_SERVER_IMAGE) + version = ".".join(feast.__version__.split(".")[:3]) + image_remote_name = f"{repository_uri}:{version}" + print( + f"Pushing local image to remote {Style.BRIGHT + Fore.GREEN}{image_remote_name}{Style.RESET_ALL}:" + ) + image.tag(image_remote_name) + docker_client.api.push(repository_uri, tag=version) class S3RegistryStore(RegistryStore): diff --git a/sdk/python/feast/infra/feature_servers/feature_server.py b/sdk/python/feast/infra/feature_servers/feature_server.py deleted file mode 100644 index e5a0ad840a8..00000000000 --- a/sdk/python/feast/infra/feature_servers/feature_server.py +++ /dev/null @@ -1 +0,0 @@ -FEATURE_SERVER_IMAGE_FOR_TYPE = {"aws_lambda": "feastdev/feature-server"} diff --git a/sdk/python/feast/infra/passthrough_provider.py b/sdk/python/feast/infra/passthrough_provider.py index 3e4c0a34858..081a4a9846c 100644 --- a/sdk/python/feast/infra/passthrough_provider.py +++ b/sdk/python/feast/infra/passthrough_provider.py @@ -52,6 +52,9 @@ def update_infra( partial=partial, ) + if self.repo_config.feature_server and self.repo_config.feature_server.enabled: + self._upload_docker_image() + def teardown_infra( self, project: str, @@ -147,3 +150,6 @@ def get_historical_features( full_feature_names=full_feature_names, ) return job + + def _upload_docker_image(self) -> None: + pass diff --git a/sdk/python/feast/infra/provider.py b/sdk/python/feast/infra/provider.py index 6147f19b9af..36aa6862ee7 100644 --- a/sdk/python/feast/infra/provider.py +++ b/sdk/python/feast/infra/provider.py @@ -143,10 +143,20 @@ def online_read( """ ... + @abc.abstractmethod + def _upload_docker_image(self) -> None: + """Upload the docker image for the feature server to the cloud.""" + pass + def get_provider(config: RepoConfig, repo_path: Path) -> Provider: if "." not in config.provider: if config.provider in {"gcp", "aws", "local"}: + if config.provider == "aws": + from feast.infra.aws import AwsProvider + + return AwsProvider(config) + from feast.infra.passthrough_provider import PassthroughProvider return PassthroughProvider(config) diff --git a/sdk/python/feast/repo_operations.py b/sdk/python/feast/repo_operations.py index 38e831597f2..a1ab7f89c24 100644 --- a/sdk/python/feast/repo_operations.py +++ b/sdk/python/feast/repo_operations.py @@ -257,10 +257,6 @@ def apply_total(repo_config: RepoConfig, repo_path: Path, skip_source_validation # Commit the update to the registry only after successful infra update registry.commit() - repo_config = store.config - if repo_config.feature_server and repo_config.feature_server.enabled: - store.upload_docker_image() - def _tag_registry_entities_for_keep_delete( project: str, registry: Registry, repo: ParsedRepo diff --git a/sdk/python/setup.py b/sdk/python/setup.py index e1c36693fdd..08a0c1db008 100644 --- a/sdk/python/setup.py +++ b/sdk/python/setup.py @@ -78,6 +78,7 @@ AWS_REQUIRED = [ "boto3==1.17.*", + "docker>=5.0.2", ] CI_REQUIRED = [ From 8f0ef0e667767a33935220e798fc2aaf6b25a40e Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Mon, 20 Sep 2021 18:02:23 -0700 Subject: [PATCH 4/5] Remove feature server logic from Provider ABC Signed-off-by: Felix Wang --- sdk/python/feast/infra/passthrough_provider.py | 1 + sdk/python/feast/infra/provider.py | 5 ----- 2 files changed, 1 insertion(+), 5 deletions(-) diff --git a/sdk/python/feast/infra/passthrough_provider.py b/sdk/python/feast/infra/passthrough_provider.py index 081a4a9846c..dddb3f79c87 100644 --- a/sdk/python/feast/infra/passthrough_provider.py +++ b/sdk/python/feast/infra/passthrough_provider.py @@ -152,4 +152,5 @@ def get_historical_features( return job def _upload_docker_image(self) -> None: + """Upload the docker image for the feature server to the cloud.""" pass diff --git a/sdk/python/feast/infra/provider.py b/sdk/python/feast/infra/provider.py index 36aa6862ee7..d2e15199f25 100644 --- a/sdk/python/feast/infra/provider.py +++ b/sdk/python/feast/infra/provider.py @@ -143,11 +143,6 @@ def online_read( """ ... - @abc.abstractmethod - def _upload_docker_image(self) -> None: - """Upload the docker image for the feature server to the cloud.""" - pass - def get_provider(config: RepoConfig, repo_path: Path) -> Provider: if "." not in config.provider: From 5f70464c3dd68c2f4feac84fe6738a9956dc7143 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Wed, 22 Sep 2021 13:30:23 -0700 Subject: [PATCH 5/5] Rename repository Signed-off-by: Felix Wang --- sdk/python/feast/infra/aws.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/python/feast/infra/aws.py b/sdk/python/feast/infra/aws.py index 90a1e4a6b59..d2bd3306466 100644 --- a/sdk/python/feast/infra/aws.py +++ b/sdk/python/feast/infra/aws.py @@ -47,7 +47,8 @@ def _upload_docker_image(self) -> None: ) docker_client.images.pull(AWS_LAMBDA_FEATURE_SERVER_IMAGE) - repository_name = "feast-python-server-test" + version = ".".join(feast.__version__.split(".")[:3]) + repository_name = f"feast-python-server-{version}" ecr_client = boto3.client("ecr") try: print( @@ -70,7 +71,6 @@ def _upload_docker_image(self) -> None: docker_client.login(username=username, password=password, registry=ecr_address) image = docker_client.images.get(AWS_LAMBDA_FEATURE_SERVER_IMAGE) - version = ".".join(feast.__version__.split(".")[:3]) image_remote_name = f"{repository_uri}:{version}" print( f"Pushing local image to remote {Style.BRIGHT + Fore.GREEN}{image_remote_name}{Style.RESET_ALL}:"