Use o SDK do FeatureStore para Python para registrar visualizações de features, gravar dados de features nos armazenamentos offline e online e gerar conjuntos de dados de treinamento para modelos.
Este guia usa o conjunto de dados de código aberto Moviedata como exemplo prático. As tabelas Movie, User e Rating correspondem, respectivamente, às tabelas de item, usuário e rótulo em um pipeline típico de recomendação.
Neste guia, você vai:
Instale o SDK e configure seu projeto
Defina entidades de features e visualizações de features
Gravar dados no armazenamento offline e publicá-los no armazenamento online
Recuperar features online e defina um seletor de features
Crie um conjunto de dados de treinamento e exportá-lo para treinar o modelo
Pré-requisitos
Antes de começar, verifique se você:
Obteve o AccessKey ID e o AccessKey Secret da sua conta Alibaba Cloud
Configure as fontes de dados no FeatureStore (offline e online)
Para obter a melhor experiência, execute os códigos deste guia em uma instância DSW.
Etapa 1: Instale o SDK e configure o projeto
Instale o SDK
Execute o seguinte comando em um ambiente Python 3:
pip install https://feature-store-py.oss-cn-beijing.aliyuncs.com/package/feature_store_py-1.3.1-py3-none-any.whl
Defina variáveis de ambiente
Armazene suas credenciais como variáveis de ambiente para evitar codificar informações sensíveis diretamente no código. No seu DSW Notebook, clique em Terminal na barra de menu superior e execute:
echo "export AccessKeyID='<your-access-key-id>'" >> ~/.bashrc
echo "export AccessKeySecret='<your-access-key-secret>'" >> ~/.bashrc
source ~/.bashrc
Substitua <your-access-key-id> e <your-access-key-secret> pelo seu AccessKey ID e AccessKey Secret reais.
Importar módulos
import unittest
import sys
import os
from os.path import dirname, join, abspath
from feature_store_py.fs_client import FeatureStoreClient
from feature_store_py.fs_project import FeatureStoreProject
from feature_store_py.fs_datasource import UrlDataSource, MaxComputeDataSource, DatahubDataSource, HologresDataSource, SparkDataSource, LabelInput, TrainingSetOutput
from feature_store_py.fs_type import FSTYPE
from feature_store_py.fs_schema import OpenSchema, OpenField
from feature_store_py.fs_feature_view import FeatureView
from feature_store_py.fs_features import FeatureSelector
from feature_store_py.fs_config import EASDeployConfig, LabelInputConfig, PartitionConfig, FeatureViewConfig, TrainSetOutputConfig, SequenceFeatureConfig, SequenceTableConfig
import logging
logger = logging.getLogger("foo")
logger.addHandler(logging.StreamHandler(stream=sys.stdout))
Conectar a um projeto do FeatureStore
Inicialize o cliente e conecte-se ao seu projeto. Você pode criar vários projetos independentes no FeatureStore. Este guia usa um projeto chamado fs_movie.
# Load credentials from environment variables
access_id = os.getenv("AccessKeyID")
access_ak = os.getenv("AccessKeySecret")
# Set the region where your FeatureStore is activated
region = 'cn-hangzhou'
fs = FeatureStoreClient(access_key_id=access_id, access_key_secret=access_ak, region=region)
cur_project_name = "fs_movie"
project = fs.get_project(cur_project_name)
if project is None:
raise ValueError("Project not found. Create the project first: fs_movie")
Para verificar a conexão, imprima os detalhes do projeto:
project = fs.get_project(cur_project_name)
print(project)
Etapa 2: Defina entidades de features
Uma entidade de feature agrupa features semanticamente relacionadas. Cada entidade possui um ID de junção que vincula features em múltiplas visualizações de features. As visualizações podem usar nomes de chave primária diferentes, mas se unem por meio desse ID de entidade compartilhado.
O exemplo do Moviedata usa três entidades:
# Movie entity
cur_entity_name_movie = "movie_data"
entity_movie = project.get_entity(cur_entity_name_movie)
if entity_movie is None:
entity_movie = project.create_entity(name=cur_entity_name_movie, join_id='movie_id')
entity_movie.print_summary()
# User entity
cur_entity_name_user = "user_data"
entity_user = project.get_entity(cur_entity_name_user)
if entity_user is None:
entity_user = project.create_entity(name=cur_entity_name_user, join_id='user_md5')
entity_user.print_summary()
# Rating entity
cur_entity_name_ratings = "rating_data"
entity_ratings = project.get_entity(cur_entity_name_ratings)
if entity_ratings is None:
entity_ratings = project.create_entity(name=cur_entity_name_ratings, join_id='rating_id')
entity_ratings.print_summary()
Etapa 3: Crie visualizações de features e gravar dados
As visualizações de features definem como os dados externos entram no FeatureStore: a fonte de dados, o esquema, o local de armazenamento e os metadados das features. O FeatureStore oferece suporte a três tipos de visualização de features:
|
Tipo |
Caso de uso |
Caminho de gravação |
|
|
Features offline ou T-1 dia |
Gravar no armazenamento offline → publicar no armazenamento online |
|
|
Features em tempo real |
Gravar diretamente no armazenamento online |
|
|
Sequências de comportamento do usuário |
Gravar offline → ler online em tempo real |
O tempo de vida (TTL) controla por quanto tempo o armazenamento online retém os dados. O valor -1 retém todos os dados; um valor positivo mantém apenas os dados dentro do período especificado.
BatchFeatureView
O BatchFeatureView gerencia features offline ou de T-1 dia. O fluxo de dados ocorre por meio de duas operações distintas:
write_table()— grava dados no armazenamento offline do MaxCompute. Esta ação registra os dados apenas no armazenamento offline; eles ainda não ficam disponíveis para consultas online.publish_table()— sincroniza os dados do armazenamento offline para o online, tornando-os disponíveis para recuperação em tempo real.
Visualização de features de filme
# Load the movie CSV from a public URL
path = 'https://feature-store-test.oss-cn-beijing.aliyuncs.com/dataset/moviedata_all/movies.csv'
ds = UrlDataSource(path, delimiter=',', omit_header=True)
# Define the schema
movie_schema = OpenSchema(
OpenField(name='movie_id', type='STRING'),
OpenField(name='name', type='STRING'),
OpenField(name='alias', type='STRING'),
OpenField(name='actors', type='STRING'),
OpenField(name='cover', type='STRING'),
OpenField(name='directors', type='STRING'),
OpenField(name='double_score', type='STRING'),
OpenField(name='double_votes', type='STRING'),
OpenField(name='genres', type='STRING'),
OpenField(name='imdb_id', type='STRING'),
OpenField(name='languages', type='STRING'),
OpenField(name='mins', type='STRING'),
OpenField(name='official_site', type='STRING'),
OpenField(name='regions', type='STRING'),
OpenField(name='release_date', type='STRING'),
OpenField(name='slug', type='STRING'),
OpenField(name='story', type='STRING'),
OpenField(name='tags', type='STRING'),
OpenField(name='year', type='STRING'),
OpenField(name='actor_ids', type='STRING'),
OpenField(name='director_ids', type='STRING'),
OpenField(name='dt', type='STRING')
)
# Create the feature view (registers metadata only — data is not written yet)
feature_view_movie_name = "feature_view_movie"
batch_feature_view = project.get_feature_view(feature_view_movie_name)
if batch_feature_view is None:
batch_feature_view = project.create_batch_feature_view(
name=feature_view_movie_name,
schema=movie_schema,
online=True,
entity=cur_entity_name_movie,
primary_key='movie_id',
partitions=['dt'],
ttl=-1
)
# Write data to the offline store
cur_task = batch_feature_view.write_table(ds, partitions={'dt': '20220830'})
cur_task.wait()
print(cur_task.task_summary)
# Publish to the online store
cur_task = batch_feature_view.publish_table({'dt': '20220830'})
cur_task.wait()
print(cur_task.task_summary)
# Verify
batch_feature_view = project.get_feature_view(feature_view_movie_name)
batch_feature_view.print_summary()
Visualizações de features de usuário e avaliação
As tabelas de usuário e avaliação seguem o mesmo padrão: carregue a fonte, defina o esquema, crie a visualização de features, grave e publique:
# User feature view
users_path = 'https://feature-store-test.oss-cn-beijing.aliyuncs.com/dataset/moviedata_all/users.csv'
ds = UrlDataSource(users_path, delimiter=',', omit_header=True)
user_schema = OpenSchema(
OpenField(name='user_md5', type='STRING'),
OpenField(name='user_nickname', type='STRING'),
OpenField(name='ds', type='STRING')
)
feature_view_user_name = "feature_view_users"
batch_feature_view = project.get_feature_view(feature_view_user_name)
if batch_feature_view is None:
batch_feature_view = project.create_batch_feature_view(
name=feature_view_user_name,
schema=user_schema,
online=True,
entity=cur_entity_name_user,
primary_key='user_md5',
partitions=['ds'],
ttl=-1
)
write_table_task = batch_feature_view.write_table(ds, {'ds': '20220830'})
write_table_task.wait()
print(write_table_task.task_summary)
cur_task = batch_feature_view.publish_table({'ds': '20220830'})
cur_task.wait()
print(cur_task.task_summary)
batch_feature_view = project.get_feature_view(feature_view_user_name)
batch_feature_view.print_summary()
# Rating feature view
ratings_path = 'https://feature-store-test.oss-cn-beijing.aliyuncs.com/dataset/moviedata_all/ratings.csv'
ds = UrlDataSource(ratings_path, delimiter=',', omit_header=True)
ratings_schema = OpenSchema(
OpenField(name='rating_id', type='STRING'),
OpenField(name='user_md5', type='STRING'),
OpenField(name='movie_id', type='STRING'),
OpenField(name='rating', type='STRING'),
OpenField(name='rating_time', type='STRING'),
OpenField(name='dt', type='STRING')
)
feature_view_rating_name = "feature_view_ratings"
batch_feature_view = project.get_feature_view(feature_view_rating_name)
if batch_feature_view is None:
batch_feature_view = project.create_batch_feature_view(
name=feature_view_rating_name,
schema=ratings_schema,
online=True,
entity=cur_entity_name_ratings,
primary_key='rating_id',
event_time='rating_time',
partitions=['dt']
)
cur_task = batch_feature_view.write_table(ds, {'dt': '20220831'})
cur_task.wait()
print(cur_task.task_summary)
batch_feature_view = project.get_feature_view(feature_view_rating_name)
batch_feature_view.print_summary()
StreamFeatureView
O StreamFeatureView lida com features em tempo real. Os dados são gravados diretamente no armazenamento online.
Primeiro, crie a tabela de dados de teste no MaxCompute ou DataWorks:
CREATE TABLE IF NOT EXISTS online_stream_test_t1 (
id STRING COMMENT 'ID',
count_value BIGINT COMMENT 'Count value',
metric_value DOUBLE COMMENT 'Metric value'
)
PARTITIONED BY (
ds string COMMENT 'Data timestamp'
)
LIFECYCLE 365;
INSERT INTO TABLE online_stream_test_t1 PARTITION (ds='20250815')
SELECT
CONCAT('str_', CAST(id AS STRING)) AS id,
CAST(FLOOR(RAND() * 1000000) AS BIGINT) AS count_value,
ROUND(RAND() * 1000, 2) AS metric_value
FROM (
SELECT SEQUENCE(1, 1000) AS id_list
) tmp
LATERAL VIEW EXPLODE(id_list) table_tmp AS id;
Após a execução bem-sucedida do SQL, a tabela online_stream_test_t1 será criada com dados na partição ds=20250815.
Em seguida, crie e publique o StreamFeatureView:
online_schema = OpenSchema(
OpenField(name='id', type='STRING'),
OpenField(name='count_value', type='INT64'),
OpenField(name='metric_value', type='DOUBLE')
)
feature_view_rating_name_stream = "feature_view_online_stream"
stream_feature_view = project.get_feature_view(feature_view_rating_name_stream)
if stream_feature_view is None:
stream_feature_view = project.create_stream_feature_view(
name=feature_view_rating_name_stream,
schema=online_schema,
online=True,
entity=cur_entity_name_user,
primary_key='id',
event_time='count_value'
)
stream_feature_view = project.get_feature_view(feature_view_rating_name_stream)
stream_feature_view.print_summary()
Em umStreamFeatureView, o campoevent_timelimpa dados expirados quando configurado. Para mais detalhes, consulte Ciclo de vida de features em tempo real .
Sincronize os dados com o armazenamento online:
# Replace offline_datasource_id with your FeatureStore project's offline store ID.
# table_name is the offline feature table to push to the online store.
stream_task = stream_feature_view.publish_table(
partitions={'ds': '20250815'},
mode='Merge',
offline_to_online=True,
publish_config={
'offline_datasource_id': project.offline_datasource_id,
'table_name': 'online_stream_test_t1'
}
)
stream_task.wait()
print(stream_task.task_summary)
Sequence FeatureView
O Sequence FeatureView armazena sequências de comportamento do usuário para treinamento offline e recuperação online em tempo real.
Primeiro, copie os dados de sequência da fonte de pai_online_project (acesso de leitura público) para o seu próprio projeto:
CREATE TABLE IF NOT EXISTS rec_sln_demo_behavior_table_preprocess_sequence_wide_seq_feature_v3
LIKE pai_online_project.rec_sln_demo_behavior_table_preprocess_sequence_wide_seq_feature_v3
STORED AS ALIORC
LIFECYCLE 90;
INSERT OVERWRITE TABLE rec_sln_demo_behavior_table_preprocess_sequence_wide_seq_feature_v3 PARTITION(ds)
SELECT *
FROM pai_online_project.rec_sln_demo_behavior_table_preprocess_sequence_wide_seq_feature_v3
WHERE ds >= '20231022' AND ds <= '20231024';
Após a execução bem-sucedida do SQL, a tabela de features de sequência será criada com dados das partições ds=20231022, ds=20231023 e ds=20231024.
Crie a visualização de features de sequência:
user_entity_name = "user"
seq_feature_view_name = "wide_seq_feature_v3"
seq_feature_view = project.get_feature_view(seq_feature_view_name)
if seq_feature_view is None:
seq_table_name = "rec_sln_demo_behavior_table_preprocess_sequence_wide_seq_feature_v3"
behavior_table_name = 'rec_sln_demo_behavior_table_preprocess_v3'
ds = MaxComputeDataSource(project.offline_datasource_id, behavior_table_name)
event_time = 'event_unix_time' # Event time field in the behavior table
item_id = 'item_id' # Item ID field in the behavior table
event = 'event' # Event type field in the behavior table
# deduplication_method=1: deduplicates on ['user_id', 'item_id', 'event']
# deduplication_method=2: deduplicates on ['user_id', 'item_id', 'event', 'event_time']
sequence_feature_config_list = [
SequenceFeatureConfig(
offline_seq_name='click_seq_50_seq', # Field name in the offline sequence table
seq_event='click', # Event type to filter
online_seq_name='click_seq_50', # Name exposed to the online Go SDK
seq_len=50 # Maximum sequence length; longer sequences are truncated
)
]
seq_table_config = SequenceTableConfig(
table_name=seq_table_name,
primary_key='user_id',
event_time='event_unix_time'
)
seq_feature_view = project.create_sequence_feature_view(
seq_feature_view_name,
datasource=ds,
event_time=event_time,
item_id=item_id,
event=event,
deduplication_method=1,
sequence_feature_config=sequence_feature_config_list,
sequence_table_config=seq_table_config,
entity=user_entity_name
)
seq_feature_view.print_summary()
Sincronize os dados com o armazenamento online:
seq_task = seq_feature_view.publish_table({'ds': '20231023'}, days_to_load=30)
seq_task.wait()
seq_task.print_summary()
Registre a tabela de rótulos para futura geração de conjuntos de treinamento:
label_table_name = 'fs_movie_feature_view_ratings_offline'
ds = MaxComputeDataSource(data_source_id=project.offline_datasource_id, table=label_table_name)
label_table = project.get_label_table(label_table_name)
if label_table is None:
label_table = project.create_label_table(datasource=ds, event_time='rating_time')
Etapa 4: Recuperar features online e defina seletores de features
Recuperar features online
Recupere features diretamente de uma visualização de features. Atualmente, o FeatureStore prioriza o FeatureDB para recuperação de features online.
feature_view_movie_name = "feature_view_movie"
batch_feature_view = project.get_feature_view(feature_view_movie_name)
# Retrieve features for a single item
ret_features_1 = batch_feature_view.list_feature_view_online_features(join_ids=['26357307'])
print("ret_features1 = ", ret_features_1)
# Retrieve features for multiple items in one call
ret_features_2 = batch_feature_view.list_feature_view_online_features(join_ids=['30444960', '3317352'])
print("ret_features2 = ", ret_features_2)
Defina um seletor de features
Um FeatureSelector especifica quais features extrair de uma visualização ao gerar um conjunto de dados de treinamento ou executar inferência. Três padrões de seleção têm suporte:
feature_view_name = 'feature_view_movie'
# Select specific features by name
feature_selector = FeatureSelector(feature_view_name, ['site_id', 'site_category'])
# Select all features
feature_selector = FeatureSelector(feature_view_name, '*')
# Select specific features and apply an alias
feature_selector = FeatureSelector(
feature_view='user1',
features=['f1', 'f2', 'f3'],
alias={"f1": "f1_1"} # Expose f1 as f1_1 in the output
)
Etapa 5: Crie e exportar um conjunto de dados de treinamento
Crie um conjunto de dados de treinamento
Um conjunto de dados de treinamento (tabela de amostras) combina uma tabela de rótulos com features de uma ou mais visualizações, unidas por chaves primárias usando junções pontuais no tempo para evitar vazamento de dados.
label_table_name = 'fs_movie_feature_view_ratings_offline'
output_ds = MaxComputeDataSource(data_source_id=project.offline_datasource_id)
train_set_output = TrainingSetOutput(output_ds)
# Select features from movie and user views
feature_movie_selector = FeatureSelector('feature_view_movie', ['name', 'actors', 'regions', 'tags'])
feature_user_selector = FeatureSelector('feature_view_users', ['user_nickname'])
train_set = project.create_training_set(
label_table_name=label_table_name,
train_set_output=train_set_output,
feature_selectors=[feature_movie_selector, feature_user_selector]
)
print("train_set = ", train_set)
Registrar um modelo
Crie uma entrada de modelo vinculada ao conjunto de dados de treinamento:
model_name = "fs_rank_v1"
cur_model = project.get_model(model_name)
if cur_model is None:
cur_model = project.create_model(model_name, train_set)
print("cur_model_train_set_table_name = ", cur_model.train_set_table_name)
Exportar o conjunto de dados de treinamento
Especifique as partições e o tempo de evento para a tabela de rótulos e cada visualização de features e execute a exportação:
# Label table configuration
label_partitions = PartitionConfig(name='dt', value='20220831')
label_input_config = LabelInputConfig(
partition_config=label_partitions,
event_time='1999-01-00 00:00:00'
)
# Feature view configurations
movie_partitions = PartitionConfig(name='dt', value='20220830')
feature_view_movie_config = FeatureViewConfig(name='feature_view_movie', partition_config=movie_partitions)
user_partitions = PartitionConfig(name='ds', value='20220830')
feature_view_user_config = FeatureViewConfig(name='feature_view_users', partition_config=user_partitions)
feature_view_config_list = [feature_view_movie_config, feature_view_user_config]
# Output partition configuration
train_set_partitions = PartitionConfig(name='dt', value='20220831')
train_set_output_config = TrainSetOutputConfig(partition_config=train_set_partitions)
# Run the export
task = cur_model.export_train_set(label_input_config, feature_view_config_list, train_set_output_config)
task.wait()
print(task.summary)
Conceitos principais
Armazenamento offline
O armazenamento offline é um data warehouse para features históricas. As features são gravadas no MaxCompute ou no Hadoop Distributed File System (HDFS) usando Apache Spark. Esse armazenamento atende a duas finalidades: gerar conjuntos de dados para treinamento de modelos e fornecer features para previsões em lote.
Armazenamento online
O armazenamento online é um data warehouse para features em tempo real, oferecendo acesso de baixa latência para inferência online. O FeatureStore oferece suporte a FeatureDB, Hologres e Tablestore como armazenamentos online.
Tipos de visualização de features
|
Tipo |
Descrição |
|
|
Features offline ou T-1 dia. Os dados offline são gravados no armazenamento offline e podem ser publicados no armazenamento online para consultas em tempo real. |
|
|
Features em tempo real. Os dados são gravados diretamente no armazenamento online. |
|
|
Features de sequência de comportamento do usuário, com suporte para gravação offline e leitura online em tempo real. |
Próximos passos
Visão geral do FeatureDB — conheça o armazenamento online que sustenta a recuperação de features em tempo real
Configure itens do FeatureStore — crie e gerencie múltiplos projetos
EasyRec no GitHub — integre o FeatureStore com o EasyRec para geração de features (FG) e treinamento de modelos
FeatureGenerator no GitHub — gere features para modelos de recomendação
Caso encontre problemas ao usar o FeatureStore, participe do grupo de suporte no DingTalk (ID do grupo: 32260796).