Este tópico explica como implantar e iniciar jobs PyFlink em streaming e batch, abordando o fluxo de desenvolvimento no Realtime Compute for Apache Flink.
Pré-requisitos
Se você utilizar um usuário RAM ou uma função RAM para acessar o console, verifique se a identidade tem as permissões necessárias. Para mais informações, consulte Gerenciamento de Permissões.
Crie um workspace. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
Etapa 1: Preparar arquivos de código Python
O console de gerenciamento do Realtime Compute for Apache Flink não oferece ambiente de desenvolvimento Python. Desenvolva seus jobs localmente. Para mais detalhes sobre depuração de jobs e conectores, consulte Desenvolver jobs PyFlink.
Garanta que a versão do Flink usada no desenvolvimento local corresponda à versão do mecanismo selecionada na Etapa 3: Implantar um job PyFlink. Para saber como usar outras dependências, como ambientes virtuais Python personalizados, pacotes Python de terceiros, pacotes JAR e arquivos de dados, consulte Usar dependências Python.
Para agilizar o início, este tópico fornece arquivos Python de exemplo para um job de contagem de palavras e um arquivo de dados de amostra. Baixe-os e use-os nas etapas a seguir.
-
Baixe o arquivo de job Python de exemplo apropriado.
Job de streaming: word_count_streaming.py.
Job batch: word_count_batch.py.
Clique em Shakespeare para baixar o arquivo de dados de exemplo.
Etapa 2: Fazer upload dos arquivos Python e de dados
Faça login no console do Realtime Compute.
Localize o workspace Flink desejado e clique em Console na coluna Actions.
No painel de navegação à esquerda, clique em Artifacts.
-
Clique em Upload Artifact para enviar os arquivos Python e de dados.
Envie os arquivos Python e de dados de exemplo baixados na Etapa 1. Para mais informações sobre caminhos de armazenamento de arquivos, consulte Artifacts.
Etapa 3: Implantar um job PyFlink
Streaming
Na página , clique em .
-
Configure os parâmetros da implantação.
Parâmetro
Descrição
Exemplo
Deployment mode
Selecione o modo stream.
stream mode
Deployment name
Insira um nome para a implantação Python.
flink-streaming-test-python
Engine version
Versão do mecanismo Flink para a implantação.
Recomendamos usar uma versão com a tag RECOMMENDED ou STABLE para garantir maior confiabilidade e desempenho. Para mais informações, consulte Notas de Lançamento e Versões do Mecanismo.
vvr-8.0.9-flink-1.17
Python URI
Baixe o arquivo de exemplo word_count_streaming.py. Em seguida, clique no ícone de upload
para selecionar e enviar o arquivo.Se o arquivo já existir em Artifacts, selecione-o diretamente sem precisar fazer novo upload.
-
Entry module
Módulo de ponto de entrada do programa.
-
Este parâmetro não é necessário se o job PyFlink for um arquivo .py.
-
Se o job PyFlink for um arquivo .zip, insira o módulo de entrada. Exemplo:
word_count.
Não obrigatório
Entry point main arguments
Argumentos a serem passados para o método principal.
Neste tutorial, insira o caminho de armazenamento do arquivo de dados de entrada, Shakespeare.
--input oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/ShakespeareCopie o caminho completo do arquivo Shakespeare na página Artifacts.
Deployment target
Na lista suspensa, selecione a queue ou o session cluster de destino. Clusters de sessão não são recomendados para produção. Para mais informações, consulte Gerenciar filas e Criar um cluster de sessão.
ImportanteImplantações em um cluster de sessão não suportam métricas de monitoramento, configuração de alertas ou Autopilot. Use clusters de sessão apenas para desenvolvimento e testes; não os utilize em ambientes de produção. Para mais informações, consulte Depurar implantações.
default-queue
Para mais detalhes sobre outros parâmetros de configuração, consulte Implantar um job.
-
Clique em Deploy.
Batch
Na página , clique em Create Deployment e selecione Python Deployment.
-
Configure os parâmetros da implantação.
Parâmetro
Descrição
Exemplo
Deployment mode
Selecione o modo batch.
batch mode
Deployment name
Insira um nome para a implantação Python.
flink-batch-test-python
Engine version
Versão do mecanismo Flink para a implantação.
Recomendamos usar uma versão com a tag RECOMMENDED ou STABLE para garantir maior confiabilidade e desempenho. Para mais informações, consulte Notas de Lançamento e Versões do Mecanismo.
vvr-8.0.9-flink-1.17
Python URI
Baixe o arquivo de exemplo word_count_batch.py. Em seguida, clique no ícone de upload
para selecionar e enviar o arquivo.-
Entry module
Módulo de ponto de entrada do programa.
-
Este parâmetro não é necessário se o job PyFlink for um arquivo .py.
-
Se o job PyFlink for um arquivo .zip, insira o módulo de entrada. Exemplo:
word_count.
Não obrigatório
Entry point main arguments
Argumentos a serem passados para o método principal.
Neste tutorial, insira os caminhos de armazenamento para o arquivo de entrada Shakespeare e o diretório de saída
python-batch-quickstart-test-output.NotaBasta especificar o caminho do diretório de saída. O diretório de saída deve estar no mesmo diretório pai do arquivo de entrada. Não é necessário criar o diretório de saída antecipadamente.
--input oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/Shakespeare--output oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/python-batch-quickstart-test-outputCopie o caminho completo do arquivo Shakespeare na página Artifacts.
Deployment target
Na lista suspensa, selecione a queue ou o session cluster de destino. Clusters de sessão não são recomendados para produção. Para mais informações, consulte Gerenciar filas e Criar um cluster de sessão.
ImportanteImplantações em um cluster de sessão não suportam métricas de monitoramento, configuração de alertas ou Autopilot. Use clusters de sessão apenas para desenvolvimento e testes; não os utilize em ambientes de produção. Para mais informações, consulte Depurar implantações.
default-queue
Para mais detalhes sobre outros parâmetros de configuração, consulte Implantar um job.
-
Clique em Deploy.
Etapa 4: Iniciar a implantação e visualizar resultados
Streaming
Na página , localize a implantação desejada e clique em Start na coluna Actions.
-
Na caixa de diálogo Start Job, selecione Initial Mode e clique em Start. Para mais informações, consulte Iniciar uma implantação.
Após clicar em Start, o status RUNNING ou FINISHED indica que a implantação está em execução conforme esperado. Se você usar o arquivo de exemplo deste tópico, o status final será FINISHED.
-
Quando o status da implantação mudar para RUNNING, visualize os resultados da implantação de streaming.
ImportanteSe usar o arquivo Python de exemplo deste tópico, os resultados serão excluídos quando a implantação de streaming entrar no estado FINISHED. Portanto, só é possível visualizar os resultados enquanto a implantação estiver no estado RUNNING.
No arquivo de log do TaskManager que termina com .out, pesquise por
shakespearepara encontrar o resultado do cálculo.Na aba Logs, clique na aba Running Task Managers. Para o TaskManager relevante, clique na subaba Log List. Abra o arquivo
flink.oute insirashakespearena caixa de pesquisa no canto superior direito para localizar o resultado da contagem de palavras, como(shakespeare,1).
Batch
-
Na página , localize a implantação desejada e clique em Start na coluna Actions.
Para filtrar a lista, selecione Batch Deployment na lista suspensa de tipos.
Na caixa de diálogo Start Job, clique em Start. Para mais informações, consulte Iniciar uma implantação.
-
Quando o status da implantação mudar para FINISHED, visualize os resultados da implantação batch.
Faça login no console do OSS. Acesse o diretório oss://<Your-OSS-Bucket-Name>/artifacts/namespaces/<Your-Workspace-Name>/python-batch-quickstart-test-output. Clique na pasta nomeada com a data e hora de início da implantação, clique no arquivo alvo e, em seguida, clique em Download no painel exibido.

A implantação batch gera um arquivo .ext. Após baixar o arquivo, abra-o com um editor de texto ou Microsoft Word para visualizar os resultados. A saída é semelhante à seguinte:
(As,40) (At,5) (Ay,1) (Be,9) (By,14) (Do,4) (He,7) (I,,4) (If,34) (In,36) (Is,10) (It,6)
(Opcional) Etapa 5: Parar uma implantação
Para aplicar alterações a um job (como modificações de código, atualizações de parâmetros WITH ou mudanças de versão), reimplemente-o, pare-o e reinicie-o. Também é necessário reiniciar para um início sem estado ou para aplicar alterações de configuração não dinâmicas. Para mais informações sobre como parar um job, consulte Parar um job.
Tópicos relacionados
Você pode configurar recursos para uma implantação antes de iniciá-la ou modificá-los após a execução. Há dois modos de configuração de recursos: básico (granularidade grossa) e especialista (granularidade fina). Para mais informações, consulte Configurar recursos da implantação.
O Realtime Compute for Apache Flink suporta atualizações dinâmicas nos parâmetros da implantação. Isso permite que as configurações entrem em vigor mais rapidamente e reduz o tempo de inatividade do service causado pela interrupção e reinício das implantações. Para mais informações, consulte Dimensionamento dinâmico e atualizações de parâmetros.
Configure níveis de log da implantação e especifique saídas diferentes para cada nível. Para mais informações, consulte Configurar Saída de Log.
Para um passo a passo do fluxo de desenvolvimento SQL, consulte Job Flink SQL.
Construir um data lakehouse de streaming com Paimon e StarRocks.