Class StreamFeatureGroup
-
Field Summary
FieldsFields inherited from class com.logicalclocks.hsfs.FeatureGroupBase
created, creator, dataSource, deltaStreamerJobConf, deprecated, description, eventTime, expectationsNames, featureGroupEngineBase, features, featureStore, featurestoreId, hudiPrecombineKey, id, location, LOGGER, name, notificationTopicName, onlineConfig, onlineEnabled, onlineIngestionApi, onlineTopicName, partitionKeys, primaryKeys, statisticColumns, statisticsConfig, subject, timeTravelFormat, topicName, type, utils, version -
Constructor Summary
ConstructorsConstructorDescriptionStreamFeatureGroup(FeatureStore featureStore, int id) StreamFeatureGroup(FeatureStore featureStore, @NonNull String name, Integer version, String description, List<String> primaryKeys, List<String> partitionKeys, String hudiPrecombineKey, boolean onlineEnabled, TimeTravelFormat timeTravelFormat, List<Feature> features, StatisticsConfig statisticsConfig, String onlineTopicName, String topicName, String notificationTopicName, String eventTime, OnlineConfig onlineConfig, StorageConnector storageConnector, String path) StreamFeatureGroup(Integer id, String description, List<Feature> features) -
Method Summary
Modifier and TypeMethodDescriptionvoidappendFeatures(Feature features) Append a single feature to the schema of the stream feature group.voidappendFeatures(List<Feature> features) Append features to the schema of the stream feature group.Get Query object to retrieve all features of the group at a point in the past.Get Query object to retrieve all features of the group at a point in the past.voidcommitDeleteRecord(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) Deprecated.voidcommitDeleteRecord(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) Deprecated.useremoveRows(Dataset, Map)instead.Retrieves commit timeline for this stream feature group.commitDetails(Integer limit) /** Retrieves commit timeline for this stream feature group.commitDetails(String wallclockTime) Return commit details as of specific point in time.commitDetails(String wallclockTime, Integer limit) Return commit details as of specific point in time.Recompute the statistics for the stream feature group and save them to the feature store.computeStatistics(String wallclockTime) Recompute the statistics for the feature group and save them to the feature store.voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) Incrementally insert data to a stream feature group or overwrite all data contained in the feature group.voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, boolean overwrite) Incrementally insert data to a stream feature group or overwrite all data contained in the feature group.voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, boolean overwrite, Map<String, String> writeOptions) Incrementally insert data to a stream feature group or overwrite all data contained in the feature group.voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, boolean overwrite, Map<String, String> writeOptions, JobConfiguration jobConfiguration) Incrementally insert data to a stream feature group or overwrite all data contained in the feature group.voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, HudiOperationType operation) voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, JobConfiguration jobConfiguration) Incrementally insert data to a stream feature group or overwrite all data contained in the feature group.voidvoidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage, boolean overwrite) voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage, boolean overwrite, HudiOperationType operation, Map<String, String> writeOptions) voidinsert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) Incrementally insert data to a stream feature group or overwrite all data contained in the feature group.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout, String checkpointLocation) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout, String checkpointLocation, Map<String, String> writeOptions) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout, String checkpointLocation, Map<String, String> writeOptions, JobConfiguration jobConfiguration) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, String checkpointLocation) org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, String checkpointLocation) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, Map<String, String> writeOptions) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.streaming.StreamingQueryinsertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) Ingest a Spark Structured Streaming Dataframe to the online feature store.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>read()Reads the feature group from the offline storage as Spark DataFrame.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>read(boolean online) Reads the stream feature group from the offline or online storage as Spark DataFrame.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>Reads the stream feature group from the offline or online storage as Spark DataFrame.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>Reads stream Feature group into a dataframe at a specific point in time.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>Reads stream Feature group into a dataframe at a specific point in time.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>Reads the stream feature group from the offline storage as Spark DataFrame.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>readChanges(String wallclockStartTime, String wallclockEndTime) Deprecated.`readChanges` method is deprecated.org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>Deprecated.voidremoveRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) Drops records present in the provided DataFrame and commits it as update to this Stream Feature group.voidremoveRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage) Drops records present in the provided DataFrame from the offline table, and optionally from the online store of an online-enabled feature group.voidremoveRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) Drops records present in the provided DataFrame and commits it as update to this Stream Feature group.voidremoveRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions, Storage storage) Drops records present in the provided DataFrame from the offline table, and optionally from the online store of an online-enabled feature group.voidsave(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) Deprecated.voidsave(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions, JobConfiguration jobConfiguration) Deprecated.Select a subset of features of the feature group and return a query object.Select all features of the feature group and return a query object.selectExcept(List<String> features) Select all features including primary key and event time feature of the feature group except provided `features` and return a query object.selectExceptFeatures(List<Feature> features) Select all features including primary key and event time feature of the feature group except provided `features` and return a query object.selectFeatures(List<Feature> features) Select a subset of features of the feature group and return a query object.voidshow(int numRows) Show the first `n` rows of the feature group.voidshow(int numRows, boolean online) Show the first `n` rows of the feature group.voidupdateFeatures(Feature feature) Update the metadata of feature.voidupdateFeatures(List<Feature> features) Update the metadata of multiple features.Methods inherited from class com.logicalclocks.hsfs.FeatureGroupBase
addTag, checkDeprecated, delete, deleteTag, getAvroSchema, getComplexFeatures, getDeserializedAvroSchema, getDeserializedEncodedAvroSchema, getEncodedAvroSchema, getFeature, getFeatureAvroSchema, getLatestOnlineIngestion, getOnlineIngestion, getPrimaryKeys, getSubject, getTag, getTags, setDataSource, setDeprecated, unloadSubject, updateDeprecated, updateDeprecated, updateDescription, updateFeatureDescription, updateNotificationTopicName, updateStatisticsConfig
-
Field Details
-
featureGroupEngine
-
-
Constructor Details
-
StreamFeatureGroup
public StreamFeatureGroup(FeatureStore featureStore, @NonNull @NonNull String name, Integer version, String description, List<String> primaryKeys, List<String> partitionKeys, String hudiPrecombineKey, boolean onlineEnabled, TimeTravelFormat timeTravelFormat, List<Feature> features, StatisticsConfig statisticsConfig, String onlineTopicName, String topicName, String notificationTopicName, String eventTime, OnlineConfig onlineConfig, StorageConnector storageConnector, String path) -
StreamFeatureGroup
public StreamFeatureGroup() -
StreamFeatureGroup
-
StreamFeatureGroup
-
-
Method Details
-
read
public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> read() throws FeatureStoreException, IOExceptionReads the feature group from the offline storage as Spark DataFrame.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // read feature group fg.read()- Returns:
- DataFrame.
- Throws:
FeatureStoreException- In case it cannot run read query on storage and/or no commit information was found for this feature group;IOException- Generic IO exception.
-
read
public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> read(boolean online) throws FeatureStoreException, IOException Reads the stream feature group from the offline or online storage as Spark DataFrame.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // read feature group data from online storage fg.read(true) // read feature group data from offline storage fg.read(false)- Parameters:
online- Set `online` to `true` to read from the online storage.- Returns:
- Spark DataFrame containing the feature data.
- Throws:
FeatureStoreException- In case it cannot run read query on storage and/or no commit information was found for this feature group;IOException- Generic IO exception.
-
read
public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> read(Map<String, String> readOptions) throws FeatureStoreException, IOExceptionReads the stream feature group from the offline storage as Spark DataFrame.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional read options (this example applies to HUDI enabled FGs) Map<String, String> readOptions = new HashMap<String, String>() {{ put("hoodie.datasource.read.end.instanttime", "20230401211015") }}; // read feature group data fg.read(readOptions)- Parameters:
readOptions- Additional read options as key/value pairs.- Returns:
- Spark DataFrame containing the feature data.
- Throws:
FeatureStoreException- In case it cannot run read query on storage and/or no commit information was found for this feature group.IOException- Generic IO exception.
-
read
public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> read(boolean online, Map<String, String> readOptions) throws FeatureStoreException, IOExceptionReads the stream feature group from the offline or online storage as Spark DataFrame.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional read options Map<String, String> readOptions = new HashMap<String, String>() {{ put("hoodie.datasource.read.end.instanttime", "20230401211015") }}; // read feature group data from offline storage fg.read(false, readOptions)- Parameters:
online- Set `online` to `true` to read from the online storage.readOptions- Additional read options as key/value pairs.- Returns:
- Spark DataFrame containing the feature data.
- Throws:
FeatureStoreException- In case it cannot run read query on storage and/or no commit information was found for this feature group;IOException- Generic IO exception.
-
read
public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> read(String wallclockTime) throws FeatureStoreException, IOException, ParseException Reads stream Feature group into a dataframe at a specific point in time.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // read feature group data as of specific point in time (Hudi commit timestamp). fg.read("20230205210923")- Parameters:
wallclockTime- Read data as of this point in time. Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.- Returns:
- Spark DataFrame containing feature data.
- Throws:
FeatureStoreException- In case it's unable to identify format of the provided wallclockTime date formatIOException- Generic IO exception.ParseException- In case it's unable to parse provided wallclockTime to date type.
-
read
public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> read(String wallclockTime, Map<String, String> readOptions) throws FeatureStoreException, IOException, ParseExceptionReads stream Feature group into a dataframe at a specific point in time.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional read options Map<String, String> readOptions = new HashMap<String, String>() {{ put("hoodie.datasource.read.end.instanttime", "20230401211015") }}; // read stream feature group data as of specific point in time (Hudi commit timestamp). fg.read("20230205210923", readOptions)- Parameters:
wallclockTime- Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.readOptions- Additional read options as key-value pairs.- Returns:
- Spark DataFrame containing feature data.
- Throws:
FeatureStoreException- In case it's unable to identify format of the provided wallclockTime date formatIOException- Generic IO exception.ParseException- In case it's unable to parse provided wallclockTime to date type.
-
show
Show the first `n` rows of the feature group.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // show top 5 lines of feature group data. fg.show(5);- Parameters:
numRows- Number of rows to show.- Throws:
FeatureStoreException- In case it cannot run read query on storage and/or no commit information was found for this feature group;IOException- Generic IO exception.
-
show
Show the first `n` rows of the feature group.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // show top 5 lines of feature data from online storage. fg.show(5, true);- Parameters:
numRows- Number of rows to show.online- If `true` read from online feature store.- Throws:
FeatureStoreException- In case it cannot run read query on storage and/or no commit information was found for this feature group;IOException- Generic IO exception.
-
readChanges
@Deprecated public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> readChanges(String wallclockStartTime, String wallclockEndTime) throws FeatureStoreException, IOException, ParseException Deprecated.`readChanges` method is deprecated. Use `asOf(wallclockEndTime, wallclockStartTime).read()` instead.Reads changes that occurred between specified points in time.- Parameters:
wallclockStartTime- start date.wallclockEndTime- end date.- Returns:
- DataFrame.
- Throws:
FeatureStoreException- FeatureStoreExceptionIOException- IOExceptionParseException- ParseException
-
readChanges
@Deprecated public org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> readChanges(String wallclockStartTime, String wallclockEndTime, Map<String, String> readOptions) throws FeatureStoreException, IOException, ParseExceptionDeprecated. -
asOf
Get Query object to retrieve all features of the group at a point in the past. This method selects all features in the feature group and returns a Query object at the specified point in time. This can then either be read into a Dataframe or used further to perform joins or construct a training dataset.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // get query object to retrieve stream feature group feature data as of // specific point in time (Hudi commit timestamp). fg.asOf("20230205210923")- Parameters:
wallclockTime- Read data as of this point in time. Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.- Returns:
- Query. The query object with the applied time travel condition
- Throws:
FeatureStoreException- In case it's unable to identify format of the provided wallclockTime date formatParseException- In case it's unable to parse provided wallclockTime to date type.
-
asOf
public Query asOf(String wallclockTime, String excludeUntil) throws FeatureStoreException, ParseException Get Query object to retrieve all features of the group at a point in the past. This method selects all features in the feature group and returns a Query object at the specified point in time. This can then either be read into a Dataframe or used further to perform joins or construct a training dataset.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // get query object to retrieve feature group feature data as of specific point in time "20230205210923" // but exclude commits until "20230204073411" (Hudi commit timestamp). fg.asOf("20230205210923", "20230204073411")- Parameters:
wallclockTime- Read data as of this point in time. Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.excludeUntil- Exclude commits until this point in time. Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.- Returns:
- Query. The query object with the applied time travel condition
- Throws:
FeatureStoreException- In case it's unable to identify format of the provided wallclockTime date formatParseException- In case it's unable to parse provided wallclockTime to date type.
-
save
@Deprecated public void save(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) throws FeatureStoreException, IOException, ParseExceptionDeprecated. -
save
@Deprecated public void save(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions, JobConfiguration jobConfiguration) throws FeatureStoreException, IOException, ParseExceptionDeprecated. -
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) throws FeatureStoreException, IOException, ParseException Incrementally insert data to a stream feature group or overwrite all data contained in the feature group. The `features` dataframe can be a Spark DataFrame or RDD. If the stream feature group doesn't exist, the insert method will create the necessary metadata the first time it is invoked and write the specified `features` dataframe as feature group to the online/offline feature store.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); //insert feature data fg.insert(featureData);- Parameters:
featureData- spark DataFrame, RDD. Features to be saved.- Throws:
IOException- Generic IO exception.FeatureStoreException- If client is not connected to Hopsworks; cannot run read query on storage and/or can't reconcile HUDI schema.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) throws FeatureStoreException, IOException, ParseExceptionIncrementally insert data to a stream feature group or overwrite all data contained in the feature group. The `features` dataframe can be a Spark DataFrame or RDD. If the stream feature group doesn't exist, the insert method will create the necessary metadata the first time it is invoked and write the specified `features` dataframe as feature group to the online/offline feature store.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // insert feature data fg.insert(featureData, writeOptions);- Parameters:
featureData- Spark DataFrame, RDD. Features to be saved.writeOptions- Additional write options as key-value pairs.- Throws:
IOException- Generic IO exception.FeatureStoreException- If client is not connected to Hopsworks; cannot run read query on storage and/or can't reconcile HUDI schema.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage) throws IOException, FeatureStoreException, ParseException -
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, boolean overwrite) throws IOException, FeatureStoreException, ParseException Incrementally insert data to a stream feature group or overwrite all data contained in the feature group. The `features` dataframe can be a Spark DataFrame or RDD. If the stream feature group doesn't exist, the insert method will create the necessary metadata the first time it is invoked and write the specified `features` dataframe as feature group to the online/offline feature store.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data and drop all data in the stream feature group before inserting new data fg.insert(featureData, true);- Parameters:
featureData- Spark DataFrame, RDD. Features to be saved.overwrite- Drop all data in the feature group before inserting new data. This does not affect metadata.- Throws:
IOException- Generic IO exception.FeatureStoreException- If client is not connected to Hopsworks; cannot run read query on storage and/or can't reconcile HUDI schema.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage, boolean overwrite) throws IOException, FeatureStoreException, ParseException -
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, boolean overwrite, Map<String, String> writeOptions) throws FeatureStoreException, IOException, ParseExceptionIncrementally insert data to a stream feature group or overwrite all data contained in the feature group. The `features` dataframe can be a Spark DataFrame or RDD. If the stream feature group doesn't exist, the insert method will create the necessary metadata the first time it is invoked and write the specified `features` dataframe as feature group to the online/offline feature store.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // insert feature data and drop all data in the stream feature group before inserting new data fg.insert(featureData, true, writeOptions);- Parameters:
featureData- Spark DataFrame, RDD. Features to be saved.overwrite- Drop all data in the feature group before inserting new data. This does not affect metadata.writeOptions- Additional write options as key-value pairs.- Throws:
IOException- Generic IO exception.FeatureStoreException- If client is not connected to Hopsworks; cannot run read query on storage and/or can't reconcile HUDI schema.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, HudiOperationType operation) throws FeatureStoreException, IOException, ParseException -
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage, boolean overwrite, HudiOperationType operation, Map<String, String> writeOptions) throws FeatureStoreException, IOException, ParseException -
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, JobConfiguration jobConfiguration) throws FeatureStoreException, IOException, ParseException Incrementally insert data to a stream feature group or overwrite all data contained in the feature group. The `features` dataframe can be a Spark DataFrame or RDD. If the stream feature group doesn't exist, the insert method will create the necessary metadata the first time it is invoked and write the specified `features` dataframe as feature group to the online/offline feature store.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // Define job configuration. JobConfiguration jobConfiguration = new JobConfiguration(); jobConfiguration.setDynamicAllocationEnabled(true); jobConfiguration.setDriverMemory(2048); // insert feature data fg.insert(featureData, jobConfiguration);- Parameters:
featureData- Spark DataFrame, RDD. Features to be saved.jobConfiguration- configure the Hopsworks Job used to write data into the stream feature group.- Throws:
IOException- Generic IO exception.FeatureStoreException- If client is not connected to Hopsworks; cannot run read query on storage and/or can't reconcile HUDI schema.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
insert
public void insert(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, boolean overwrite, Map<String, String> writeOptions, JobConfiguration jobConfiguration) throws FeatureStoreException, IOException, ParseExceptionIncrementally insert data to a stream feature group or overwrite all data contained in the feature group. The `features` dataframe can be a Spark DataFrame or RDD. If the stream feature group doesn't exist, the insert method will create the necessary metadata the first time it is invoked and write the specified `features` dataframe as feature group to the online/offline feature store.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // Define job configuration. JobConfiguration jobConfiguration = new JobConfiguration(); jobConfiguration.setDynamicAllocationEnabled(true); jobConfiguration.setDriverMemory(2048); // insert feature data fg.insert(featureData, false, writeOptions, jobConfiguration);- Parameters:
featureData- Spark DataFrame, RDD. Features to be saved.overwrite- Drop all data in the feature group before inserting new data. This does not affect metadata.writeOptions- Additional write options as key-value pairs.jobConfiguration- configure the Hopsworks Job used to write data into the stream feature group.- Throws:
IOException- Generic IO exception.FeatureStoreException- If client is not connected to Hopsworks; cannot run read query on storage and/or can't reconcile HUDI schema.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data fg.insertStream(featureData);- Specified by:
insertStreamin classFeatureGroupBase<org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>>- Parameters:
featureData- Features in Streaming Dataframe to be saved.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data fg.insertStream(featureData, queryName);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UI- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // insert feature data fg.insertStream(featureData, writeOptions);- Specified by:
insertStreamin classFeatureGroupBase<org.apache.spark.sql.Dataset<org.apache.spark.sql.Row>>- Parameters:
featureData- Features in Streaming Dataframe to be saved.writeOptions- Additional write options as key-value pairs.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, Map<String, String> writeOptions) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // insert feature data fg.insertStream(featureData, queryName, writeOptions);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIwriteOptions- Additional write options as key-value pairs.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data String queryName = "electricity_prices_streaming_query"; String outputMode = "append"; fg.insertStream(featureData, queryName, outputMode);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIoutputMode- Specifies how data of a streaming DataFrame/Dataset is written to a streaming sink. (1) `"append"`: Only the new rows in the streaming DataFrame/Dataset will be written to the sink. (2) `"complete"`: All the rows in the streaming DataFrame/Dataset will be written to the sink every time there is some update. (3) `"update"`: only the rows that were updated in the streaming DataFrame/Dataset will be written to the sink every time there are some updates. If the query doesn’t contain aggregations, it will be equivalent to append mode. Default behaviour is `"append"`.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, String checkpointLocation) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data String queryName = "electricity_prices_streaming_query"; String outputMode = "append"; String checkpointLocation = "path_to_checkpoint_dir"; fg.insertStream(featureData, queryName outputMode, checkpointLocation);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIoutputMode- Specifies how data of a streaming DataFrame/Dataset is written to a streaming sink. (1) `"append"`: Only the new rows in the streaming DataFrame/Dataset will be written to the sink. (2) `"complete"`: All the rows in the streaming DataFrame/Dataset will be written to the sink every time there is some update. (3) `"update"`: only the rows that were updated in the streaming DataFrame/Dataset will be written to the sink every time there are some updates. If the query doesn’t contain aggregations, it will be equivalent to append mode.checkpointLocation- Checkpoint directory location. This will be used to as a reference to from where to resume the streaming job.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data String queryName = "electricity_prices_streaming_query"; String outputMode = "append"; fg.insertStream(featureData, queryName, outputMode, outputMode, true, 1000);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIoutputMode- Specifies how data of a streaming DataFrame/Dataset is written to a streaming sink. (1) `"append"`: Only the new rows in the streaming DataFrame/Dataset will be written to the sink. (2) `"complete"`: All the rows in the streaming DataFrame/Dataset will be written to the sink every time there is some update. (3) `"update"`: only the rows that were updated in the streaming DataFrame/Dataset will be written to the sink every time there are some updates. If the query doesn’t contain aggregations, it will be equivalent to append mode.awaitTermination- Waits for the termination of this query, either by query.stop() or by an exception. If the query has terminated with an exception, then the exception will be thrown. If timeout is set, it returns whether the query has terminated or not within the timeout secondstimeout- Only relevant in combination with `awaitTermination=true`.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout, String checkpointLocation) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // insert feature data String queryName = "electricity_prices_streaming_query"; String outputMode = "append"; String checkpointLocation = "path_to_checkpoint_dir"; fg.insertStream(featureData, queryName, outputMode, outputMode, true, 1000, checkpointLocation);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIoutputMode- Specifies how data of a streaming DataFrame/Dataset is written to a streaming sink. (1) `"append"`: Only the new rows in the streaming DataFrame/Dataset will be written to the sink. (2) `"complete"`: All the rows in the streaming DataFrame/Dataset will be written to the sink every time there is some update. (3) `"update"`: only the rows that were updated in the streaming DataFrame/Dataset will be written to the sink every time there are some updates. If the query doesn’t contain aggregations, it will be equivalent to append mode.awaitTermination- Waits for the termination of this query, either by query.stop() or by an exception. If the query has terminated with an exception, then the exception will be thrown. If timeout is set, it returns whether the query has terminated or not within the timeout secondstimeout- Only relevant in combination with `awaitTermination=true`.checkpointLocation- Checkpoint directory location. This will be used to as a reference to from where to resume the streaming job.- Returns:
- Streaming Query object.
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout, String checkpointLocation, Map<String, String> writeOptions) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // insert feature data String queryName = "electricity_prices_streaming_query"; String outputMode = "append"; String checkpointLocation = "path_to_checkpoint_dir"; fg.insertStream(featureData, queryName, outputMode, outputMode, true, 1000, checkpointLocation, writeOptions);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIoutputMode- Specifies how data of a streaming DataFrame/Dataset is written to a streaming sink. (1) `"append"`: Only the new rows in the streaming DataFrame/Dataset will be written to the sink. (2) `"complete"`: All the rows in the streaming DataFrame/Dataset will be written to the sink every time there is some update. (3) `"update"`: only the rows that were updated in the streaming DataFrame/Dataset will be written to the sink every time there are some updates. If the query doesn’t contain aggregations, it will be equivalent to append mode.awaitTermination- Waits for the termination of this query, either by query.stop() or by an exception. If the query has terminated with an exception, then the exception will be thrown. If timeout is set, it returns whether the query has terminated or not within the timeout secondstimeout- Only relevant in combination with `awaitTermination=true`.checkpointLocation- Checkpoint directory location. This will be used to as a reference to from where to resume the streaming job.writeOptions- Additional write options as key-value pairs.- Returns:
- Streaming Query object.
-
insertStream
-
insertStream
public org.apache.spark.sql.streaming.StreamingQuery insertStream(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, String queryName, String outputMode, boolean awaitTermination, Long timeout, String checkpointLocation, Map<String, String> writeOptions, JobConfiguration jobConfiguration) Ingest a Spark Structured Streaming Dataframe to the online feature store. This method creates a long-running Spark Streaming Query, you can control the termination of the query through the arguments// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // Define job configuration. JobConfiguration jobConfiguration = new JobConfiguration(); jobConfiguration.setDynamicAllocationEnabled(true); jobConfiguration.setDriverMemory(2048); String queryName = "electricity_prices_streaming_query"; String outputMode = "append"; String checkpointLocation = "path_to_checkpoint_dir"; // insert feature data fg.insertStream(featureData, queryName, outputMode, outputMode, true, 1000, checkpointLocation, writeOptions, jobConfiguration);- Parameters:
featureData- Features in Streaming Dataframe to be saved.queryName- Specify a name for the query to make it easier to recognise in the Spark UIoutputMode- Specifies how data of a streaming DataFrame/Dataset is written to a streaming sink. (1) `"append"`: Only the new rows in the streaming DataFrame/Dataset will be written to the sink. (2) `"complete"`: All the rows in the streaming DataFrame/Dataset will be written to the sink every time there is some update. (3) `"update"`: only the rows that were updated in the streaming DataFrame/Dataset will be written to the sink every time there are some updates. If the query doesn’t contain aggregations, it will be equivalent to append mode.awaitTermination- Waits for the termination of this query, either by query.stop() or by an exception. If the query has terminated with an exception, then the exception will be thrown. If timeout is set, it returns whether the query has terminated or not within the timeout secondstimeout- Only relevant in combination with `awaitTermination=true`.checkpointLocation- Checkpoint directory location. This will be used to as a reference to from where to resume the streaming job.writeOptions- Additional write options as key-value pairs.jobConfiguration- configure the Hopsworks Job used to write data into the stream feature group.- Returns:
- Streaming Query object.
-
selectFeatures
Select a subset of features of the feature group and return a query object. The query can be used to construct joins of feature groups or create a feature view with a subset of features of the feature group.- Parameters:
features- List of Feature meta data objects.- Returns:
- Query object.
-
select
Select a subset of features of the feature group and return a query object. The query can be used to construct joins of feature groups or create a feature view with a subset of features of the feature group.- Parameters:
features- List of Feature names.- Returns:
- Query object.
-
selectAll
Select all features of the feature group and return a query object. The query can be used to construct joins of feature groups or create a feature view with a subset of features of the feature group.- Returns:
- Query object.
-
selectExceptFeatures
Select all features including primary key and event time feature of the feature group except provided `features` and return a query object. The query can be used to construct joins of feature groups or create a feature view with a subset of features of the feature group.- Parameters:
features- List of Feature meta data objects.- Returns:
- Query object.
-
selectExcept
Select all features including primary key and event time feature of the feature group except provided `features` and return a query object. The query can be used to construct joins of feature groups or create a feature view with a subset of features of the feature group.- Parameters:
features- List of Feature names.- Returns:
- Query object.
-
removeRows
public void removeRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) throws FeatureStoreException, IOException, ParseException Drops records present in the provided DataFrame and commits it as update to this Stream Feature group.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // drop records of feature data and commit fg.removeRows(featureData);When the feature group is online-enabled the records are deleted from the online store as well. Use
removeRows(Dataset, Storage)withStorage.OFFLINEto delete offline only.- Parameters:
featureData- Spark DataFrame, RDD. Feature data to be deleted.- Throws:
FeatureStoreException- If Client is not connected to Hopsworks and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
removeRows
public void removeRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) throws FeatureStoreException, IOException, ParseExceptionDrops records present in the provided DataFrame and commits it as update to this Stream Feature group.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // define additional write options Map<String, String> writeOptions = = new HashMap<String, String>() {{ put("hoodie.bulkinsert.shuffle.parallelism", "5"); put("hoodie.insert.shuffle.parallelism", "5"); put("hoodie.upsert.shuffle.parallelism", "5");} }; // drop records of feature data and commit fg.removeRows(featureData, writeOptions);When the feature group is online-enabled the records are deleted from the online store as well. Use
removeRows(Dataset, Map, Storage)withStorage.OFFLINEto delete offline only.- Parameters:
featureData- Spark DataFrame, RDD. Feature data to be deleted.writeOptions- Additional write options as key-value pairs.- Throws:
FeatureStoreException- If Client is not connected to Hopsworks and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
removeRows
public void removeRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Storage storage) throws FeatureStoreException, IOException, ParseException Drops records present in the provided DataFrame from the offline table, and optionally from the online store of an online-enabled feature group.featureDataneeds to carry the key columns the offline delete matches on: the primary key, plus the event time and any partition columns when the feature group has them. The online delete matches on the primary key alone and ignores every other column, so a value passed for a non-key feature has no effect on it.A stream feature group's inserts reach the offline table through the materialization job. Deleting a row whose insert has not been materialized yet deletes nothing offline, and the materialization job then writes the row, so it stays in the offline table while the online store has it deleted. Re-run the delete after the materialization job to reconcile the two stores.
- Parameters:
featureData- Spark DataFrame, RDD. Feature data to be deleted.storage- The storage to delete from. Null follows the feature group: both stores when it is online-enabled, offline alone when it is not.Storage.OFFLINEdeletes from the offline table only andStorage.ONLINEfrom the online store only, as on insert. A single-store delete is not reconciled later: the rows stay in the other store until they are deleted there too.- Throws:
FeatureStoreException- If storage isStorage.ONLINEand the feature group is not online-enabled; or on other client/commit errors.IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
removeRows
public void removeRows(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions, Storage storage) throws FeatureStoreException, IOException, ParseExceptionDrops records present in the provided DataFrame from the offline table, and optionally from the online store of an online-enabled feature group.featureDataneeds to carry the key columns the offline delete matches on: the primary key, plus the event time and any partition columns when the feature group has them. The online delete matches on the primary key alone and ignores every other column, so a value passed for a non-key feature has no effect on it.A stream feature group's inserts reach the offline table through the materialization job. Deleting a row whose insert has not been materialized yet deletes nothing offline, and the materialization job then writes the row, so it stays in the offline table while the online store has it deleted. Re-run the delete after the materialization job to reconcile the two stores.
- Parameters:
featureData- Spark DataFrame, RDD. Feature data to be deleted.writeOptions- Additional write options as key-value pairs.storage- The storage to delete from. Null follows the feature group: both stores when it is online-enabled, offline alone when it is not.Storage.OFFLINEdeletes from the offline table only andStorage.ONLINEfrom the online store only, as on insert. A single-store delete is not reconciled later: the rows stay in the other store until they are deleted there too.- Throws:
FeatureStoreException- If storage isStorage.ONLINEand the feature group is not online-enabled; or on other client/commit errors.IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
commitDeleteRecord
@Deprecated public void commitDeleteRecord(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData) throws FeatureStoreException, IOException, ParseException Deprecated.useremoveRows(Dataset)instead.Drops records present in the provided DataFrame.- Parameters:
featureData- Spark DataFrame, RDD. Feature data to be deleted.- Throws:
FeatureStoreException- on client/commit errors.IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
commitDeleteRecord
@Deprecated public void commitDeleteRecord(org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> featureData, Map<String, String> writeOptions) throws FeatureStoreException, IOException, ParseExceptionDeprecated.useremoveRows(Dataset, Map)instead.Drops records present in the provided DataFrame.- Parameters:
featureData- Spark DataFrame, RDD. Feature data to be deleted.writeOptions- Additional write options as key-value pairs.- Throws:
FeatureStoreException- on client/commit errors.IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
commitDetails
public Map<Long,Map<String, commitDetails() throws IOException, FeatureStoreException, ParseExceptionString>> Retrieves commit timeline for this stream feature group.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // get commit timeline fg.commitDetails();- Returns:
- commit details.
- Throws:
FeatureStoreException- If Client is not connected to Hopsworks and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
commitDetails
public Map<Long,Map<String, commitDetailsString>> (Integer limit) throws IOException, FeatureStoreException, ParseException /** Retrieves commit timeline for this stream feature group.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // get latest 10 commit details fg.commitDetails(10);- Parameters:
limit- number of commits to return.- Returns:
- commit details.
- Throws:
FeatureStoreException- If Client is not connected to Hopsworks and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
commitDetails
public Map<Long,Map<String, commitDetailsString>> (String wallclockTime) throws IOException, FeatureStoreException, ParseException Return commit details as of specific point in time.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); //get commit details as of 20230206 fg.commitDetails("20230206");- Parameters:
wallclockTime- Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.- Throws:
FeatureStoreException- If Client is not connected to Hopsworks and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
commitDetails
public Map<Long,Map<String, commitDetailsString>> (String wallclockTime, Integer limit) throws IOException, FeatureStoreException, ParseException Return commit details as of specific point in time.// get feature store handle FeatureStore fs = HopsworksConnection.builder().build().getFeatureStore(); // get feature group handle StreamFeatureGroup fg = fs.getStreamFeatureGroup("electricity_prices", 1); // get top 10 commit details as of 20230206 fg.commitDetails("20230206", 10);- Parameters:
wallclockTime- Datetime string. The String should be formatted in one of the following formats `yyyyMMdd`, `yyyyMMddHH`, `yyyyMMddHHmm`, or `yyyyMMddHHmmss`.limit- number of commits to return.- Returns:
- commit details.
- Throws:
FeatureStoreException- If Client is not connected to Hopsworks and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse HUDI commit date string to date type.
-
updateFeatures
public void updateFeatures(List<Feature> features) throws FeatureStoreException, IOException, ParseException Update the metadata of multiple features. Currently only feature description updates are supported.- Parameters:
features- List of Feature metadata objects- Throws:
FeatureStoreException- If Client is not connected to Hopsworks, unable to identify date format and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse date string to date type.
-
updateFeatures
public void updateFeatures(Feature feature) throws FeatureStoreException, IOException, ParseException Update the metadata of feature. Currently only feature description updates are supported.- Parameters:
feature- Feature metadata object- Throws:
FeatureStoreException- If Client is not connected to Hopsworks, unable to identify date format and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse date string to date type.
-
appendFeatures
public void appendFeatures(List<Feature> features) throws FeatureStoreException, IOException, ParseException Append features to the schema of the stream feature group. It is only possible to append features to a feature group. Removing features is considered a breaking change.- Parameters:
features- list of Feature metadata objects- Throws:
FeatureStoreException- If Client is not connected to Hopsworks, unable to identify date format and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse date string to date type.
-
appendFeatures
public void appendFeatures(Feature features) throws FeatureStoreException, IOException, ParseException Append a single feature to the schema of the stream feature group. It is only possible to append features to a feature group. Removing features is considered a breaking change.- Parameters:
features- List of Feature metadata objects- Throws:
FeatureStoreException- If Client is not connected to Hopsworks, unable to identify date format and/or no commit information was found for this feature group;IOException- Generic IO exception.ParseException- In case it's unable to parse date string to date type.
-
computeStatistics
Recompute the statistics for the stream feature group and save them to the feature store.- Returns:
- statistics object of computed statistics
- Throws:
FeatureStoreException- If Client is not connected to Hopsworks,IOException- Generic IO exception.
-
computeStatistics
public Statistics computeStatistics(String wallclockTime) throws FeatureStoreException, IOException, ParseException Recompute the statistics for the feature group and save them to the feature store.- Parameters:
wallclockTime- number of commits to return.- Returns:
- statistics object of computed statistics
- Throws:
FeatureStoreExceptionIOExceptionParseException
-
getStatistics
- Throws:
FeatureStoreExceptionIOException
-
removeRows(Dataset)instead.