Airflow gerenciado (Geração 3) | Airflow gerenciado (Geração 2) | Airflow gerenciado (Geração 1 legada)
Nesta página, descrevemos como usar as funções do Cloud Run para acionar DAGs do Serviço Gerenciado para Apache Airflow em resposta a eventos.
O Apache Airflow foi criado para executar DAGs regularmente, mas também é possível acionar DAGs em resposta a eventos. Uma forma de fazer isso é usando as funções do Cloud Run para acionar DAGs do Airflow gerenciado quando ocorre um evento especificado.
Você também pode:
- Acionar DAGs usando apenas a API REST do Airflow.
- Criar uma função que aciona um DAG quando uma mensagem é enviada a um tópico do Pub/Sub.
O exemplo neste guia demonstra uma função que aciona um DAG em resposta a um evento:
- Você configura acionadores para sua função no Cloud Run functions.
- Quando a função é acionada, ela faz uma solicitação para acionar um DAG pela API REST do Airflow do ambiente do Airflow gerenciado. A solicitação contém o identificador e o tipo do evento, além do payload dele.
- O Airflow processa essa solicitação e executa o DAG especificado nela. O DAG gera os dados transmitidos a ele pela função.
Antes de começar
Esta seção lista as etapas preparatórias.
Verificar a configuração de rede do ambiente
Essa solução não funciona nas configurações de IP particular e VPC Service Controls porque não é possível configurar a conectividade das funções do Cloud Run com o servidor da Web do Airflow nessas configurações.
Ativar as APIs do projeto
Console
Ative as APIs do Airflow gerenciado e do Cloud Run functions.
Funções necessárias para ativar APIs
Para ativar as APIs, é necessário ter a permissão serviceusage.services.enable. Se você
criou o projeto, provavelmente já tem essa permissão pelo
papel de proprietário (roles/owner). Caso contrário, você pode receber essa permissão pelo
papel de administrador de uso do serviço (roles/serviceusage.serviceUsageAdmin).
Saiba como conceder papéis.
gcloud
Ative as APIs do Airflow gerenciado e do Cloud Run functions:
Funções necessárias para ativar APIs
Para ativar as APIs, é necessário ter a permissão serviceusage.services.enable. Se você
criou o projeto, provavelmente já tem essa permissão pelo
papel de proprietário (roles/owner). Caso contrário, você pode receber essa permissão pelo
papel de administrador de uso do serviço (roles/serviceusage.serviceUsageAdmin).
Saiba como conceder papéis.
gcloud services enable cloudfunctions.googleapis.comcomposer.googleapis.com
Ativar a API REST do Airflow
Dependendo da sua versão do Airflow, faça o seguinte:
- Para o Airflow 2, a API REST estável já está ativada por padrão. Se a API estável estiver desativada no ambiente, ative a API REST estável.
- Para o Airflow 1, ative a API REST experimental.
Permitir chamadas de API para a API REST do Airflow usando o controle de acesso à rede do servidor da Web
As funções do Cloud Run podem acessar a API REST do Airflow por um endereço IPv4 ou IPv6.
Se você não tiver certeza de qual será o intervalo de IP de chamada, use uma opção de configuração padrão no Controle de acesso ao servidor da Web que seja All IP addresses have access (default) para não bloquear acidentalmente as funções do Cloud Run. É possível
configurar o acesso à rede do servidor da Web mais tarde.
Ver o URL do servidor da Web do Airflow
Este exemplo faz solicitações da API REST para o endpoint do servidor da Web do Airflow.
Use a parte do URL da interface da Web do Airflow antes de .appspot.com no
código da função do Cloud.
Console
No Google Cloud console, acesse a página Ambientes.
Clique no nome do seu ambiente.
Na página Detalhes do ambiente, acesse a guia Configuração do ambiente.
O URL do servidor da Web do Airflow está listado no item IU da Web do Airflow.
gcloud
Execute este comando:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
Substitua:
ENVIRONMENT_NAMEpelo nome do ambienteLOCATIONpela região em que o ambiente está localizado
Acessar o client_id do proxy do IAM
Para fazer uma solicitação ao endpoint de API REST do Airflow, a função exige o ID do cliente do proxy do Identity and Access Management que protege o servidor da Web do Airflow.
O Airflow gerenciado não fornece essas informações diretamente. Em vez disso, faça uma solicitação não autenticada no servidor da Web do Airflow e capture o ID do cliente do URL de redirecionamento:
cURL
curl -v AIRFLOW_URL 2>&1 >/dev/null | grep -o "client_id\=[A-Za-z0-9-]*\.apps\.googleusercontent\.com"
Substitua AIRFLOW_URL pelo URL da interface da Web do Airflow.
Na saída, procure a string após client_id. Exemplo:
client_id=836436932391-16q2c5f5dcsfnel77va9bvf4j280t35c.apps.googleusercontent.com
Python
Salve o código a seguir em um arquivo chamado get_client_id.py. Preencha os valores de project_id, location e composer_environment. Em seguida, execute o código no Cloud Shell ou no ambiente local.
Fazer upload de um DAG para o ambiente
Faça o upload de um DAG para seu ambiente. O exemplo a seguir mostra a configuração de execução do DAG recebida. Você acionará esse DAG a partir de uma função que vai criar neste guia depois.
import datetime
import airflow
from airflow.operators.bash_operator import BashOperator
with airflow.DAG(
'composer_sample_trigger_response_dag',
start_date=datetime.datetime(2026, 1, 1),
# Not scheduled, trigger only
schedule=None) as dag:
# Print the dag_run's configuration, which includes information about the
# Cloud Storage object change.
print_gcs_info = BashOperator(
task_id='print_gcs_info', bash_command='echo {{ dag_run.conf }}}}')
Implantar uma função que aciona o DAG
É possível implantar uma função usando a linguagem preferida com suporte do Cloud Run functions ou do Cloud Run. Este tutorial demonstra uma função do Cloud implementada em Python e Java.
Especificar parâmetros de configuração de função
Acionador: selecione um ou vários acionadores do Eventarc para sua função.
Para mais informações sobre como criar acionadores, consulte Criar acionadores com o Eventarc. Por exemplo, é possível acionar funções do Cloud Storage usando o Eventarc.
Conta de serviço: a conta de serviço especificada para o acionador precisa ter permissões suficientes para acionar DAGs em ambientes do Airflow gerenciado.
Recomendamos seguir o princípio de privilégio mínimo e conceder apenas o papel Usuário do Composer (
composer.user). Para mais informações sobre como configurar permissões, consulte Papéis e permissões para destinos do Cloud Run.Ponto de entrada da função:
(Python) Ao adicionar o código deste exemplo, selecione o ambiente de execução Python 3.10 ou posterior e especifique
trigger_dagcomo o ponto de entrada.
Adicionar requisitos
Especifique as dependências no arquivo requirements.txt:
Adicionar código da função
Coloque o código a seguir no arquivo main.py e faça as seguintes substituições:
Substitua o valor da variável
client_idpelo valorclient_idrecebido anteriormente.Substitua o valor da variável
webserver_idpelo ID do projeto de locatário, que faz parte do URL da interface da Web do Airflow antes de.appspot.com. Você já recebeu o URL da interface da Web do Airflow.Especifique a versão da API REST do Airflow usada:
- Se você usa a API REST estável do Airflow, defina a variável
USE_EXPERIMENTAL_APIcomoFalse. - Se você usa a API REST experimental do Airflow, não é necessário fazer alterações. A variável
USE_EXPERIMENTAL_APIjá está definida comoTrue.
- Se você usa a API REST estável do Airflow, defina a variável
Testar a função
Para verificar se a função e o DAG funcionam conforme o esperado:
- Aguarde até que a função seja implantada.
- Acione a função de acordo com o acionador especificado. Também é possível acionar a função manualmente selecionando a ação Testar a função no Google Cloud console.
- Verifique a página do DAG na interface da Web do Airflow. O DAG precisa ter uma execução ativa ou já concluída.
- Na IU do Airflow, verifique os registros de tarefas desta execução. Você verá que a tarefa
print_gcs_infogera os dados recebidos da função para os registros:
Exemplo de saída:
[2021-04-04 18:25:44,778] {bash_operator.py:154} INFO - Output:
[2021-04-04 18:25:44,781] {bash_operator.py:158} INFO - Triggered from GCF:
{bucket: example-storage-for-gcf-triggers, contentType: text/plain,
crc32c: dldNmg==, etag: COW+26Sb5e8CEAE=, generation: 1617560727904101,
... }
[2021-04-04 18:25:44,781] {bash_operator.py:162} INFO - Command exited with
return code 0h
A seguir
- Acessar a IU do Airflow
- Acessar a API REST do Airflow
- Gravar DAGs
- Gravar funções do Cloud Run
- Gatilhos do Cloud Storage