O mecanismo de computação do Lindorm oferece uma API RESTful para envio de jobs Python do Apache Spark. Use-a para executar tarefas de streaming e em lote, machine learning e computação em grafos. Este guia detalha todo o fluxo de trabalho: definição do job, empacotamento de dependências, upload de arquivos para o Object Storage Service (OSS) e envio do job.
Pré-requisitos
Antes de começar, verifique se você tem:
Um mecanismo de computação do Lindorm ativado. Consulte Ativar o serviço.
Um bucket do OSS para armazenar os arquivos do projeto Python, o ambiente de execução e o script de inicialização.
Um AccessKey ID e um AccessKey secret com permissão de leitura e gravação no bucket do OSS. Consulte Criar um par de AccessKey.
Um ambiente Linux (necessário para empacotar arquivos binários Python compatíveis com o mecanismo de computação do Lindorm).
Como funciona
O fluxo de envio de jobs Python do Spark tem quatro etapas:
Defina o job: estruture o projeto Python com os arquivos de ponto de entrada necessários.
Empacote o job: agrupe o código do projeto e o ambiente de execução Python separadamente.
Faça upload dos arquivos: armazene todos os artefatos no OSS.
Envie o job: execute-o pelo console do Lindorm ou pelo Data Management Service (DMS).
Etapa 1: Definir o job
Baixe o pacote de exemplo de job Spark e extraia-o. A pasta extraída chama-se lindorm-spark-examples. Consulte o diretório lindorm-spark-examples/python como referência da estrutura do projeto.
A raiz do projeto (your_project no exemplo) exige três alterações estruturais antes do envio.
1. Adicionar __init__.py
Crie um arquivo __init__.py vazio no diretório your_project. Isso transforma o diretório em um módulo Python importável pelo launcher.py.
2. Preparar main.py
Abra o arquivo your_project/main.py e faça duas alterações:
Adicione o diretório do projeto ao sys.path para resolver as importações corretamente durante a execução:
current_dir = os.path.abspath(os.path.dirname(__file__))
sys.path.append(current_dir)
Envolva a lógica de entrada em uma função main(argv):
def main(argv):
# Write your job logic here
pass
if __name__ == "__main__":
main(sys.argv)
O exemplo abaixo inicializa uma SparkSession:
from pyspark.sql import SparkSession
def main(argv):
spark = SparkSession \
.builder \
.appName("PythonImportTest") \
.getOrCreate()
print(spark.conf)
spark.stop()
if __name__ == "__main__":
main(sys.argv)
3. Criar launcher.py
No diretório raiz de your_project, crie um arquivo chamado launcher.py. Copie o conteúdo de lindorm-spark-examples/python/launcher.py. Esse arquivo atua como ponto de entrada chamado pelo mecanismo de computação do Lindorm. Ele adiciona o diretório do projeto ao sys.path e invoca a função main(argv).
Etapa 2: Empacotar o job
O empacotamento gera dois artefatos separados: um arquivo .zip (ou .egg) com o código do projeto e um arquivo tar com o ambiente de execução Python.
Empacotar o código do projeto
Compacte your_project em um arquivo .zip:
zip -r your_project.zip your_project
Alternativamente, crie um arquivo .egg. Consulte Building Eggs.
Empacotar o ambiente de execução Python
Use Conda ou Virtualenv para empacotar o runtime Python e bibliotecas de terceiros em um arquivo tar. Em seguida, passe-o pelo parâmetro spark.archives.
|
Ferramenta |
Quando usar |
|
Conda |
Use quando o job exigir uma versão específica do Python ou precisar executar em nós sem Python pré-instalado. O Conda inclui o interpretador Python no arquivo tar. |
|
Virtualenv |
Use quando a versão do Python nos nós do cluster já corresponder à do projeto. O Virtualenv não inclui o interpretador; ele usa a instalação Python existente no nó. |
Execute a etapa de empacotamento no Linux. O mecanismo de computação do Lindorm exige arquivos binários Python compilados para Linux.
Exemplo: uso do Conda
conda create -y -n pyspark_conda_env -c conda-forge numpy conda-pack
conda activate pyspark_conda_env
conda pack -f -o pyspark_conda_env.tar.gz
Para outras opções de empacotamento, consulte Python Package Management.
Etapa 3: Fazer upload dos arquivos para o OSS
Faça upload dos três artefatos para o bucket do OSS. Consulte Upload simples.
launcher.py: ponto de entrada criado na etapa 1.your_project.zip(ou.egg): pacote de código do projeto da etapa 2.pyspark_conda_env.tar.gz: ambiente de execução Python da etapa 2.
Etapa 4: Enviar o job
O mecanismo de computação do Lindorm aceita dois métodos de envio:
Console do Lindorm: consulte Gerenciar jobs no console.
DMS: consulte Gerenciar jobs usando o DMS.
Independentemente do método, configure os seguintes parâmetros no campo configs da solicitação de envio do job.
Parâmetros do ambiente de execução
Defina estes três parâmetros para apontar o job para os artefatos no OSS:
|
Parâmetro |
Descrição |
Exemplo |
|
|
Caminho para o arquivo |
|
|
|
Caminho para o arquivo tar do runtime Python. Use |
|
|
|
Caminho para o executável Python dentro do arquivo tar extraído |
|
Parâmetros de acesso ao OSS
Configure estes parâmetros para permitir que o mecanismo de computação leia os arquivos no OSS:
|
Parâmetro |
Descrição |
Exemplo |
|
|
Endpoint do bucket do OSS que armazena os arquivos Python |
|
|
|
AccessKey ID |
|
|
|
AccessKey secret |
|
|
|
Classe usada para acessar o OSS |
|
Para mais parâmetros do OSS no Hadoop, consulte a documentação do Hadoop-Aliyun.
Exemplo: valor de configs montado
O exemplo abaixo mostra todos os parâmetros de ambiente de execução e de acesso ao OSS combinados em um único objeto configs:
{
"spark.archives": "oss://testBucketName/pyspark_conda_env.tar.gz#environment",
"spark.kubernetes.driverEnv.PYSPARK_PYTHON": "./environment/bin/python",
"spark.submit.pyFiles": "oss://testBucketName/your_project.zip",
"spark.hadoop.fs.oss.endpoint": "oss-cn-beijing-internal.aliyuncs.com",
"spark.hadoop.fs.oss.accessKeyId": "<your-access-key-id>",
"spark.hadoop.fs.oss.accessKeySecret": "<your-access-key-secret>",
"spark.hadoop.fs.oss.impl": "org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystem"
}
Substitua <your-access-key-id> e <your-access-key-secret> pelas credenciais AccessKey reais.
Diagnóstico de jobs
Após enviar o job, visualize o status e o endereço da UI do Spark na página Jobs. Consulte Visualizar um job.
Se o envio falhar, abra um ticket e forneça o ID do job e o endereço da UI do Spark à equipe de suporte.