Complete AWS SageMaker Feature Store Tutorial: Feature Management for ML
Amazon SageMaker Feature Store is a fully managed repository for storing, sharing, and managing ML features. It provides a centralized store for features that can be used across training and inference, ensuring consistency and reusability.
Why Feature Store?
Key Benefits:- Consistency: Same features for training and inference
- Reusability: Share features across teams and models
- Versioning: Track feature changes over time
- Low latency: Real-time feature serving
- Offline storage: Historical data for training
- Feature Groups
- Online Store (real-time)
- Offline Store (batch/training)
- Feature Definitions
Prerequisites
pip install sagemaker boto3 pandas
SageMaker SDK >= 2.0
python -c "import sagemaker; print(sagemaker.version)"
Quick Start
1. Setup
import boto3
import sagemaker
from sagemaker.featurestore.featuregroup import FeatureGroup
import pandas as pd
import time
session = sagemaker.Session()
bucket = session.defaultbucket()
role = sagemaker.getexecutionrole()
region = session.botoregionname
featurestoresession = sagemaker.Session()
2. Prepare Data
import pandas as pd
import numpy as np
from datetime import datetime
Create sample customer data
customerdata = pd.DataFrame({
"customerid": [f"C{i:04d}" for i in range(1, 101)],
"age": np.random.randint(18, 70, 100),
"income": np.random.randint(30000, 150000, 100),
"creditscore": np.random.randint(300, 850, 100),
"accountbalance": np.random.uniform(0, 50000, 100).round(2),
"numproducts": np.random.randint(1, 5, 100),
"isactive": np.random.choice([0, 1], 100)
})
Add event time (required for Feature Store)
customerdata["eventtime"] = datetime.now().strftime("%Y-%m-%dT%H:%M:%SZ")
print(customerdata.head())
Creating Feature Groups
1. Define Feature Group
from sagemaker.featurestore.featuredefinition import (
FeatureDefinition,
FeatureTypeEnum
)
Feature group name
featuregroupname = "customer-features"
Create feature group
customerfeaturegroup = FeatureGroup(
name=featuregroupname,
sagemakersession=featurestoresession
)
Load feature definitions from DataFrame
customerfeaturegroup.loadfeaturedefinitions(dataframe=customerdata)
Or define manually
featuredefinitions = [
FeatureDefinition(featurename="customerid", featuretype=FeatureTypeEnum.STRING),
FeatureDefinition(featurename="age", featuretype=FeatureTypeEnum.INTEGRAL),
FeatureDefinition(featurename="income", featuretype=FeatureTypeEnum.INTEGRAL),
FeatureDefinition(featurename="creditscore", featuretype=FeatureTypeEnum.INTEGRAL),
FeatureDefinition(featurename="accountbalance", featuretype=FeatureTypeEnum.FRACTIONAL),
FeatureDefinition(featurename="numproducts", featuretype=FeatureTypeEnum.INTEGRAL),
FeatureDefinition(featurename="isactive", featuretype=FeatureTypeEnum.INTEGRAL),
FeatureDefinition(featurename="eventtime", featuretype=FeatureTypeEnum.STRING)
]
2. Create Feature Group
# Create feature group with online and offline stores
customerfeaturegroup.create(
s3uri=f"s3://{bucket}/feature-store/",
recordidentifiername="customerid",
eventtimefeaturename="eventtime",
rolearn=role,
enableonlinestore=True,
description="Customer features for churn prediction"
)
Wait for feature group to be created
status = customerfeaturegroup.describe().get("FeatureGroupStatus")
while status == "Creating":
print(f"Status: {status}")
time.sleep(5)
status = customerfeaturegroup.describe().get("FeatureGroupStatus")
print(f"Feature group status: {status}")
3. Ingest Features
# Ingest data into feature store
customerfeaturegroup.ingest(
dataframe=customerdata,
maxworkers=4,
wait=True
)
print("Feature ingestion complete!")
Online Store (Real-time)
1. Get Single Record
# Get single record by identifier
record = customerfeaturegroup.getrecord(
recordidentifiervalueasstring="C0001"
)
print("Record for C0001:")
for feature in record:
print(f" {feature['FeatureName']}: {feature['ValueAsString']}")
2. Batch Get Records
from sagemaker.featurestore.featurestore import FeatureStore
feature
store = FeatureStore(sagemakersession=featurestoresession)
Batch get records
identifiers = [
{"FeatureGroupName": feature
groupname, "RecordIdentifiersValueAsString": ["C0001", "C0002", "C0003"]}
]
response = feature
store.batchgetrecord(identifiers=identifiers)
for record in response["Records"]:
customerid = next(f["ValueAsString"] for f in record["Record"] if f["FeatureName"] == "customerid")
print(f"Customer: {customerid}")
3. Real-time Inference Integration
import json
def getfeaturesforinference(customerid):
"""Get features for real-time inference."""
record = customerfeaturegroup.getrecord(
recordidentifiervalueasstring=customerid
)
# Convert to feature vector
features = {f["FeatureName"]: f["ValueAsString"] for f in record}
return [
float(features["age"]),
float(features["income"]),
float(features["creditscore"]),
float(features["accountbalance"]),
float(features["numproducts"])
]
Use in inference
features = getfeaturesforinference("C0001")
print(f"Features for inference: {features}")
Offline Store (Training)
1. Query with Athena
from sagemaker.featurestore.featuregroup import AthenaQuery
Create Athena query
athena
query = customerfeaturegroup.athenaquery()
Get table name
table
name = athenaquery.tablename
databasename = athenaquery.database
print(f"Database: {databasename}")
print(f"Table: {tablename}")
Run query
querystring = f"""
SELECT *
FROM "{tablename}"
WHERE isactive = 1
LIMIT 100
"""
athenaquery.run(
querystring=querystring,
outputlocation=f"s3://{bucket}/athena-results/"
)
Wait and get results
athenaquery.wait()
df = athenaquery.asdataframe()
print(df.head())
2. Join Multiple Feature Groups
# Create another feature group for transactions
transactiondata = pd.DataFrame({
"customerid": [f"C{i:04d}" for i in np.random.randint(1, 101, 500)],
"transactionid": [f"T{i:06d}" for i in range(1, 501)],
"amount": np.random.uniform(10, 1000, 500).round(2),
"transactiontype": np.random.choice(["purchase", "refund", "transfer"], 500),
"eventtime": datetime.now().strftime("%Y-%m-%dT%H:%M:%SZ")
})
Create and ingest transaction feature group
transactionfeaturegroup = FeatureGroup(
name="transaction-features",
sagemakersession=featurestoresession
)
transactionfeaturegroup.loadfeaturedefinitions(dataframe=transactiondata)
transactionfeaturegroup.create(
s3uri=f"s3://{bucket}/feature-store/",
recordidentifiername="transactionid",
eventtimefeaturename="eventtime",
rolearn=role,
enableonlinestore=True
)
Wait for creation
time.sleep(60)
transactionfeaturegroup.ingest(dataframe=transactiondata, wait=True)
# Join query
joinquery = f"""
SELECT
c.customerid,
c.age,
c.income,
c.creditscore,
t.amount,
t.transactiontype
FROM "{customerfeaturegroup.athenaquery().tablename}" c
JOIN "{transactionfeaturegroup.athenaquery().tablename}" t
ON c.customerid = t.customerid
WHERE c.isactive = 1
"""
athenaquery.run(
querystring=joinquery,
outputlocation=f"s3://{bucket}/athena-results/"
)
athenaquery.wait()
trainingdf = athenaquery.asdataframe()
3. Create Training Dataset
from sagemaker.featurestore.datasetbuilder import DatasetBuilder
Build training dataset
builder = DatasetBuilder(
sagemakersession=featurestoresession,
base=customerfeaturegroup,
outputpath=f"s3://{bucket}/training-data/"
)
Add point-in-time join
builder = builder.withfeaturegroup(
featuregroup=transactionfeaturegroup,
targetfeaturenameinbase="customerid",
includedfeaturenames=["amount", "transactiontype"]
)
Build dataset
trainingdataset = builder.todataframe()
print(f"Training dataset shape: {trainingdataset.shape}")
Feature Updates
1. Update Records
# Update existing customer data
updateddata = pd.DataFrame({
"customerid": ["C0001"],
"age": [35],
"income": [85000],
"creditscore": [750],
"accountbalance": [25000.00],
"numproducts": [3],
"isactive": [1],
"eventtime": datetime.now().strftime("%Y-%m-%dT%H:%M:%SZ")
})
Ingest updated data
customerfeaturegroup.ingest(
dataframe=updateddata,
maxworkers=1,
wait=True
)
print("Record updated!")
Verify update
record = customerfeaturegroup.getrecord(
recordidentifiervalueasstring="C0001"
)
print("Updated record:", record)
2. Delete Records
# Delete record (soft delete with deletion marker)
customerfeaturegroup.deleterecord(
recordidentifiervalueasstring="C0099",
eventtime=datetime.now().strftime("%Y-%m-%dT%H:%M:%SZ")
)
Feature Store SDK
1. Feature Processor
from sagemaker.featurestore.featureprocessor import (
FeatureProcessor,
CSVDataSource,
FeatureGroupDataSource
)
Create feature processor
@feature
processor(
inputs=[
CSVDataSource(s3uri=f"s3://{bucket}/raw-data/")
],
output=featuregroupname
)
def processfeatures(inputdf):
# Feature engineering
df = inputdf.copy()
# Add derived features
df["incomeperproduct"] = df["income"] / df["numproducts"]
df["creditincomeratio"] = df["creditscore"] / (df["income"] / 10000)
return df
2. Scheduled Ingestion
from sagemaker.featurestore.featureprocessor import (
FeatureProcessorPipelineEvents,
topipeline
)
Create pipeline for scheduled ingestion
pipeline = topipeline(
pipelinename="feature-ingestion-pipeline",
step=processfeatures,
role=role
)
Schedule execution
pipeline.puttriggers(
triggers=[
FeatureProcessorPipelineEvents(
eventpattern={
"source": ["aws.s3"],
"detail-type": ["Object Created"],
"detail": {
"bucket": {"name": [bucket]},
"object": {"key": [{"prefix": "raw-data/"}]}
}
}
)
]
)
Monitoring and Management
1. Describe Feature Group
# Get feature group details
description = customerfeaturegroup.describe()
print(f"Name: {description['FeatureGroupName']}")
print(f"Status: {description['FeatureGroupStatus']}")
print(f"Record Identifier: {description['RecordIdentifierFeatureName']}")
print(f"Event Time: {description['EventTimeFeatureName']}")
List feature definitions
for feature in description['FeatureDefinitions']:
print(f" {feature['FeatureName']}: {feature['FeatureType']}")
2. List Feature Groups
# List all feature groups
sagemakerclient = boto3.client("sagemaker")
response = sagemakerclient.listfeaturegroups()
for fg in response["FeatureGroupSummaries"]:
print(f"{fg['FeatureGroupName']}: {fg['FeatureGroupStatus']}")
3. Delete Feature Group
# Delete feature group
customerfeaturegroup.delete()
Note: Offline store data in S3 is not deleted automatically
Best Practices
1. Feature Naming Convention
# Good: Descriptive, consistent naming
featuredefinitions = [
FeatureDefinition("customerid", FeatureTypeEnum.STRING),
FeatureDefinition("customerageyears", FeatureTypeEnum.INTEGRAL),
FeatureDefinition("customertotalspendusd", FeatureTypeEnum.FRACTIONAL),
FeatureDefinition("customerispremiumflag", FeatureTypeEnum.INTEGRAL)
]
2. Event Time Management
from datetime import datetime, timezone
def geteventtime():
"""Generate consistent event time."""
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
Always include eventtime in your data
data["eventtime"] = geteventtime()
3. Batch Ingestion
def batchingest(featuregroup, df, batchsize=1000):
"""Ingest data in batches for better reliability."""
for i in range(0, len(df), batch
size):
batch = df.iloc[i:i+batchsize]
featuregroup.ingest(
dataframe=batch,
maxworkers=4,
wait=True
)
print(f"Ingested batch {i//batch_size + 1}")
Conclusion
SageMaker Feature Store provides:
Key takeaways:
- Design features with reusability in mind
- Use online store for real-time inference
- Leverage Athena for offline training queries
- Implement proper event time management
- Monitor feature freshness and quality