diff --git a/go/embedded/online_features.go b/go/embedded/online_features.go index de7e6d11a2e..2563e0f43b4 100644 --- a/go/embedded/online_features.go +++ b/go/embedded/online_features.go @@ -269,12 +269,12 @@ func (s *OnlineFeatureService) StartGprcServerWithLogging(host string, port int, go func() { // As soon as these signals are received from OS, try to gracefully stop the gRPC server <-s.grpcStopCh - fmt.Println("Stopping the gRPC server...") + log.Println("Stopping the gRPC server...") grpcServer.GracefulStop() if loggingService != nil { loggingService.Stop() } - fmt.Println("gRPC server terminated") + log.Println("gRPC server terminated") }() err = grpcServer.Serve(lis) @@ -314,11 +314,15 @@ func (s *OnlineFeatureService) StartHttpServerWithLogging(host string, port int, go func() { // As soon as these signals are received from OS, try to gracefully stop the gRPC server <-s.httpStopCh - fmt.Println("Stopping the HTTP server...") + log.Println("Stopping the HTTP server...") err := ser.Stop() if err != nil { - fmt.Printf("Error when stopping the HTTP server: %v\n", err) + log.Printf("Error when stopping the HTTP server: %v\n", err) } + if loggingService != nil { + loggingService.Stop() + } + log.Println("HTTP server terminated") }() return ser.Serve(host, port) diff --git a/protos/feast/core/FeatureService.proto b/protos/feast/core/FeatureService.proto index c04fa97507d..2654703cc59 100644 --- a/protos/feast/core/FeatureService.proto +++ b/protos/feast/core/FeatureService.proto @@ -54,7 +54,6 @@ message FeatureServiceMeta { message LoggingConfig { float sample_rate = 1; - google.protobuf.Duration partition_interval = 2; oneof destination { FileDestination file_destination = 3; diff --git a/protos/feast/core/Registry.proto b/protos/feast/core/Registry.proto index 1978f41064a..2c31101510b 100644 --- a/protos/feast/core/Registry.proto +++ b/protos/feast/core/Registry.proto @@ -30,9 +30,10 @@ import "feast/core/OnDemandFeatureView.proto"; import "feast/core/RequestFeatureView.proto"; import "feast/core/DataSource.proto"; import "feast/core/SavedDataset.proto"; +import "feast/core/ValidationProfile.proto"; import "google/protobuf/timestamp.proto"; -// Next id: 13 +// Next id: 14 message Registry { repeated Entity entities = 1; repeated FeatureTable feature_tables = 2; @@ -42,6 +43,7 @@ message Registry { repeated RequestFeatureView request_feature_views = 9; repeated FeatureService feature_services = 7; repeated SavedDataset saved_datasets = 11; + repeated ValidationReference validation_references = 13; Infra infra = 10; string registry_schema_version = 3; // to support migrations; incremented when schema is changed diff --git a/protos/feast/core/ValidationProfile.proto b/protos/feast/core/ValidationProfile.proto index 673a792fdf8..b660e449bd2 100644 --- a/protos/feast/core/ValidationProfile.proto +++ b/protos/feast/core/ValidationProfile.proto @@ -39,9 +39,24 @@ message GEValidationProfile { } message ValidationReference { - SavedDataset dataset = 1; - + // Unique name of validation reference within the project + string name = 1; + // Name of saved dataset used as reference dataset + string reference_dataset_name = 2; + // Name of Feast project that this object source belongs to + string project = 3; + // Description of the validation reference + string description = 4; + // User defined metadata + map tags = 5; + + // validation profiler oneof profiler { - GEValidationProfiler ge_profiler = 2; + GEValidationProfiler ge_profiler = 6; + } + + // (optional) cached validation profile (to avoid constant recalculation) + oneof cached_profile { + GEValidationProfile ge_profile = 7; } } diff --git a/sdk/python/feast/cli.py b/sdk/python/feast/cli.py index b1281d297f1..9f3cf26dee0 100644 --- a/sdk/python/feast/cli.py +++ b/sdk/python/feast/cli.py @@ -11,7 +11,7 @@ # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. - +import json import logging import warnings from datetime import datetime @@ -23,6 +23,7 @@ import yaml from colorama import Fore, Style from dateutil import parser +from pygments import formatters, highlight, lexers from feast import flags, flags_helper, utils from feast.constants import DEFAULT_FEATURE_TRANSFORMATION_SERVER_PORT @@ -758,5 +759,61 @@ def disable_alpha_features(ctx: click.Context): store.config.write_to_path(Path(repo_path)) +@cli.command("validate") +@click.option( + "--feature-service", "-f", help="Specify a feature service name", +) +@click.option( + "--reference", "-r", help="Specify a validation reference name", +) +@click.option( + "--no-profile-cache", is_flag=True, help="Do not store cached profile in registry", +) +@click.argument("start_ts") +@click.argument("end_ts") +@click.pass_context +def validate( + ctx: click.Context, + feature_service: str, + reference: str, + start_ts: str, + end_ts: str, + no_profile_cache, +): + """ + Perform validation of logged features (produced by a given feature service) against provided reference. + + START_TS and END_TS should be in ISO 8601 format, e.g. '2021-07-16T19:20:01' + """ + repo = ctx.obj["CHDIR"] + cli_check_repo(repo) + store = FeatureStore(repo_path=str(repo)) + + feature_service = store.get_feature_service(name=feature_service) + reference = store.get_validation_reference(reference) + + result = store.validate_logged_features( + source=feature_service, + reference=reference, + start=datetime.fromisoformat(start_ts), + end=datetime.fromisoformat(end_ts), + throw_exception=False, + cache_profile=not no_profile_cache, + ) + + if not result: + print(f"{Style.BRIGHT + Fore.GREEN}Validation successful!{Style.RESET_ALL}") + return + + errors = [e.to_dict() for e in result.report.errors] + formatted_json = json.dumps(errors, indent=4) + colorful_json = highlight( + formatted_json, lexers.JsonLexer(), formatters.TerminalFormatter() + ) + print(f"{Style.BRIGHT + Fore.RED}Validation failed!{Style.RESET_ALL}") + print(colorful_json) + exit(1) + + if __name__ == "__main__": cli() diff --git a/sdk/python/feast/diff/registry_diff.py b/sdk/python/feast/diff/registry_diff.py index 197bdfcefaf..33cd3df0edc 100644 --- a/sdk/python/feast/diff/registry_diff.py +++ b/sdk/python/feast/diff/registry_diff.py @@ -20,6 +20,9 @@ from feast.protos.feast.core.RequestFeatureView_pb2 import ( RequestFeatureView as RequestFeatureViewProto, ) +from feast.protos.feast.core.ValidationProfile_pb2 import ( + ValidationReference as ValidationReferenceProto, +) from feast.registry import FEAST_OBJECT_TYPES, FeastObjectType, Registry from feast.repo_contents import RepoContents @@ -103,6 +106,7 @@ def tag_objects_for_keep_delete_update_add( FeatureServiceProto, OnDemandFeatureViewProto, RequestFeatureViewProto, + ValidationReferenceProto, ) @@ -120,9 +124,9 @@ def diff_registry_objects( current_spec: FeastObjectSpecProto new_spec: FeastObjectSpecProto - if isinstance(current_proto, DataSourceProto) or isinstance( - new_proto, DataSourceProto - ): + if isinstance( + current_proto, (DataSourceProto, ValidationReferenceProto) + ) or isinstance(new_proto, (DataSourceProto, ValidationReferenceProto)): assert type(current_proto) == type(new_proto) current_spec = cast(DataSourceProto, current_proto) new_spec = cast(DataSourceProto, new_proto) diff --git a/sdk/python/feast/dqm/profilers/ge_profiler.py b/sdk/python/feast/dqm/profilers/ge_profiler.py index 93c8b7d5de8..81e2d81d8c3 100644 --- a/sdk/python/feast/dqm/profilers/ge_profiler.py +++ b/sdk/python/feast/dqm/profilers/ge_profiler.py @@ -1,4 +1,5 @@ import json +from types import FunctionType from typing import Any, Callable, Dict, List import dill @@ -140,9 +141,12 @@ def analyze_dataset(self, df: pd.DataFrame) -> Profile: return GEProfile(expectation_suite=self.user_defined_profiler(dataset)) def to_proto(self): + # keep only the code and drop context for now + # ToDo (pyalex): include some context, but not all (dill tries to pull too much) + udp = FunctionType(self.user_defined_profiler.__code__, {}) return GEValidationProfilerProto( profiler=GEValidationProfilerProto.UserDefinedProfiler( - body=dill.dumps(self.user_defined_profiler, recurse=True) + body=dill.dumps(udp, recurse=False) ) ) diff --git a/sdk/python/feast/dqm/profilers/profiler.py b/sdk/python/feast/dqm/profilers/profiler.py index 5d2e9d36bc1..e65bf8601f0 100644 --- a/sdk/python/feast/dqm/profilers/profiler.py +++ b/sdk/python/feast/dqm/profilers/profiler.py @@ -69,6 +69,7 @@ class ValidationError: missing_count: Optional[int] missing_percent: Optional[float] + observed_value: Optional[float] def __init__( self, @@ -77,12 +78,24 @@ def __init__( check_config: Optional[Any] = None, missing_count: Optional[int] = None, missing_percent: Optional[float] = None, + observed_value: Optional[float] = None, ): self.check_name = check_name self.column_name = column_name self.check_config = check_config self.missing_count = missing_count self.missing_percent = missing_percent + self.observed_value = observed_value def __repr__(self): return f"" + + def to_dict(self): + return dict( + check_name=self.check_name, + column_name=self.column_name, + check_config=self.check_config, + missing_count=self.missing_count, + missing_percent=self.missing_percent, + observed_value=self.observed_value, + ) diff --git a/sdk/python/feast/errors.py b/sdk/python/feast/errors.py index e680337d98c..07deb6401bd 100644 --- a/sdk/python/feast/errors.py +++ b/sdk/python/feast/errors.py @@ -94,6 +94,13 @@ def __init__(self, name: str, project: str): super().__init__(f"Saved dataset {name} does not exist in project {project}") +class ValidationReferenceNotFound(FeastObjectNotFoundException): + def __init__(self, name: str, project: str): + super().__init__( + f"Validation reference {name} does not exist in project {project}" + ) + + class FeastProviderLoginError(Exception): """Error class that indicates a user has not authenticated with their provider.""" diff --git a/sdk/python/feast/feast_object.py b/sdk/python/feast/feast_object.py index 4ffd693c44f..0ac0446f5f6 100644 --- a/sdk/python/feast/feast_object.py +++ b/sdk/python/feast/feast_object.py @@ -11,7 +11,11 @@ from .protos.feast.core.FeatureView_pb2 import FeatureViewSpec from .protos.feast.core.OnDemandFeatureView_pb2 import OnDemandFeatureViewSpec from .protos.feast.core.RequestFeatureView_pb2 import RequestFeatureViewSpec +from .protos.feast.core.ValidationProfile_pb2 import ( + ValidationReference as ValidationReferenceProto, +) from .request_feature_view import RequestFeatureView +from .saved_dataset import ValidationReference # Convenience type representing all Feast objects FeastObject = Union[ @@ -21,6 +25,7 @@ Entity, FeatureService, DataSource, + ValidationReference, ] FeastObjectSpecProto = Union[ @@ -30,4 +35,5 @@ EntitySpecV2, FeatureServiceSpec, DataSourceProto, + ValidationReferenceProto, ] diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index edd0f5a46c7..140bd76a835 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -608,6 +608,7 @@ def apply( OnDemandFeatureView, RequestFeatureView, FeatureService, + ValidationReference, List[FeastObject], ], objects_to_delete: Optional[List[FeastObject]] = None, @@ -669,6 +670,9 @@ def apply( data_sources_set_to_update = { ob for ob in objects if isinstance(ob, DataSource) } + validation_references_to_update = [ + ob for ob in objects if isinstance(ob, ValidationReference) + ] for fv in views_to_update: data_sources_set_to_update.add(fv.batch_source) @@ -719,6 +723,10 @@ def apply( self._registry.apply_feature_service( feature_service, project=self.project, commit=False ) + for validation_references in validation_references_to_update: + self._registry.apply_validation_reference( + validation_references, project=self.project, commit=False + ) if not partial: # Delete all registry objects that should not exist. @@ -740,6 +748,9 @@ def apply( data_sources_to_delete = [ ob for ob in objects_to_delete if isinstance(ob, DataSource) ] + validation_references_to_delete = [ + ob for ob in objects_to_delete if isinstance(ob, ValidationReference) + ] for data_source in data_sources_to_delete: self._registry.delete_data_source( @@ -765,6 +776,10 @@ def apply( self._registry.delete_feature_service( service.name, project=self.project, commit=False ) + for validation_references in validation_references_to_delete: + self._registry.delete_validation_reference( + validation_references.name, project=self.project, commit=False + ) self._get_provider().update_infra( project=self.project, @@ -2039,6 +2054,7 @@ def serve_transformations(self, port: int) -> None: def _teardown_go_server(self): self._go_server = None + @log_exceptions_and_usage def write_logged_features( self, logs: Union[pa.Table, Path], source: Union[FeatureService] ): @@ -2066,6 +2082,7 @@ def write_logged_features( registry=self._registry, ) + @log_exceptions_and_usage def validate_logged_features( self, source: Union[FeatureService], @@ -2073,6 +2090,7 @@ def validate_logged_features( end: datetime, reference: ValidationReference, throw_exception: bool = True, + cache_profile: bool = True, ) -> Optional[ValidationFailed]: """ Load logged features from an offline store and validate them against provided validation reference. @@ -2083,6 +2101,7 @@ def validate_logged_features( end: upper bound for loading logged features reference: validation reference throw_exception: throw exception or return it as a result + cache_profile: store cached profile in Feast registry Returns: Throw or return (depends on parameter) ValidationFailed exception if validation was not successful @@ -2116,8 +2135,27 @@ def validate_logged_features( return exc + if cache_profile: + self.apply(reference) + return None + @log_exceptions_and_usage + def get_validation_reference( + self, name: str, allow_cache: bool = False + ) -> ValidationReference: + """ + Retrieves a validation reference. + + Raises: + ValidationReferenceNotFoundException: The validation reference could not be found. + """ + ref = self._registry.get_validation_reference( + name, project=self.project, allow_cache=allow_cache + ) + ref._dataset = self.get_saved_dataset(ref.dataset_name) + return ref + def _validate_entity_values(join_key_values: Dict[str, List[Value]]): set_of_row_lengths = {len(v) for v in join_key_values.values()} diff --git a/sdk/python/feast/infra/passthrough_provider.py b/sdk/python/feast/infra/passthrough_provider.py index f01fd9bac61..a53788dc85a 100644 --- a/sdk/python/feast/infra/passthrough_provider.py +++ b/sdk/python/feast/infra/passthrough_provider.py @@ -278,6 +278,6 @@ def retrieve_feature_service_logs( join_key_columns=[], feature_name_columns=columns, timestamp_field=ts_column, - start_date=start_date, - end_date=end_date, + start_date=make_tzaware(start_date), + end_date=make_tzaware(end_date), ) diff --git a/sdk/python/feast/registry.py b/sdk/python/feast/registry.py index be009566d05..c46eba8a5d6 100644 --- a/sdk/python/feast/registry.py +++ b/sdk/python/feast/registry.py @@ -38,6 +38,7 @@ FeatureViewNotFoundException, OnDemandFeatureViewNotFoundException, SavedDatasetNotFound, + ValidationReferenceNotFound, ) from feast.feature_service import FeatureService from feast.feature_view import FeatureView @@ -49,7 +50,7 @@ from feast.repo_config import RegistryConfig from feast.repo_contents import RepoContents from feast.request_feature_view import RequestFeatureView -from feast.saved_dataset import SavedDataset +from feast.saved_dataset import SavedDataset, ValidationReference REGISTRY_SCHEMA_VERSION = "1" @@ -795,7 +796,7 @@ def apply_saved_dataset( self, saved_dataset: SavedDataset, project: str, commit: bool = True, ): """ - Registers a single entity with Feast + Stores a saved dataset metadata with Feast Args: saved_dataset: SavedDataset that will be added / updated to registry @@ -870,6 +871,85 @@ def list_saved_datasets( if saved_dataset.spec.project == project ] + def apply_validation_reference( + self, + validation_reference: ValidationReference, + project: str, + commit: bool = True, + ): + """ + Persist a validation reference + + Args: + validation_reference: ValidationReference that will be added / updated to registry + project: Feast project that this dataset belongs to + commit: Whether the change should be persisted immediately + """ + validation_reference_proto = validation_reference.to_proto() + validation_reference_proto.project = project + + registry_proto = self._prepare_registry_for_changes() + for idx, existing_validation_reference in enumerate( + registry_proto.validation_references + ): + if ( + existing_validation_reference.name == validation_reference_proto.name + and existing_validation_reference.project == project + ): + del registry_proto.validation_references[idx] + break + + registry_proto.validation_references.append(validation_reference_proto) + if commit: + self.commit() + + def get_validation_reference( + self, name: str, project: str, allow_cache: bool = False + ) -> ValidationReference: + """ + Retrieves a validation reference. + + Args: + name: Name of dataset + project: Feast project that this dataset belongs to + allow_cache: Whether to allow returning this dataset from a cached registry + + Returns: + Returns either the specified ValidationReference, or raises an exception if + none is found + """ + registry_proto = self._get_registry_proto(allow_cache=allow_cache) + for validation_reference in registry_proto.validation_references: + if ( + validation_reference.name == name + and validation_reference.project == project + ): + return ValidationReference.from_proto(validation_reference) + raise ValidationReferenceNotFound(name, project=project) + + def delete_validation_reference(self, name: str, project: str, commit: bool = True): + """ + Deletes a validation reference or raises an exception if not found. + + Args: + name: Name of validation reference + project: Feast project that this object belongs to + commit: Whether the change should be persisted immediately + """ + registry_proto = self._prepare_registry_for_changes() + for idx, existing_validation_reference in enumerate( + registry_proto.validation_references + ): + if ( + existing_validation_reference.name == name + and existing_validation_reference.project == project + ): + del registry_proto.validation_references[idx] + if commit: + self.commit() + return + raise ValidationReferenceNotFound(name, project=project) + def commit(self): """Commits the state of the registry cache to the remote registry store.""" if self.cached_registry_proto: diff --git a/sdk/python/feast/saved_dataset.py b/sdk/python/feast/saved_dataset.py index aead7fe8eff..e2004d15f4c 100644 --- a/sdk/python/feast/saved_dataset.py +++ b/sdk/python/feast/saved_dataset.py @@ -13,6 +13,9 @@ from feast.protos.feast.core.SavedDataset_pb2 import ( SavedDatasetStorage as SavedDatasetStorageProto, ) +from feast.protos.feast.core.ValidationProfile_pb2 import ( + ValidationReference as ValidationReferenceProto, +) if TYPE_CHECKING: from feast.infra.offline_stores.offline_store import RetrievalJob @@ -178,8 +181,8 @@ def to_proto(self) -> SavedDatasetProto: if self.feature_service_name: spec.feature_service_name = self.feature_service_name - feature_service_proto = SavedDatasetProto(spec=spec, meta=meta) - return feature_service_proto + saved_dataset_proto = SavedDatasetProto(spec=spec, meta=meta) + return saved_dataset_proto def with_retrieval_job(self, retrieval_job: "RetrievalJob") -> "SavedDataset": self._retrieval_job = retrieval_job @@ -203,21 +206,123 @@ def to_arrow(self) -> pyarrow.Table: return self._retrieval_job.to_arrow() - def as_reference(self, profiler: "Profiler") -> "ValidationReference": - return ValidationReference(profiler=profiler, dataset=self) + def as_reference(self, name: str, profiler: "Profiler") -> "ValidationReference": + return ValidationReference.from_saved_dataset( + name=name, profiler=profiler, dataset=self + ) def get_profile(self, profiler: Profiler) -> Profile: return profiler.analyze_dataset(self.to_df()) class ValidationReference: - dataset: SavedDataset + name: str + dataset_name: str + description: str + tags: Dict[str, str] profiler: Profiler - def __init__(self, dataset: SavedDataset, profiler: Profiler): - self.dataset = dataset + _profile: Optional[Profile] = None + _dataset: Optional[SavedDataset] = None + + def __init__( + self, + name: str, + dataset_name: str, + profiler: Profiler, + description: str = "", + tags: Optional[Dict[str, str]] = None, + ): + """ + Validation reference combines a reference dataset (currently only a saved dataset object can be used as + a reference) and a profiler function to generate a validation profile. + The validation profile can be cached in this object, and in this case + the saved dataset retrieval and the profiler call will happen only once. + + Validation reference is being stored in the Feast registry and can be retrieved by its name, which + must be unique within one project. + + Args: + name: the unique name for validation reference + dataset_name: the name of the saved dataset used as a reference + description: a human-readable description + tags: a dictionary of key-value pairs to store arbitrary metadata + profiler: the profiler function used to generate profile from the saved dataset + """ + self.name = name + self.dataset_name = dataset_name self.profiler = profiler + self.description = description + self.tags = tags or {} + + @classmethod + def from_saved_dataset(cls, name: str, dataset: SavedDataset, profiler: Profiler): + """ + Internal constructor to create validation reference object with actual saved dataset object + (regular constructor requires only its name). + """ + ref = ValidationReference(name, dataset.name, profiler) + ref._dataset = dataset + return ref @property def profile(self) -> Profile: - return self.profiler.analyze_dataset(self.dataset.to_df()) + if not self._profile: + if not self._dataset: + raise RuntimeError( + "In order to calculate a profile validation reference must be instantiated from a saved dataset. " + "Use ValidationReference.from_saved_dataset constructor or FeatureStore.get_validation_reference " + "to get validation reference object." + ) + + self._profile = self.profiler.analyze_dataset(self._dataset.to_df()) + return self._profile + + @classmethod + def from_proto(cls, proto: ValidationReferenceProto) -> "ValidationReference": + profiler_attr = proto.WhichOneof("profiler") + if profiler_attr == "ge_profiler": + from feast.dqm.profilers.ge_profiler import GEProfiler + + profiler = GEProfiler.from_proto(proto.ge_profiler) + else: + raise RuntimeError("Unrecognized profiler") + + profile_attr = proto.WhichOneof("cached_profile") + if profile_attr == "ge_profile": + from feast.dqm.profilers.ge_profiler import GEProfile + + profile = GEProfile.from_proto(proto.ge_profile) + elif not profile_attr: + profile = None + else: + raise RuntimeError("Unrecognized profile") + + ref = ValidationReference( + name=proto.name, + dataset_name=proto.reference_dataset_name, + profiler=profiler, + description=proto.description, + tags=dict(proto.tags), + ) + ref._profile = profile + + return ref + + def to_proto(self) -> ValidationReferenceProto: + from feast.dqm.profilers.ge_profiler import GEProfile, GEProfiler + + proto = ValidationReferenceProto( + name=self.name, + reference_dataset_name=self.dataset_name, + tags=self.tags, + description=self.description, + ge_profiler=self.profiler.to_proto() + if isinstance(self.profiler, GEProfiler) + else None, + ge_profile=self._profile.to_proto() + if isinstance(self._profile, GEProfile) + else None, + ) + + return proto diff --git a/sdk/python/requirements/py3.10-ci-requirements.txt b/sdk/python/requirements/py3.10-ci-requirements.txt index 120f0d1158e..e4b7e5447be 100644 --- a/sdk/python/requirements/py3.10-ci-requirements.txt +++ b/sdk/python/requirements/py3.10-ci-requirements.txt @@ -491,6 +491,7 @@ pyflakes==2.4.0 # via flake8 pygments==2.12.0 # via + # feast (setup.py) # ipython # sphinx pyjwt[crypto]==2.4.0 diff --git a/sdk/python/requirements/py3.10-requirements.txt b/sdk/python/requirements/py3.10-requirements.txt index 725b17f8caf..717982f012c 100644 --- a/sdk/python/requirements/py3.10-requirements.txt +++ b/sdk/python/requirements/py3.10-requirements.txt @@ -111,6 +111,8 @@ pydantic==1.9.0 # via # fastapi # feast (setup.py) +pygments==2.12.0 + # via feast (setup.py) pyparsing==3.0.9 # via packaging pyrsistent==0.18.1 diff --git a/sdk/python/requirements/py3.7-ci-requirements.txt b/sdk/python/requirements/py3.7-ci-requirements.txt index b445f86ea06..1d8c31808d4 100644 --- a/sdk/python/requirements/py3.7-ci-requirements.txt +++ b/sdk/python/requirements/py3.7-ci-requirements.txt @@ -505,6 +505,7 @@ pyflakes==2.4.0 # via flake8 pygments==2.12.0 # via + # feast (setup.py) # ipython # sphinx pyjwt[crypto]==2.4.0 diff --git a/sdk/python/requirements/py3.7-requirements.txt b/sdk/python/requirements/py3.7-requirements.txt index b0e1511d9c6..c175fff7893 100644 --- a/sdk/python/requirements/py3.7-requirements.txt +++ b/sdk/python/requirements/py3.7-requirements.txt @@ -117,6 +117,8 @@ pydantic==1.9.0 # via # fastapi # feast (setup.py) +pygments==2.12.0 + # via feast (setup.py) pyparsing==3.0.9 # via packaging pyrsistent==0.18.1 diff --git a/sdk/python/requirements/py3.8-ci-requirements.txt b/sdk/python/requirements/py3.8-ci-requirements.txt index 202f3ac2c71..af34dbbc2ff 100644 --- a/sdk/python/requirements/py3.8-ci-requirements.txt +++ b/sdk/python/requirements/py3.8-ci-requirements.txt @@ -497,6 +497,7 @@ pyflakes==2.4.0 # via flake8 pygments==2.12.0 # via + # feast (setup.py) # ipython # sphinx pyjwt[crypto]==2.4.0 diff --git a/sdk/python/requirements/py3.8-requirements.txt b/sdk/python/requirements/py3.8-requirements.txt index 98e1c6a76a9..eff4bae2689 100644 --- a/sdk/python/requirements/py3.8-requirements.txt +++ b/sdk/python/requirements/py3.8-requirements.txt @@ -113,6 +113,8 @@ pydantic==1.9.0 # via # fastapi # feast (setup.py) +pygments==2.12.0 + # via feast (setup.py) pyparsing==3.0.9 # via packaging pyrsistent==0.18.1 diff --git a/sdk/python/requirements/py3.9-ci-requirements.txt b/sdk/python/requirements/py3.9-ci-requirements.txt index d3ecdc34bf3..4147a391dce 100644 --- a/sdk/python/requirements/py3.9-ci-requirements.txt +++ b/sdk/python/requirements/py3.9-ci-requirements.txt @@ -491,6 +491,7 @@ pyflakes==2.4.0 # via flake8 pygments==2.12.0 # via + # feast (setup.py) # ipython # sphinx pyjwt[crypto]==2.4.0 diff --git a/sdk/python/requirements/py3.9-requirements.txt b/sdk/python/requirements/py3.9-requirements.txt index 3eded689a58..80199a1f5ba 100644 --- a/sdk/python/requirements/py3.9-requirements.txt +++ b/sdk/python/requirements/py3.9-requirements.txt @@ -111,6 +111,8 @@ pydantic==1.9.0 # via # fastapi # feast (setup.py) +pygments==2.12.0 + # via feast (setup.py) pyparsing==3.0.9 # via packaging pyrsistent==0.18.1 diff --git a/sdk/python/setup.cfg b/sdk/python/setup.cfg index e2d707e2720..ebb933f69de 100644 --- a/sdk/python/setup.cfg +++ b/sdk/python/setup.cfg @@ -10,7 +10,7 @@ known_first_party=feast,feast_serving_server,feast_core_server default_section=THIRDPARTY [flake8] -ignore = E203, E266, E501, W503 +ignore = E203, E266, E501, W503, C901 max-line-length = 88 max-complexity = 20 select = B,C,E,F,W,T4 diff --git a/sdk/python/tests/conftest.py b/sdk/python/tests/conftest.py index 627fda524d9..671acb3b92a 100644 --- a/sdk/python/tests/conftest.py +++ b/sdk/python/tests/conftest.py @@ -57,6 +57,14 @@ def pytest_configure(config): config.addinivalue_line( "markers", "goserver: mark tests that use the go feature server" ) + config.addinivalue_line( + "markers", + "universal_online_stores: mark tests that can be run against different online stores", + ) + config.addinivalue_line( + "markers", + "universal_offline_stores: mark tests that can be run against different offline stores", + ) def pytest_addoption(parser): diff --git a/sdk/python/tests/integration/e2e/test_go_feature_server.py b/sdk/python/tests/integration/e2e/test_go_feature_server.py index c87d79564fd..11526132ac8 100644 --- a/sdk/python/tests/integration/e2e/test_go_feature_server.py +++ b/sdk/python/tests/integration/e2e/test_go_feature_server.py @@ -103,8 +103,9 @@ def server_port(environment, server_type: str): embedded.stop_grpc_server() else: embedded.stop_http_server() + # wait for graceful stop - time.sleep(2) + time.sleep(5) @pytest.fixture diff --git a/sdk/python/tests/integration/e2e/test_validation.py b/sdk/python/tests/integration/e2e/test_validation.py index b78a8bde093..338dc77d236 100644 --- a/sdk/python/tests/integration/e2e/test_validation.py +++ b/sdk/python/tests/integration/e2e/test_validation.py @@ -1,4 +1,5 @@ import datetime +import shutil import pandas as pd import pyarrow as pa @@ -15,6 +16,7 @@ LoggingConfig, ) from feast.protos.feast.serving.ServingService_pb2 import FieldStatus +from feast.utils import make_tzaware from feast.wait import wait_retry_backoff from tests.integration.feature_repos.repo_configuration import ( construct_universal_feature_views, @@ -24,6 +26,7 @@ driver, location, ) +from tests.utils.cli_utils import CliRunner from tests.utils.logged_features import prepare_logs _features = [ @@ -128,7 +131,7 @@ def test_historical_retrieval_with_validation(environment, universal_data_source saved_dataset = store.get_saved_dataset("my_training_dataset") # If validation pass there will be no exceptions on this point - reference = saved_dataset.as_reference(profiler=configurable_profiler) + reference = saved_dataset.as_reference(name="ref", profiler=configurable_profiler) job.to_df(validation_reference=reference) @@ -161,7 +164,7 @@ def test_historical_retrieval_fails_on_validation(environment, universal_data_so job.to_df( validation_reference=store.get_saved_dataset( "my_other_dataset" - ).as_reference(profiler=profiler_with_unrealistic_expectations) + ).as_reference(name="ref", profiler=profiler_with_unrealistic_expectations) ) failed_expectations = exc_info.value.report.errors @@ -175,6 +178,7 @@ def test_historical_retrieval_fails_on_validation(environment, universal_data_so @pytest.mark.integration +@pytest.mark.universal_offline_stores def test_logged_features_validation(environment, universal_data_sources): store = environment.feature_store @@ -244,7 +248,7 @@ def validate(): start=logs_df[LOG_TIMESTAMP_FIELD].min(), end=logs_df[LOG_TIMESTAMP_FIELD].max() + datetime.timedelta(seconds=1), reference=reference_dataset.as_reference( - profiler=profiler_with_feature_metadata + name="ref", profiler=profiler_with_feature_metadata ), ) except ValidationFailed: @@ -257,3 +261,88 @@ def validate(): success = wait_retry_backoff(validate, timeout_secs=30) assert success, "Validation failed (unexpectedly)" + + +@pytest.mark.integration +def test_e2e_validation_via_cli(environment, universal_data_sources): + runner = CliRunner() + store = environment.feature_store + + (_, datasets, data_sources) = universal_data_sources + feature_views = construct_universal_feature_views(data_sources) + feature_service = FeatureService( + name="test_service", + features=[ + feature_views.customer[ + ["current_balance", "avg_passenger_count", "lifetime_trip_count"] + ], + ], + logging_config=LoggingConfig( + destination=environment.data_source_creator.create_logged_features_destination() + ), + ) + store.apply([customer(), feature_service, feature_views.customer]) + + entity_df = datasets.entity_df.drop( + columns=["order_id", "origin_id", "destination_id", "driver_id"] + ) + retrieval_job = store.get_historical_features( + entity_df=entity_df, features=feature_service, full_feature_names=True + ) + logs_df = prepare_logs(retrieval_job.to_df(), feature_service, store) + saved_dataset = store.create_saved_dataset( + from_=retrieval_job, + name="reference_for_validating_logged_features", + storage=environment.data_source_creator.create_saved_dataset_destination(), + ) + reference = saved_dataset.as_reference( + name="test_reference", profiler=configurable_profiler + ) + + schema = FeatureServiceLoggingSource( + feature_service=feature_service, project=store.project + ).get_schema(store._registry) + store.write_logged_features( + pa.Table.from_pandas(logs_df, schema=schema), source=feature_service + ) + + with runner.local_repo(example_repo_py="", offline_store="file") as local_repo: + local_repo.apply( + [customer(), feature_views.customer, feature_service, reference] + ) + local_repo._registry.apply_saved_dataset(saved_dataset, local_repo.project) + validate_args = [ + "validate", + "--feature-service", + feature_service.name, + "--reference", + reference.name, + (datetime.datetime.utcnow() - datetime.timedelta(days=7)).isoformat(), + datetime.datetime.utcnow().isoformat(), + ] + p = runner.run(validate_args, cwd=local_repo.repo_path) + + assert p.returncode == 0, p.stderr.decode() + assert "Validation successful" in p.stdout.decode(), p.stderr.decode() + + # make sure second validation will use cached profile + shutil.rmtree(saved_dataset.storage.file_options.uri) + + # Add some invalid data that would lead to failed validation + invalid_data = pd.DataFrame( + data={ + "customer_id": [0], + "current_balance": [0], + "avg_passenger_count": [0], + "lifetime_trip_count": [0], + "event_timestamp": [make_tzaware(datetime.datetime.utcnow())], + } + ) + invalid_logs = prepare_logs(invalid_data, feature_service, store) + store.write_logged_features( + pa.Table.from_pandas(invalid_logs, schema=schema), source=feature_service + ) + + p = runner.run(validate_args, cwd=local_repo.repo_path) + assert p.returncode == 1, p.stdout.decode() + assert "Validation failed" in p.stdout.decode(), p.stderr.decode() diff --git a/sdk/python/tests/integration/online_store/test_universal_online.py b/sdk/python/tests/integration/online_store/test_universal_online.py index 4afcd61c70b..d05045e2953 100644 --- a/sdk/python/tests/integration/online_store/test_universal_online.py +++ b/sdk/python/tests/integration/online_store/test_universal_online.py @@ -325,7 +325,6 @@ def get_online_features_dict( @pytest.mark.integration -@pytest.mark.universal def test_online_retrieval_with_shared_batch_source(environment, universal_data_sources): # Addresses https://github.com/feast-dev/feast/issues/2576 diff --git a/sdk/python/tests/utils/cli_utils.py b/sdk/python/tests/utils/cli_utils.py index 5d6d5722eb1..f2478a4a5ee 100644 --- a/sdk/python/tests/utils/cli_utils.py +++ b/sdk/python/tests/utils/cli_utils.py @@ -33,7 +33,9 @@ class CliRunner: """ def run(self, args: List[str], cwd: Path) -> subprocess.CompletedProcess: - return subprocess.run([sys.executable, cli.__file__] + args, cwd=cwd) + return subprocess.run( + [sys.executable, cli.__file__] + args, cwd=cwd, capture_output=True + ) def run_with_output(self, args: List[str], cwd: Path) -> Tuple[int, bytes]: try: diff --git a/sdk/python/tests/utils/logged_features.py b/sdk/python/tests/utils/logged_features.py index 155f0b27b12..dc844a60b42 100644 --- a/sdk/python/tests/utils/logged_features.py +++ b/sdk/python/tests/utils/logged_features.py @@ -9,7 +9,7 @@ import pandas as pd import pyarrow -from feast import FeatureService, FeatureStore +from feast import FeatureService, FeatureStore, FeatureView from feast.errors import FeatureViewNotFoundException from feast.feature_logging import LOG_DATE_FIELD, LOG_TIMESTAMP_FIELD, REQUEST_ID_FIELD from feast.protos.feast.serving.ServingService_pb2 import FieldStatus @@ -28,17 +28,6 @@ def prepare_logs( logs_df[LOG_DATE_FIELD] = logs_df[LOG_TIMESTAMP_FIELD].dt.date for projection in feature_service.feature_view_projections: - for feature in projection.features: - logs_df[f"{projection.name_to_use()}__{feature.name}"] = source_df[ - feature.name - ] - logs_df[ - f"{projection.name_to_use()}__{feature.name}__timestamp" - ] = source_df["event_timestamp"].dt.floor("s") - logs_df[ - f"{projection.name_to_use()}__{feature.name}__status" - ] = FieldStatus.PRESENT - try: view = store.get_feature_view(projection.name) except FeatureViewNotFoundException: @@ -51,6 +40,31 @@ def prepare_logs( entity = store.get_entity(entity_name) logs_df[entity.join_key] = source_df[entity.join_key] + for feature in projection.features: + source_field = ( + feature.name + if feature.name in source_df.columns + else f"{projection.name_to_use()}__{feature.name}" + ) + destination_field = f"{projection.name_to_use()}__{feature.name}" + logs_df[destination_field] = source_df[source_field] + logs_df[f"{destination_field}__timestamp"] = source_df[ + "event_timestamp" + ].dt.floor("s") + if logs_df[f"{destination_field}__timestamp"].dt.tz: + logs_df[f"{destination_field}__timestamp"] = logs_df[ + f"{destination_field}__timestamp" + ].dt.tz_convert(None) + logs_df[f"{destination_field}__status"] = FieldStatus.PRESENT + if isinstance(view, FeatureView) and view.ttl: + logs_df[f"{destination_field}__status"] = logs_df[ + f"{destination_field}__status" + ].mask( + logs_df[f"{destination_field}__timestamp"] + < (datetime.datetime.utcnow() - view.ttl), + FieldStatus.OUTSIDE_MAX_AGE, + ) + return logs_df diff --git a/setup.py b/setup.py index f5f092e3f50..e0eaf4f5262 100644 --- a/setup.py +++ b/setup.py @@ -64,6 +64,7 @@ "proto-plus<1.19.7", "pyarrow>=4,<7", "pydantic>=1,<2", + "pygments==2.12.0", "PyYAML>=5.4.*,<7", "tabulate==0.8.*", "tenacity>=7,<9",