A engenharia de features do FeatureStore é utilizada em cenários de recomendação, publicidade, controle de riscos e aprendizado de máquina para reduzir a complexidade da engenharia de features. Funções padronizadas viabilizam a engenharia de features por meio de configurações simples.
Pré-requisitos
Certifique-se de que os preparativos descritos na tabela a seguir estejam concluídos.
service | Operação |
PAI |
|
MaxCompute |
|
DataWorks |
|
1. Preparativos
Preparar os dados brutos
Este tópico utiliza quatro tabelas source:
Tabela de usuários (rec_sln_demo_user_table_preprocess_v1): contém features básicas de usuários, como gênero, idade, cidade e número de seguidores.
Tabela de comportamento (rec_sln_demo_behavior_table_preprocess_v1): contém dados comportamentais, como cliques de usuários em itens em horários específicos.
Tabela de itens (rec_sln_demo_item_table_preprocess_v1): contém features básicas de itens, como categoria, autor, número acumulado de cliques e número acumulado de curtidas.
Tabela ampla de comportamento (rec_sln_demo_behavior_table_preprocess_wide_v1): formada pelo join das três tabelas anteriores.
As tabelas de dados estão armazenadas no workspace pai_online_project, visível para todos os usuários. Essas tabelas contêm apenas dados de simulação. É necessário executar instruções SQL no DataWorks para sincronizar os dados das tabelas mencionadas do workspace pai_online_project para o seu projeto MaxCompute. Siga os passos abaixo:
Faça login no console do DataWorks.
No painel de navegação à esquerda, clique em Data Development and O&M > Data Development.
Selecione o workspace do DataWorks criado e clique em Go to Data Studio.
-
Passe o cursor sobre Create e escolha Create Node > MaxCompute > ODPS SQL. Na página exibida, configure os parâmetros do nó.
Parâmetro
Valor sugerido
Node Type
ODPS SQL
Path
Business Flow/Workflow/MaxCompute
Name
Insira um nome personalizado.
Clique em Confirm.
-
No editor do novo nó, execute os comandos SQL a seguir para sincronizar os dados das tabelas de usuário, item, comportamento e comportamento ampla do projeto pai_online_project para o seu projeto MaxCompute. Para o grupo de recursos, selecione o seu grupo de recursos exclusivo.
-
Execute as instruções SQL a seguir para sincronizar os dados da tabela de usuários rec_sln_demo_user_table_preprocess_v1:
CREATE TABLE IF NOT EXISTS rec_sln_demo_user_table_preprocess_v1 like pai_online_project.rec_sln_demo_user_table_preprocess_v1 STORED AS ALIORC LIFECYCLE 90; INSERT OVERWRITE TABLE rec_sln_demo_user_table_preprocess_v1 PARTITION(ds) SELECT * FROM pai_online_project.rec_sln_demo_user_table_preprocess_v1 WHERE ds >= '20240530' and ds <='20240605'; -
Execute as instruções SQL a seguir para sincronizar os dados da tabela de comportamento rec_sln_demo_behavior_table_preprocess_v1:
CREATE TABLE IF NOT EXISTS rec_sln_demo_behavior_table_preprocess_v1 like pai_online_project.rec_sln_demo_behavior_table_preprocess_v1 STORED AS ALIORC LIFECYCLE 90; INSERT OVERWRITE TABLE rec_sln_demo_behavior_table_preprocess_v1 PARTITION(ds) SELECT * FROM pai_online_project.rec_sln_demo_behavior_table_preprocess_v1 WHERE ds >= '20240530' and ds <='20240605'; -
Execute as instruções SQL a seguir para sincronizar os dados da tabela de itens rec_sln_demo_item_table_preprocess_v1:
CREATE TABLE IF NOT EXISTS rec_sln_demo_item_table_preprocess_v1 like pai_online_project.rec_sln_demo_item_table_preprocess_v1 STORED AS ALIORC LIFECYCLE 90; INSERT OVERWRITE TABLE rec_sln_demo_item_table_preprocess_v1 PARTITION(ds) SELECT * FROM pai_online_project.rec_sln_demo_item_table_preprocess_v1 WHERE ds >= '20240530' and ds <='20240605'; -
Sincronize a tabela ampla de comportamento: rec_sln_demo_behavior_table_preprocess_wide_v1
CREATE TABLE IF NOT EXISTS rec_sln_demo_behavior_table_preprocess_wide_v1 like pai_online_project.rec_sln_demo_behavior_table_preprocess_wide_v1 STORED AS ALIORC LIFECYCLE 90; INSERT OVERWRITE TABLE rec_sln_demo_behavior_table_preprocess_wide_v1 PARTITION(ds) SELECT * FROM pai_online_project.rec_sln_demo_behavior_table_preprocess_wide_v1 WHERE ds >= '20240530' and ds <='20240605';
-
Instale o SDK do FeatureStore
Execute o comando a seguir em um Jupyter Notebook (Python 3):
%pip install https://feature-store-py.oss-cn-beijing.aliyuncs.com/package/feature_store_py-2.0.2-py3-none-any.whl
Em seguida, importe os módulos necessários ao longo deste guia:
import os
from feature_store_py import FeatureStoreClient
from feature_store_py.fs_datasource import MaxComputeDataSource
from feature_store_py.feature_engineering import (
TableTransform, Condition, DayOf, ComboTransform, Feature,
AggregationTransform, auto_count_feature_transform,
WindowTransform, auto_window_feature_transform
)
2. Processo de transformação de tabelas e features
Execute o código a seguir em um ambiente Jupyter Notebook.
-
Defina a transformação de tabela.
-
Inicialize o cliente.
access_key_id=os.environ.get ("ALIBABA_CLOUD_ACCESS_KEY_ID") # Enter your AccessKey ID. access_key_secret=os.environ.get ("ALIBABA_CLOUD_ACCESS_KEY_SECRET") # Enter your AccessKey secret. project='project_name' # Enter your project name. region='cn-hangzhou' # Enter the region in which your project resides. For example, if your project resides in China (Hangzhou), enter cn-hangzhou. fs_client = FeatureStoreClient(access_key_id=access_key_id, access_key_secret=access_key_secret, region=region) -
Especifique a source de dados.
input_table_name = "rec_sln_demo_behavior_table_preprocess_v1" ds = MaxComputeDataSource(table=input_table_name, project=project) -
Especifique o nome da tabela de saída para a transformação.
output_table_name = "rec_sln_demo_v1_fs_test_v1" -
Defina a transformação de tabela.
trans_name = "drop_duplicates" # Name of the table transformation. keys = ["user_id", "item_id"] # Fields for deduplication. sort_keys = ["event_unix_time"] # Sort fields. sort_order = ["desc"] # Order definition. tran_i = TableTransform(trans_name, keys, sort_keys, sort_order)
-
-
Defina a transformação de features.
feature1 = Feature( name="page_net_type", input=['page', 'net_type'], transform=ComboTransform( separator='_' ) ) feature2 = Feature( name="trim_playtime", type="double", transform="playtime/10" ) -
Gere o pipeline.
pipeline = fs_client.create_pipeline(ds, output_table_name).add_table_transform(tran_i).add_feature_transform([feature1, feature2], keep_input_columns=True) -
Gere e execute a transformação.
execute_date = '20240605' output_table = pipeline.execute(execute_date, drop_table=True)O código acima envolve duas etapas:
Gerar as configurações de transformação. Essas configurações definem as instruções SQL e as informações necessárias para a transformação, como entradas, saídas, parâmetros e dependências.
Executar a transformação. O sistema executa a transformação com base nas configurações da etapa anterior e armazena os resultados na tabela de saída.
-
Visualize os resultados.
-
Confira os resultados na tabela gerada. Os resultados são renderizados diretamente no formato pandas.DataFrame.
pd_ret = output_table.to_pandas(execute_date, limit=20) -
Exiba o conteúdo de pd_ret.
pd_ret -
Consulte as configurações geradas. Elas incluem a definição da tabela de entrada, o SQL de transformação, as dependências, os parâmetros e a definição da tabela de saída. Após salvar, essas configurações podem ser usadas para depuração e tarefas rotineiras online subsequentes.
transform_info = output_table.transform_info -
Exiba o conteúdo de transform_info.
transform_info -
Consulte as configurações de entrada do primeiro estágio.
pipeline_config = pipeline.pipeline_config -
Exiba o conteúdo de pipeline_config.
pipeline_config
-
3. Transformação de features estatísticas
Features estatísticas são um método comum de pré-processamento de dados em aprendizado de máquina e análise de dados, usado para gerar features mais representativas e interpretáveis. Essas transformações resumem, calculam e extraem informações dos dados brutos, permitindo que os modelos compreendam melhor tendências temporais, periodicidade e anomalias. As vantagens são as seguintes:
Captura de tendências temporais: em dados comportamentais de usuários, comportamentos recentes geralmente têm maior impacto no estado atual.
Redução de ruído: dados brutos podem conter ruído. Transformações estatísticas utilizam operações de agregação para minimizar esse impacto.
Enriquecimento de features: transformações estatísticas geram novas features, aumentando o poder expressivo do modelo.
Melhoria no desempenho do modelo: features estatísticas podem melhorar significativamente o desempenho de previsão.
Maior interpretabilidade: features estatísticas são mais fáceis de interpretar, o que simplifica o diagnóstico e a análise de problemas.
Compressão de dados: features estatísticas reduzem efetivamente a dimensionalidade dos dados.
Embora a implementação de features estatísticas seja complexa, é possível criar diversas features estatísticas utilizando as definições simples descritas a seguir.
Definição e execução de uma transformação de feature estatística individual
-
Defina os nomes das tabelas de entrada e saída.
input_agg_table_name = "rec_sln_demo_behavior_table_preprocess_wide_v1" ds_agg = MaxComputeDataSource(table=input_agg_table_name, project=project) output_agg_table_name = "rec_sln_demo_behavior_test_agg_v1" -
Defina uma feature estatística.
feature_agg1 = Feature( name="user_avg_praise_count_1d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", # Name of the aggregate function. Optional values are 'avg', 'sum', 'min', 'max'. condition=Condition(field="event", value="expr", operator="<>"), # Condition: the "event" field is not equal to "expr". group_by_keys="user_id", # The key corresponding to group by. window_size=DayOf(1), # The window size, which is 1 day here. ), ) -
Crie um pipeline e execute a transformação de feature estatística.
agg_pipeline = fs_client.create_pipeline(ds_agg, output_agg_table_name).add_feature_transform([feature_agg1]) -
Gere e execute a transformação.
execute_date = '20240605' print("transform_info = ", agg_pipeline.transform_info) output_agg_table = agg_pipeline.execute(execute_date, drop_table=True) -
Exiba o conteúdo de transform_info.
transform_info_agg = output_agg_table.transform_info transform_info_agg -
Visualize os resultados.
pd_ret = output_agg_table.to_pandas(execute_date, limit=20) pd_ret
JOIN automático para transformações de features estatísticas com janelas diferentes
-
Defina o nome da tabela de saída.
output_agg_table_name_2 = "rec_sln_demo_behavior_test_agg_v2" -
Defina features estatísticas para múltiplas janelas diferentes. Os tamanhos de janela definidos são 1, 3, 7, 15 e 30 dias.
feature_agg1 = Feature( name="user_avg_praise_count_1d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(1), ), ) feature_agg2 = Feature( name="user_avg_praise_count_3d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(3), ), ) feature_agg3 = Feature( name="user_avg_praise_count_7d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(7), ), ) feature_agg4 = Feature( name="user_avg_praise_count_15d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(15), ), ) feature_agg5 = Feature( name="user_avg_praise_count_30d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(30), ), ) -
Crie um pipeline.
agg_pipeline_2 = fs_client.create_pipeline(ds_agg, output_agg_table_name_2).add_feature_transform([feature_agg1, feature_agg2, feature_agg3, feature_agg4, feature_agg5]) -
Gere e execute o pipeline.
execute_date = '20240605' output_agg_table_2 = agg_pipeline_2.execute(execute_date, drop_table=True) -
Consulte os resultados da transformação.
transform_info_agg_2 = output_agg_table_2.transform_info transform_info_agg_2 -
Visualize os resultados da execução da tabela.
pd_ret_2 = output_agg_table_2.to_pandas(execute_date, limit=20) pd_ret_2
Mesclagem automática e derivação de tipo para múltiplos processos de transformação de features estatísticas
Para otimizar os cálculos, features com o mesmo tamanho de janela são automaticamente mescladas e calculadas no mesmo bloco de janela de grupo. O processo de cálculo envolve alterações de tipo. Por exemplo, avg transforma o tipo bigint para o tipo double. Como é difícil lembrar os tipos de todas as features de entrada, o processo de transformação de features estatísticas suporta derivação automática de tipo. O tipo da feature resultante é derivado automaticamente durante a definição da feature.
-
Defina o nome da tabela de saída.
output_agg_table_name_3 = "rec_sln_demo_behavior_test_agg_v3" -
Defina mais features de diferentes tipos.
feature_agg6 = Feature( name="user_expr_cnt_1d", transform=AggregationTransform( agg_func="count", condition=Condition(field="event", value="expr", operator="="), group_by_keys="user_id", window_size=DayOf(1), ) ) feature_agg7 = Feature( name="user_expr_item_id_dcnt_1d", input=['item_id'], transform=AggregationTransform( agg_func="count", condition=Condition(field="event", value="expr", operator="="), group_by_keys="user_id", window_size=DayOf(1), ), ) feature_agg8 = Feature( name="user_sum_praise_count_1d", input=["praise_count"], transform=AggregationTransform( agg_func="sum", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(1), ), ) feature_agg9 = Feature( name="user_sum_praise_count_3d", input=["praise_count"], transform=AggregationTransform( agg_func="sum", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(3), ), ) -
Crie um pipeline.
agg_pipeline_3 = fs_client.create_pipeline(ds_agg, output_agg_table_name_3).add_feature_transform([feature_agg1, feature_agg2, feature_agg3, feature_agg4, feature_agg5, feature_agg6, feature_agg7, feature_agg8, feature_agg9]) -
Gere e execute o pipeline.
execute_date = '20240605' output_agg_table_3 = agg_pipeline_3.execute(execute_date, drop_table=True) -
Consulte os resultados da transformação.
transform_info_agg_3 = output_agg_table_3.transform_info transform_info_agg_3 -
Visualize os resultados da execução da tabela.
pd_ret_3 = output_agg_table_3.to_pandas(execute_date, limit=20) pd_ret_3
Funções de extensão automática integradas para transformações de features estatísticas
Implementar manualmente cada feature estatística é complexo devido ao grande número de features, que inclui diferentes tamanhos de janela e inúmeras combinações de funções de agregação. O sistema oferece funções de extensão automática integradas. Basta especificar as features de entrada a serem contabilizadas, e o sistema gera automaticamente as definições para centenas de features estatísticas.
-
Especifique as features de entrada a serem contabilizadas.
name_prefix = "user_" input_list = ["playtime", "duration", "click_count", "praise_count"] event_name = 'event' event_type = 'expr' group_by_key = "user_id" count_feature_list = auto_count_feature_transform(name_prefix, input_list, event_name, event_type, group_by_key) print("len_count_feature_list = ", len(count_feature_list)) print("count_feature_list = ", count_feature_list) -
Defina o nome da tabela de saída e crie um pipeline.
output_agg_table_name_4 = "rec_sln_demo_behavior_test_agg_v4" agg_pipeline_4 =fs_client.create_pipeline(ds_agg, output_agg_table_name_4).add_feature_transform(count_feature_list) -
Gere e execute o pipeline.
execute_date = '20240605' output_agg_table_4 = agg_pipeline_4.execute(execute_date, drop_table=True) -
Consulte os resultados da transformação.
transform_info_agg_4 = output_agg_table_4.transform_info transform_info_agg_4 -
Visualize os resultados da execução da tabela.
pd_ret_4 = output_agg_table_4.to_pandas(execute_date, limit=20) pd_ret_4
Suporte a transformação simultânea de group keys diferentes
As seções anteriores descrevem o método de processamento quando todas as group keys são iguais. O sistema também suporta operações de transformação em group keys diferentes. Veja o exemplo a seguir:
-
Defina o nome da tabela de saída.
input_agg_table_name = "rec_sln_demo_behavior_table_preprocess_wide_v1" ds_agg = MaxComputeDataSource(table=input_agg_table_name, project=project) output_agg_table_name_5 = "rec_sln_demo_behavior_test_agg_v5" -
Defina features para group keys diferentes.
feature_agg1 = Feature( name="item__sum_follow_cnt_15d", input=['follow_cnt'], transform=AggregationTransform( agg_func="sum", condition=Condition(field="event", value="expr", operator="="), group_by_keys="item_id", window_size=DayOf(1), ) ) feature_agg2 = Feature( name="author__max_follow_cnt_15d", input=['follow_cnt'], transform=AggregationTransform( agg_func="max", condition=Condition(field="event", value="expr", operator="="), group_by_keys="author", window_size=DayOf(15), ), ) -
Crie um pipeline.
agg_pipeline_5 = fs_client.create_pipeline(ds_agg, output_agg_table_name_5).add_feature_transform([feature_agg1, feature_agg2]) -
Gere e execute o pipeline.
execute_date = '20240605' output_agg_table_5 = agg_pipeline_5.execute(execute_date, drop_table=True) -
Consulte os resultados da transformação.
transform_info_agg_5 = output_agg_table_5.transform_info transform_info_agg_5
4. Transformação de features WindowTransform
A transformação de features estatísticas descrita anteriormente é suficiente para cenários comuns de engenharia de features. No entanto, em cenários de recomendação em larga escala, existem requisitos mais avançados. O FeatureStore suporta a transformação de features WindowTransform, que permite obter features no formato KV e utilizar tabelas intermediárias diárias para otimizar o processo de cálculo. Isso reduz o tempo de cálculo de features e economiza custos computacionais. As vantagens são as seguintes:
Captura de interações não lineares complexas: features simples (como idade e gênero do usuário) não conseguem expressar preferências complexas do usuário. O cruzamento de features captura relacionamentos de interação não linear mais complexos entre usuários e itens.
Melhoria na precisão de previsão: features cruzadas podem melhorar significativamente o desempenho de sistemas de recomendação e de publicidade.
Redução do espaço de armazenamento: para grandes conjuntos de usuários e itens, armazenar diretamente as features de interação de cada par usuário-item não é viável. A extração e a transformação de features reduzem efetivamente o número de features a armazenar.
Melhoria na eficiência de inferência: ao pré-calcular e armazenar features cruzadas, é possível recuperá-las rapidamente durante a inferência em tempo real, o que melhora a velocidade de resposta.
A transformação de features WindowTransform é apresentada nas seções a seguir:
Processo de cálculo de funções de agregação simples
Funções de agregação simples incluem count, sum, max e min. O processo de cálculo dessas funções é relativamente direto. Após uma sumarização em nível diário, uma sumarização adicional ao longo de múltiplos dias é realizada para obter o resultado final. Esta seção também apresenta tabelas intermediárias diárias e User-Defined Functions (UDFs) (MaxCompute UDF overview) para obter o resultado final do cálculo. No processo real de cálculo, a execução da engenharia de features funciona da mesma forma que a transformação convencional descrita anteriormente. O SDK Python do FeatureStore gerencia automaticamente a criação de tabelas intermediárias diárias, a geração de UDFs, o upload de recursos e o registro automático de funções.
-
Defina os nomes das tabelas de entrada e saída.
input_window_table_name = "rec_sln_demo_behavior_table_preprocess_wide_v1" ds_window_1 = MaxComputeDataSource(table=input_agg_table_name, project=project) output_window_table_name_1 = "rec_sln_demo_behavior_test_window_v1" -
Defina as features WindowTransform.
win_feature1 = Feature( name="item__kv_gender_click_7d", input=["gender"], transform=WindowTransform( agg_func="count", condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(7), ), ) win_feature2 = Feature( name="item__kv_gender_click_cnt_7d", input=["gender"], transform=WindowTransform( agg_func="sum", # Aggregate function. Optional values are 'sum', 'avg', 'max', 'min'. agg_field="click_count", # Perform the aggregate function calculation on this feature. condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(7), ), ) -
Crie um pipeline e execute a transformação de features WindowTransform.
window_pipeline_1 = fs_client.create_pipeline(ds_window_1, output_window_table_name_1).add_feature_transform([win_feature1, win_feature2], keep_input_columns=True) -
Gere e execute a transformação.
Este processo de geração cria uma tabela intermediária temporária diária. Neste ponto,
DROP TABLEexclua apenas o resultado final, sem excluir a tabela temporária intermediária.execute_date = '20240605' print("transform_info = ", window_pipeline_1.transform_info) output_window_table_1 = window_pipeline_1.execute(execute_date, drop_table=True)Além disso, como estatísticas de múltiplos dias estão envolvidas (por exemplo, o exemplo anterior contabiliza dados de 7 dias), a tabela temporária intermediária normalmente calcula dados apenas para a partição mais recente. Por esse motivo, o sistema fornece o parâmetro
backfill_partitions. Na primeira execução, defina este parâmetro comoTruepara que o sistema preencha automaticamente os dados dos dias dependentes. Por exemplo, se a contagem envolve 7 dias de dados, o sistema completa automaticamente os 7 dias. Nas execuções rotineiras subsequentes, defina o parâmetro comoFalsepara completar apenas os dados da partição do dia mais recente.execute_date = '20240506' output_window_table_1 = window_pipeline_1.execute(execute_date, backfill_partitions=True)Quando o parâmetro
backfill_partitionsé definido comoTrue, o sistema completa automaticamente os dados dos dias dependentes da tabela intermediária temporária. Faça isso na primeira execução rotineira.Se o número de dias a serem contabilizados for grande, a execução do código acima pode demorar bastante.
-
Visualize os dados da tabela de destino.
window_ret_1 = output_window_table_1.to_pandas(execute_date, limit=50) window_ret_1 -
Consulte o processo real de cálculo.
window_pipeline_1.transform_infoComo é possível observar no processo de cálculo, o sistema gera uma tabela intermediária temporária diária chamada
rec_sln_demo_behavior_table_preprocess_wide_v1_tmp_daily. Essa tabela resume os resultados diários e os armazena em uma partição fixa, evitando cálculos repetidos.Além disso, uma UDF chamada
count_kvcalcula o resultado final. Essa UDF classifica e resume automaticamente os resultados estatísticos em um mapa de resultado, armazenado em formato string. Isso facilita a inferência de resultados subsequentes, tanto offline quanto online.
O conteúdo acima apresenta o processo de cálculo para funções de agregação simples, usando count e sum como exemplos. Embora esse processo envolva conceitos como tabelas intermediárias temporárias diárias e UDFs, o fluxo principal é o mesmo de uma transformação convencional de dados e não adiciona complexidade operacional. Outras funções de agregação simples, como max e min, seguem o mesmo princípio.
Processo de cálculo da função de agregação avg
Como calcular a média dos resultados de médias diárias leva a cálculos imprecisos, a função de agregação avg possui um processo de cálculo específico. O método correto é calcular primeiro a soma total (sum_v) e a contagem total (count_v) para todo o período e então obter a média pela fórmula sum_v/count_v.
Embora essa função de agregação esteja documentada separadamente, seus detalhes complexos de cálculo estão encapsulados em transform_info. É possível usar essa função como uma feature convencional para gerar o resultado final.
-
Defina os nomes das tabelas de entrada e saída.
input_window_table_name = "rec_sln_demo_behavior_table_preprocess_wide_v1" ds_window_1 = MaxComputeDataSource(table=input_window_table_name, project=project) output_window_table_name_2 = "rec_sln_demo_behavior_test_window_v2" -
Defina a feature WindowTransform.
win_feature1 = Feature( name="item__kv_gender_click_avg_7d", input=["gender"], transform=WindowTransform( agg_func="avg", agg_field="click_count", condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(7), ), ) win_feature2 = Feature( name="item__kv_gender_click_avg_15d", input=["gender"], transform=WindowTransform( agg_func="avg", agg_field="click_count", condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(15), ), ) -
Crie um pipeline e execute a transformação de features WindowTransform.
window_pipeline_2 = fs_client.create_pipeline(ds_window_1, output_window_table_name_2).add_feature_transform([win_feature1, win_feature2]) -
Gere e execute a transformação.
execute_date = '20240605' print("transform_info = ", window_pipeline_2.transform_info) output_window_table_2 = window_pipeline_2.execute(execute_date, drop_table=True) -
Visualize os dados da tabela de destino.
window_ret_2 = output_window_table_2.to_pandas(execute_date, limit=50) window_ret_2
Processo de cálculo de funções para múltiplas group keys
Da mesma forma, o WindowTransform suporta cálculos simultâneos com múltiplas group keys. O resultado é então associado à tabela de entrada por meio de left join. Veja o exemplo a seguir:
-
Defina os nomes das tabelas de entrada e saída.
input_window_table_name = "rec_sln_demo_behavior_table_preprocess_wide_v1" ds_window_1 = MaxComputeDataSource(table=input_window_table_name, project=project) output_window_table_name_3 = "rec_sln_demo_behavior_test_window_v3" -
Defina as features WindowTransform.
win_feature1 = Feature( name="item__kv_gender_click_7d", input=["gender"], transform=WindowTransform( agg_func="count", condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(7), ), ) win_feature2 = Feature( name="item__kv_gender_click_cnt_7d", input=["gender"], transform=WindowTransform( agg_func="sum", agg_field="click_count", condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(7), ), ) win_feature3 = Feature( name="author__kv_gender_click_15d", input=["gender"], transform=WindowTransform( agg_func="count", condition=Condition(field="event", value="click", operator="="), group_by_keys="author", window_size=DayOf(7), ), ) win_feature4 = Feature( name="author__kv_gender_click_cnt_15d", input=["gender"], transform=WindowTransform( agg_func="sum", agg_field="click_count", condition=Condition(field="event", value="click", operator="="), group_by_keys="author", window_size=DayOf(7), ), ) -
Crie um pipeline e execute a transformação de features WindowTransform.
window_pipeline_3 = fs_client.create_pipeline(ds_window_1, output_window_table_name_3).add_feature_transform([win_feature1, win_feature2, win_feature3, win_feature4]) -
Gere e execute a transformação.
execute_date = '20240605' print("transform_info = ", window_pipeline_3.transform_info) output_window_table_3 = window_pipeline_3.execute(execute_date, drop_table=True) -
Visualize os dados da tabela de destino.
window_ret_3 = output_window_table_3.to_pandas(execute_date, limit=50) window_ret_3
Funções de extensão automática integradas para features WindowTransform
De forma semelhante às transformações de features estatísticas, implementar manualmente cada feature estatística WindowTransform é complexo devido ao grande número de features, que inclui diferentes tamanhos de janela e inúmeras combinações de cálculos de funções de agregação. O sistema oferece funções de extensão automática integradas. Basta especificar as features de entrada a serem contabilizadas, e o sistema gera e completa automaticamente as definições para centenas de features estatísticas.
-
Especifique as features de entrada a serem contabilizadas.
name_prefix = "item" input_list = ['gender'] agg_field = ["click_count"] event_name = 'event' event_type = 'click' group_by_key = "item_id" window_size = [7, 15, 30, 45] window_transform_feature_list = auto_window_feature_transform(name_prefix, input_list, agg_field, event_name, event_type, group_by_key, window_size) print("len_window_transform_feature_list = ", len(window_transform_feature_list)) print("window_transform_feature_list = ", window_transform_feature_list) -
Defina o nome da tabela de saída e crie um pipeline.
input_window_table_name = "rec_sln_demo_behavior_table_preprocess_wide_v1" ds_window_1 = MaxComputeDataSource(table=input_window_table_name, project=project) output_window_table_name_4 = "rec_sln_demo_behavior_test_window_v4" window_pipeline_4 =fs_client.create_pipeline(ds_window_1, output_window_table_name_4).add_feature_transform(window_transform_feature_list) -
Gere e execute a transformação.
execute_date = '20240605' print("transform_info = ", window_pipeline_4.transform_info) output_window_table_4 = window_pipeline_4.execute(execute_date, drop_table=True)
Transformação JoinTransform
No processo de engenharia de features descrito anteriormente, especialmente para AggregationTransform e WindowTransform, as entradas são tabelas de comportamento e os resultados de saída também são armazenados em uma tabela de comportamento. No entanto, a tabela final para uso online normalmente não é uma tabela de comportamento, mas sim uma tabela de features gerada pelo join da tabela de comportamento com outras tabelas, como a tabela de usuários ou a tabela de itens.
Por esse motivo, o JoinTransform foi introduzido para permitir o join das features de AggregationTransform e WindowTransform com tabelas existentes de usuários ou itens.
Associar JoinTransform com WindowTransform
-
Defina a tabela de entrada do WindowTransform.
window_table_name = 'rec_sln_demo_behavior_table_preprocess_wide_v1' ds_window_1 = MaxComputeDataSource(table=window_table_name, project=project) -
Defina as features WindowTransform.
win_fea1 = Feature( name="item__kv_gender_click_7d", input=["gender"], transform=WindowTransform( agg_func="count", condition=Condition(field="event", value="click", operator="="), group_by_keys="item_id", window_size=DayOf(7), ) ) -
Crie um pipeline.
NotaComo outras tabelas precisam ser associadas posteriormente, a tabela de saída não é especificada aqui.
win_pipeline_1 = fs_client.create_pipeline(ds_window_1).add_feature_transform([win_fea1]) -
Defina as tabelas de entrada e saída do JoinTransform.
item_table_name = 'rec_sln_demo_item_table_preprocess_v1' ds_join_1 = MaxComputeDataSource(table=item_table_name, project=project) output_table_name = 'rec_sln_demo_item_table_v1_fs_window_debug_v1' -
Crie um pipeline JoinTransform e conecte-o ao pipeline WindowTransform.
join_pipeline_1 = fs_client.create_pipeline(ds_join_1, output_table_name).merge(win_pipeline_1) -
Gere e execute a transformação.
execute_date = '20240605' output_join_table_1 = join_pipeline_1.execute(execute_date, drop_table=True) -
Visualize os dados da tabela de destino.
join_ret_1 = output_join_table_1.to_pandas(execute_date, limit = 50) join_ret_1 -
Consulte o processo real de cálculo.
output_join_table_1.transform_info
Associar JoinTransform com AggregationTransform
-
Defina a tabela de entrada do AggregationTransform.
agg_table_name = 'rec_sln_demo_behavior_table_preprocess_wide_v1' ds_agg_1 = MaxComputeDataSource(table=agg_table_name, project=project) -
Defina as features AggregationTransform.
agg_fea1 = Feature( name="user_avg_praise_count_1d", input=["praise_count"], transform=AggregationTransform( agg_func="avg", condition=Condition(field="event", value="expr", operator="<>"), group_by_keys="user_id", window_size=DayOf(1), ), ) -
Crie um pipeline.
NotaComo outras tabelas precisam ser associadas posteriormente, a tabela de saída não é especificada aqui.
agg_pipeline_1 = fs_client.create_pipeline(ds_agg_1).add_feature_transform([agg_fea1]) -
Defina as tabelas de entrada e saída do JoinTransform.
user_table_name = 'rec_sln_demo_user_table_preprocess_v1' ds_join_2 = MaxComputeDataSource(table=user_table_name, project=project) output_table_name_2 = 'rec_sln_demo_user_table_v1_fs_window_debug_v1' -
Crie um pipeline JoinTransform e conecte-o ao pipeline AggregationTransform.
join_pipeline_2 = fs_client.create_pipeline(ds_join_2, output_table_name_2).merge(agg_pipeline_1, keep_input_columns=False) -
Gere e execute a transformação.
execute_date = '20240605' output_join_table_2 = join_pipeline_2.execute(execute_date, drop_table=True) -
Visualize os dados da tabela de destino.
join_ret_2 = output_join_table_2.to_pandas(execute_date, limit = 50) join_ret_2 -
Consulte o processo real de cálculo.
output_join_table_2.transform_info
Referências
Para cenários de aplicação, consulte Best practices for feature engineering.
O FeatureStore é adequado para todos os cenários que necessitam de features, como recomendação, controle de risco financeiro e crescimento de usuários. Integrado com os mecanismos de source de dados e de service de recomendação comuns da Alibaba Cloud, o FeatureStore oferece uma plataforma ponta a ponta eficiente e conveniente, desde o registro de features até o desenvolvimento e a aplicação de modelos. Para mais informações sobre o FeatureStore, consulte Overview.
Se tiver dúvidas ao utilizar o FeatureStore, participe do grupo DingTalk 34415007523 para obter assistência técnica.