Para permitir a leitura e a gravação em tabelas do MaxCompute em jobs do Deep Learning Containers (DLC), a equipe do Platform for AI (PAI) desenvolveu o módulo PAIIO. O PAIIO oferece três tipos de interfaces: TableRecordDataset, TableReader e TableWriter. Este tópico explica como usar essas interfaces para ler e gravar dados em tabelas do MaxCompute e fornece exemplos de código.
Limitações
O PAIIO é compatível apenas com jobs do DLC que usam imagens do TensorFlow 1.12, 1,15 ou 2,0.
O PAIIO não suporta imagens personalizadas.
Configure informações da conta
Antes de usar o módulo paiio para ler ou gravar em tabelas do MaxCompute, configure o AccessKey da sua conta do MaxCompute. O PAI lê essa configuração de um arquivo. Coloque esse arquivo em um sistema de arquivos montado e referencie-o no seu código por meio de uma variável de ambiente.
-
Crie um arquivo de configuração com o seguinte conteúdo:
access_id=xxxx access_key=xxxx end_point=http://xxxxParâmetro
Descrição
access_id
O AccessKey ID da sua conta Alibaba Cloud.
access_key
O AccessKey secret da sua conta Alibaba Cloud.
end_point
O endpoint do MaxCompute. Por exemplo, o endpoint da região China (Shanghai) é
http://service.cn-shanghai.maxcompute.aliyun.com/api. Para mais informações, consulte Endpoints. -
No seu código, especifique o caminho do arquivo de configuração da seguinte forma:
os.environ['ODPS_CONFIG_FILE_PATH'] = '<your MaxCompute config file path>'Substitua <your MaxCompute config file path> pelo caminho real do seu arquivo de configuração.
TableRecordDataset
API
A comunidade TensorFlow recomenda o uso da interface Dataset no TensorFlow 1.2 e versões posteriores para construir pipelines de entrada. Essa interface substitui as antigas interfaces de threads e filas. Combine e transforme vários objetos Dataset para gerar dados destinados à computação e simplificar o código de entrada de dados.
-
Definição em Python
class TableRecordDataset(Dataset): def __init__(self, filenames, record_defaults, selected_cols=None, excluded_cols=None, slice_id=0, slice_count=1, num_threads=0, capacity=0): -
Parâmetros
Parâmetro
Obrigatório
Tipo
Padrão
Descrição
filenames
Sim
STRING
-
Lista de tabelas para leitura. Todas as tabelas devem ter o mesmo schema. O nome da tabela deve seguir o formato
odps://${your_projectname}/tables/${table_name}/${pt_1}/${pt_2}/....record_defaults
Sim
LIST ou TUPLE
-
Lista ou tupla que define o tipo de dado e o valor padrão de cada coluna a ser lida. O método lança uma exceção se a quantidade de elementos não corresponder ao número de colunas ou se houver incompatibilidade de tipos.
Os tipos de dados suportados incluem
FLOAT32,FLOAT64,INT32,INT64,BOOLeSTRING. Para valores padrão do tipoINT64, usenp.array(0, np.int64).selected_cols
Não
STRING
None
String separada por vírgulas com os nomes das colunas a serem lidas. Se este parâmetro for None, todas as colunas serão lidas. Não use este parâmetro com excluded_cols.
excluded_cols
Não
STRING
None
String separada por vírgulas com os nomes das colunas a serem excluídas. Se este parâmetro for None, nenhuma coluna será excluída. Não use este parâmetro com selected_cols.
slice_id
Não
INT
0
Em leituras distribuídas, define o índice (base zero) do shard de dados a ser lido. O sistema divide a tabela na quantidade de shards especificada por slice_count e lê o shard correspondente a este slice_id.
Se slice_id for 0 (padrão) e slice_count for 1, a tabela inteira será lida. Caso slice_count seja maior que 1, apenas o primeiro shard (índice 0) será lido.
slice_count
Não
INT
1
Indica o número total de shards nos quais os dados serão divididos durante a leitura distribuída. Geralmente, esse valor corresponde ao número de workers. O valor padrão 1 significa que a tabela não será fragmentada e o leitor processará a tabela inteira.
num_threads
Não
INT
0
Define quantas threads paralelas o leitor usará para pré-carregar dados de cada tabela. Essas threads funcionam independentemente das threads de computação. O valor deve ser um inteiro entre 1 e 64. Se num_threads estiver definido como 0, o sistema ajustará automaticamente o número de threads de pré-carregamento para um quarto do número de threads de computação.
NotaAumentar o número de threads de pré-carregamento não garante um treinamento de modelo mais rápido, pois o impacto de I/O varia conforme o modelo.
capacity
Não
INT
0
Especifica o total de linhas a serem pré-carregadas de uma tabela. Quando num_threads for maior que 1, a capacidade de pré-carregamento de cada thread será de capacity/num_threads linhas, arredondada para cima. Se capacity estiver definido como 0, o Reader interno configurará automaticamente a capacidade total de pré-carregamento com base no tamanho médio das primeiras N linhas da tabela, onde N assume o valor padrão de 256. Isso assegura que a quantidade de dados pré-carregados por thread seja de aproximadamente 64 MB.
NotaSe um campo em uma tabela do MaxCompute for do tipo DOUBLE, mapeie-o para
np.float64no TensorFlow. -
Valor de retorno
Retorna um objeto
Datasetque pode ser usado para construir um pipeline de dados.
Exemplo
Considere uma tabela chamada test no projeto myproject com o seguinte conteúdo parcial.
|
itemid (BIGINT) |
name (STRING) |
price (DOUBLE) |
virtual (BOOL) |
|
25 |
"Apple" |
5,0 |
False |
|
38 |
"Pear" |
4,5 |
False |
|
17 |
"Watermelon" |
2,2 |
False |
O código abaixo demonstra como usar a interface TableRecordDataset para ler as colunas itemid e price da tabela test.
import os
import tensorflow as tf
import paiio
# Specify the path to the configuration file. Replace this with the actual file path.
os.environ['ODPS_CONFIG_FILE_PATH'] = "/mnt/data/odps_config.ini"
# Define the table(s) to read. Replace with your actual project and table names.
table = ["odps://${your_projectname}/tables/${table_name}"]
# Define the TableRecordDataset to read the 'itemid' and 'price' columns.
dataset = paiio.data.TableRecordDataset(table,
record_defaults=[0, 0.0],
selected_cols="itemid,price",
num_threads=1,
capacity=10)
# Set 2 epochs, a batch size of 3, and prefetch 100 batches.
dataset = dataset.repeat(2).batch(3).prefetch(100)
ids, prices = tf.compat.v1.data.make_one_shot_iterator(dataset).get_next()
with tf.compat.v1.Session() as sess:
sess.run(tf.compat.v1.global_variables_initializer())
sess.run(tf.compat.v1.local_variables_initializer())
try:
while True:
batch_ids, batch_prices = sess.run([ids, prices])
print("batch_ids:", batch_ids)
print("batch_prices:", batch_prices)
except tf.errors.OutOfRangeError:
print("End of dataset")
TableReader
Referência da API
O TableReader foi desenvolvido sobre o SDK do MaxCompute e opera independentemente do framework TensorFlow. Ele permite acessar diretamente as tabelas do MaxCompute e obter resultados de I/O em tempo real.
-
Crie um leitor e abra uma tabela
Sintaxe
Parâmetros
-
Valor de retorno
Retorna um objeto Reader.
reader = paiio.python_io.TableReader(table, selected_cols="", excluded_cols="", slice_id=0, slice_count=1):Parâmetro
Obrigatório
Tipo
Padrão
Descrição
table
Sim
STRING
N/A
Nome da tabela do MaxCompute a ser aberta. O nome deve seguir o formato:
odps://${your_projectname}/tables/${table_name}/${pt_1}/${pt_2}/...selected_cols
Não
STRING
String vazia ("")
String separada por vírgulas com os nomes das colunas a selecionar. Se uma string vazia ("") for fornecida, todas as colunas serão lidas. Não use este parâmetro com excluded_cols.
excluded_cols
Não
STRING
String vazia ("")
String separada por vírgulas com os nomes das colunas a excluir. Se uma string vazia ("") for fornecida, todas as colunas serão lidas. Não use este parâmetro com selected_cols.
slice_id
Não
INT
0
Em cenários de leitura distribuída, indica o índice do shard atual. O valor pode variar entre [0, slice_count-1]. Na leitura distribuída, o sistema divide a tabela em múltiplos shards conforme slice_count e lê o shard especificado por slice_id. O valor padrão 0 indica que a tabela não será fragmentada e todas as linhas serão lidas.
slice_count
Não
INT
1
Em cenários de leitura distribuída, define o número total de shards, que geralmente corresponde ao número de workers.
-
Ler registros
Sintaxe
-
Parâmetros
num_records define a quantidade de linhas a serem lidas sequencialmente. O valor padrão é 1, o que resulta na leitura de uma única linha. Se num_records exceder o número de linhas não lidas, todas as linhas restantes serão retornadas. Caso nenhum registro seja lido, uma exceção
paiio.python_io.OutOfRangeExceptionserá lançada. -
Valor de retorno
Retorna um ndarray (ou recarray) do NumPy. Cada elemento do array é uma tupla que representa uma linha da tabela.
reader.read(num_records=1) -
Posicionar em uma linha específica
Sintaxe
Parâmetros
-
Valor de retorno
Nenhum. Uma exceção será lançada caso ocorra um erro.
reader.seek(offset=0)offset indica a linha para a qual o ponteiro deve ser movido (a indexação começa em 0). A próxima operação de leitura iniciará nessa linha. Se slice_id e slice_count estiverem configurados, o posicionamento será relativo à posição dentro do shard. Se offset ultrapassar o número total de linhas da tabela, uma
OutOfRangeExceptionserá lançada. Tentar posicionar novamente quando a posição de leitura já estiver além do final da tabela também resultará em umapaiio.python_io.OutOfRangeException.ImportanteAo ler um lote, se o número de linhas restantes for menor que o
batch_size, a operaçãoreadretorna as linhas restantes sem lançar exceções. Nesse cenário, tentar outra operaçãoseekgerará uma exceção. -
Obter a contagem total de linhas
Sintaxe
-
Parâmetros
Nenhum
-
Valor de retorno
Retorna o número de linhas da tabela. Se slice_id e slice_count estiverem configurados, retorna o tamanho do shard.
reader.get_row_count() -
Obter o schema da tabela
Sintaxe
-
Parâmetros
Nenhum
Valor de retorno
reader.get_schema()Retorna um ndarray estruturado unidimensional. Cada elemento descreve uma coluna selecionada da tabela do MaxCompute e contém os três campos a seguir.
Parâmetro
Descrição
colname
Nome da coluna.
typestr
Nome do tipo de dado do MaxCompute.
pytype
Tipo de dado Python correspondente a typestr.
A tabela a seguir descreve o mapeamento entre typestr e pytype.
typestr
pytype
BIGINT
INT
DOUBLE
FLOAT
BOOLEAN
BOOL
STRING
OBJECT
DATETIME
INT
MAP
NotaO PAI-TensorFlow não suporta dados do tipo MAP.
OBJECT
-
Fechar a tabela
Sintaxe
-
Parâmetros
Nenhum
-
Valor de retorno
Nenhum. Uma exceção será lançada caso ocorra um erro.
reader.close() -
Crie um gravador e abra uma tabela
-
Sintaxe
writer = paiio.python_io.TableWriter(table, slice_id=0)NotaEsta operação anexa dados a uma tabela e não apaga os dados existentes.
Os dados recém-gravados só poderão ser lidos após o fechamento da tabela.
-
Parâmetros
Parâmetro
Obrigatório
Tipo
Padrão
Descrição
table
Sim
STRING
None
Nome da tabela do MaxCompute a ser aberta. O nome deve seguir o formato:
odps://${your_projectname}/tables/${table_name}/${pt_1}/${pt_2}/...slice_id
Não
INT
0
ID do shard onde os dados serão gravados. No modo distribuído, gravar em shards diferentes evita conflitos de escrita. No modo standalone, use o valor padrão 0. No modo distribuído, a operação de gravação falhará se múltiplos workers, incluindo nós de servidor de parâmetros (PS), tentarem gravar no mesmo shard usando o mesmo slice_id.
-
Valor de retorno
Retorna um objeto Writer.
-
-
Gravar registros
-
Sintaxe
writer.write(values, indices) -
Parâmetros
Parâmetro
Obrigatório
Tipo
Padrão
Descrição
values
Sim
STRING
None
Dados a serem gravados, especificados como um único registro ou múltiplos registros:
-
Para gravar um único registro, passe uma TUPLE, LIST ou ndarray 1D de escalares para o parâmetro values. Ao passar uma LIST ou ndarray, todas as colunas do registro devem ter o mesmo tipo de dado.
-
Para gravar um ou mais registros, passe uma LIST ou ndarray 1D para o parâmetro values. Cada elemento deve ser uma TUPLE, LIST ou elemento de ndarray estruturado representando um único registro.
indices
Sim
INT
None
Índices das colunas onde os dados serão gravados. Pode ser uma TUPLE, LIST ou ndarray 1D de inteiros. Cada índice em indices refere-se ao número da coluna baseado em zero.
-
-
Valor de retorno
Número de registros gravados com sucesso. Se a operação falhar, uma exceção será lançada.
-
-
Fechar a tabela
-
Sintaxe
writer.close()NotaNão é necessário chamar explicitamente o método close() ao usar uma instrução with.
-
Parâmetros
Nenhum
-
Valor de retorno
Nenhum. Se ocorrer um erro, uma exceção será lançada.
-
Exemplo
O código abaixo demonstra como usar o TableWriter com uma instrução with.
with paiio.python_io.TableWriter(table) as writer: # Prepare values for writing. writer.write(values, indices) # The writer is closed automatically when the 'with' block is exited.
-
Crie um conjunto de dados e envie seus arquivos de configuração e código para a source de dados. Para mais informações, consulte Criar e gerencie conjuntos de dados.
-
Crie um job do DLC. Os principais parâmetros estão descritos abaixo. Para outros parâmetros, consulte Criar um job de treinamento.
Node Image: Em Alibaba Cloud Images, selecione uma imagem para TensorFlow 1.12, TensorFlow 1,15 ou TensorFlow 2,0.
Dataset Configuration: Para Dataset, selecione o conjunto de dados criado na etapa 1 e defina Mount Path como
/mnt/data/.Job Command: Insira
python /mnt/data/xxx.py. Substitua xxx.py pelo nome do arquivo de código enviado na etapa 1.
-
Clique em Confirm.
Após enviar o job de treinamento, visualize os resultados nos logs do job. Para mais informações, consulte Visualizar logs do job.
Exemplo
Este exemplo usa uma tabela chamada test no projeto myproject com os seguintes dados.
|
uid (BIGINT) |
name (STRING) |
price (DOUBLE) |
virtual (BOOL) |
|
25 |
"Apple" |
5,0 |
False |
|
38 |
"Pear" |
4,5 |
False |
|
17 |
"Watermelon" |
2,2 |
False |
O código a seguir mostra como usar o TableReader para ler dados das colunas uid, name e price.
import os
import paiio
# Specify the path of the configuration file. Replace the value with the actual path.
os.environ['ODPS_CONFIG_FILE_PATH'] = "/mnt/data/odps_config.ini"
# Open a table. Replace myproject and test with your project and table names.
reader = paiio.python_io.TableReader("odps://myproject/tables/test", selected_cols="uid,name,price")
# Get the total number of rows in the table.
total_records_num = reader.get_row_count() # return 3
batch_size = 2
# Read the table. The return value is a recarray in the format [(uid, name, price)*2].
records = reader.read(batch_size) # Returns [(25, "Apple", 5.0), (38, "Pear", 4.5)]
records = reader.read(batch_size) # Returns [(17, "Watermelon", 2.2)]
# Reading again throws an OutOfRangeException.
# Close the reader.
reader.close()
Uso do TableWriter
O TableWriter baseia-se no SDK do MaxCompute e não depende do framework TensorFlow, permitindo gravar dados diretamente nas tabelas do MaxCompute.
API
Exemplo
import paiio
import os
# Specify the path of the configuration file. Replace the value with the actual path.
os.environ['ODPS_CONFIG_FILE_PATH'] = "/mnt/data/odps_config.ini"
# Prepare the data.
values = [(25, "Apple", 5.0, False),
(38, "Pear", 4.5, False),
(17, "Watermelon", 2.2, False)]
# Open a table to get a writer object. Replace the project and table names with your actual values.
writer = paiio.python_io.TableWriter("odps://project/tables/test")
# Write records to columns 0 to 3 of the table.
records = writer.write(values, indices=[0, 1, 2, 3])
# Close the writer.
writer.close()
Próximas etapas
Após configurar o código, siga estas etapas para usar o PAIIO na leitura e gravação de tabelas do MaxCompute: