O DataWorks oferece um nó PyODPS 3 para escrever e executar periodicamente jobs do MaxCompute em Python. Este tópico descreve como configurar e agendar jobs Python usando o DataWorks.
Pré-requisitos
Um nó PyODPS 3 deve estar criado. Para mais informações, consulte Crie e gerencie nós do MaxCompute.
Informações de fundo
O PyODPS é o SDK Python para o MaxCompute. Ele fornece uma interface de programação Python para escrever jobs do MaxCompute, consultar tabelas e visualizações, além de gerenciar recursos. Para mais detalhes, consulte PyODPS. No DataWorks, use um nó PyODPS para agendar e executar jobs Python, integrando-os a outros tipos de jobs.
Notas de uso
-
Caso seu código PyODPS exija pacotes de terceiros, instale-os utilizando um serverless resource group e um custom image.
NotaSe o código incluir uma Função Definida pelo Usuário (UDF) que referencie um pacote de terceiros, esse método não é suportado. Para a configuração correta, consulte Exemplo de UDF: Use pacotes de terceiros em UDFs Python.
Para atualizar a versão do PyODPS, utilize um custom image para executar o comando
/home/tops/bin/pip3 install pyodps==0.12.1em um grupo de recursos serverless (substitua0.12.1pela versão desejada do PyODPS) ou use O&M Assistant para executar o mesmo comando em um grupo de recursos exclusivo para agendamento.
Quando o job PyODPS precisar acessar um ambiente de rede especial, como uma source de dados ou service em uma VPC ou um data center local (IDC), utilize um grupo de recursos serverless e estabeleça uma conexão de rede entre o grupo de recursos e o ambiente de destino. Para mais informações, consulte Soluções de conectividade de rede.
Para obter mais detalhes sobre a sintaxe do PyODPS, consulte a Documentação do PyODPS.
Os nós PyODPS estão disponíveis em dois tipos: PyODPS 2 e PyODPS 3. Eles utilizam versões subjacentes diferentes do Python: os nós PyODPS 2 usam Python 2, enquanto os nós PyODPS 3 usam Python 3. Certifique-se de criar o tipo de nó correspondente à sua versão do Python.
-
Se a execução de SQL em um nó PyODPS falhar ao gerar a linhagem de dados correta, impedindo sua aparição na Data Map, resolva o problema definindo manualmente os parâmetros de agendamento e tempo de execução do DataWorks no seu código. Para saber como visualizar a linhagem de dados, consulte Visualize a linhagem de dados. Para configurações de parâmetros, veja Defina dicas de parâmetros de tempo de execução. Obtenha os parâmetros de tempo de execução necessários usando o seguinte código de exemplo:
import os ... # get DataWorks sheduler runtime parameters skynet_hints = {} for k, v in os.environ.items(): if k.startswith('SKYNET_'): skynet_hints[k] = v ... # setting hints while submiting a task o.execute_sql('INSERT OVERWRITE TABLE XXXX SELECT * FROM YYYY WHERE ***', hints=skynet_hints) ...
O tamanho máximo para o log de saída de um nó PyODPS é de 4 MB. Evite imprimir grandes volumes de resultados de dados no log. Em vez disso, exiba apenas informações essenciais de alerta e progresso.
Limitações
Ao executar um nó PyODPS em um grupo de recursos exclusivo para agendamento, não processe mais de 50 MB de dados locais. Isso se deve às especificações de recursos do grupo exclusivo. Processar uma grande quantidade de dados locais que exceda o limiar do sistema operacional pode causar um erro de falta de memória (OOM), indicado pela mensagem
Got Killed. Evite escrever códigos extensos de processamento de dados diretamente no nó PyODPS.-
Ao executar um nó PYODPS usando um grupo de recursos serverless, configure as CUs para o nó com base na quantidade de dados que ele precisa processar.
NotaAo executar uma tarefa em um grupo de recursos serverless, uma única tarefa suporta uma configuração máxima de
64CU, mas recomendamos não exceder16CUpara evitar escassez de recursos causada por um valor excessivo de CU, o que pode impactar o início da tarefa. Um erro Got killed indica que o uso de memory excedeu o limite, causando o encerramento do processo. Para evitar isso, evite operações de dados locais. Essa limitação não se aplica a jobs SQL ou DataFrame (exceto para
to_pandas) iniciados via PyODPS.Utilize as bibliotecas Numpy e Pandas pré-instaladas para códigos que não envolvam funções personalizadas. Outros pacotes de terceiros contendo código binário não são suportados.
Por motivos de compatibilidade, options.tunnel.use_instance_tunnel é definido como False por padrão no DataWorks. Se você precisar ativar o instance tunnel globalmente, defina manualmente esse valor como True.
-
A definição de bytecode difere entre versões menores do Python 3, como Python 3.8 e Python 3.7.
Atualmente, o MaxCompute usa Python 3.7. Ocorrerá um erro de execução se você usar sintaxe de outras versões do Python 3, como o bloco
finallydo Python 3.8. Recomendamos o uso do Python 3.7. O PyODPS 3 suporta execução em um serverless resource group. Para adquirir e utilizar um, consulte Use grupos de recursos serverless.
Não há suporte para executar múltiplos jobs Python simultaneamente dentro de um único nó PyODPS.
Edite o código: Exemplo básico
Após criar um nó PyODPS, edite e execute seu código. Para mais informações sobre a sintaxe do PyODPS, consulte Visão geral das operações básicas.
-
Ponto de entrada ODPS
Um nó PyODPS do DataWorks fornece uma variável global, chamada odps ou o, como ponto de entrada ODPS. Não é necessário defini-la manualmente.
print(odps.exist_table('PyODPS_iris')) -
Execute SQL
Execute instruções SQL em um nó PyODPS. Para mais informações, consulte SQL.
-
Por padrão, o instance tunnel está desativado no DataWorks. Isso significa que instance.open_reader usa a interface Result, que lê no máximo 10.000 registros. Utilize reader.count para obter o número de registros. Para iterar por todos os dados, desative o
limit. Use as seguintes instruções para ativar o instance tunnel globalmente e desativar olimit.options.tunnel.use_instance_tunnel = True options.tunnel.limit_instance_tunnel = False # Disable the limit to read all data. with instance.open_reader() as reader: # All data can be read through the instance tunnel. -
Adicione
tunnel=Trueao para ativar o instance tunnel na chamada atual do open_reader. Da mesma forma, adicionelimit=Falsepara desativar a restrição delimitna chamada atual.# Use the Instance Tunnel interface for the current open_reader operation to read all data. with instance.open_reader(tunnel=True, limit=False) as reader:
-
-
Parâmetros de tempo de execução
-
Defina os parâmetros de tempo de execução usando o parâmetro hints, que é um dict. Para mais informações sobre hints, consulte Operações SET.
o.execute_sql('select * from PyODPS_iris', hints={'odps.sql.mapper.split.size': 16}) -
Se você definir sql.settings na configuração global, esses parâmetros de tempo de execução serão adicionados a cada execução.
from odps import options options.sql.settings = {'odps.sql.mapper.split.size': 16} o.execute_sql('select * from PyODPS_iris') # This call includes hints from the global configuration.
-
-
Resultados da execução
Uma instância de execução SQL pode executar diretamente a operação open_reader nos dois cenários seguintes:
-
A instrução SQL retorna dados estruturados.
with o.execute_sql('select * from dual').open_reader() as reader: for record in reader: # Process each record. -
Ao executar instruções como
desc, recupere o resultado bruto da execução SQL usando a propriedade reader.raw.with o.execute_sql('desc dual').open_reader() as reader: print(reader.raw)NotaSe você usar parâmetros de agendamento personalizados, deverá codificar rigidamente (hardcode) a hora ao acionar diretamente a execução de um nó PyODPS 3 a partir da página. O nó PyODPS não consegue substituir esse valor diretamente.
-
-
DataFrame
Processe dados também usando um DataFrame (não recomendado).
-
Execução
No ambiente DataWorks, as operações DataFrame devem ser acionadas explicitamente chamando um método de execução imediata.
from odps.df import DataFrame iris = DataFrame(o.get_table('pyodps_iris')) for record in iris[iris.sepal_width < 3].execute(): # Call an immediately executed method to process each record.Se precisar acionar uma execução imediata durante a impressão, ative
options.interactive.from odps import options from odps.df import DataFrame options.interactive = True # Enable the switch at the beginning. iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepal_width.sum()) # An immediate execution is triggered when printing. -
Imprimir informações detalhadas
Configure a opção
options.verbose. Essa opção vem ativada por padrão no DataWorks e imprime informações detalhadas, como a URL do Logview, durante a execução.
-
Exemplo
O exemplo a seguir mostra como usar um nó PyODPS:
Prepare o conjunto de dados e crie a tabela de amostra pyodps_iris. Para detalhes, consulte Processe dados usando DataFrame.
Crie um DataFrame. Para detalhes, consulte Crie um DataFrame a partir de uma tabela do MaxCompute.
-
Insira e execute o seguinte código no nó PyODPS.
from odps.df import DataFrame # Create a DataFrame from an ODPS table. iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepallength.head(5))O seguinte resultado é retornado:
sepallength 0 4.5 1 5.5 2 4.9 3 5.0 4 6.0
Edite o código: Exemplo avançado
Se o nó precisar ser executado periodicamente, defina suas propriedades de agendamento. Para mais informações, consulte Configure propriedades de agendamento para um nó.
Parâmetros de agendamento
No painel direito do editor de nós, clique em Scheduling Settings. Na seção Parameter, configure parâmetros personalizados. A forma como as variáveis são definidas em um nó PyODPS difere da definição em um nó SQL. Para mais informações, consulte Configure parâmetros de agendamento.
Diferentemente dos nós SQL no DataWorks, os nós PyODPS não substituem strings como ${param_name} no código. Em vez disso, um dict chamado args é adicionado às variáveis globais antes da execução do código. Recupere os parâmetros de agendamento desse dict. Por exemplo, se você definir ds=${yyyymmdd} em Parameter, use o seguinte método para recuperar esse parâmetro no seu código.
print('ds=' + args['ds'])
ds=20161116
Se precisar obter a partição chamada ds, utilize o seguinte método.
o.get_table('table_name').get_partition('ds=' + args['ds'])
Para mais informações sobre o desenvolvimento de jobs PyODPS para outros casos de uso, consulte os seguintes tópicos:
Próximos passos
Determine se um script Shell personalizado foi executado com êxito: A lógica para determinar se um script Python personalizado foi executado com êxito é a mesma usada para scripts Shell. Utilize este método para verificação.
Implante um job: Se estiver usando um workspace no modo padrão, implante o job no ambiente de produção antes que ele possa ser executado periodicamente.
O&M para jobs executados periodicamente: Após implantar e agendar um job no ambiente de produção, realize O&M no job em Operation Center.
FAQ do PyODPS: Encontre respostas para perguntas frequentes sobre a execução de jobs PyODPS para ajudar na rápida solução de problemas.
FAQ
P: Estou usando um nó PyODPS3 para coletar dados de uma API de terceiros, como o Lark, e importá-los para o DataWorks. O código roda sem problemas no meu ambiente de desenvolvimento local, mas relata um erro de tempo limite de resposta quando enviado ao ambiente de produção e executado no Operation Center. Por quê?
R: Em . Na lista de permissões de sandbox, adicione o nome de domínio da API de terceiros para conceder acesso ao job PyODPS 3. Por exemplo:
Adicione o nome de domínio da API do Lark open.feishu.cn e defina a porta como 443.