> For the complete documentation index, see [llms.txt](https://aurimrv.gitbook.io/pratica-devops-com-docker-para-machine-learning/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://aurimrv.gitbook.io/pratica-devops-com-docker-para-machine-learning/id-9-model-serving-batch/9-3-airflow.md).

# 9.3 Airflow

## 9.3.1 História e Princípios de Design

A **Airbnb** iniciou como uma startup e tomou proporções globais rapidamente. Esse crescimento exponencial gerou uma complexidade imensa na **orquestração** e no monitoramento de seus fluxos de dados internos. Em 2014, Maxime Beauchemin iniciou a criação do Airflow para facilitar a **coordenação** e agendamento desses pipelines, sendo a ferramenta publicada oficialmente no segundo semestre de 2015.

* *Publicação original:* [Airflow: a workflow management platform (Maxime Beauchemin, 2015)](https://medium.com/airbnb-engineering/airflow-a-workflow-management-platform-46318b977fd8)

Desde a sua publicação, a ferramenta ganhou grande tração na comunidade de engenharia de dados, sendo **incubada** pela Apache Software Foundation em 2016 e promovida a projeto de nível superior (*Top-Level Project*) em 2019.

O design do Apache Airflow foi estruturado em quatro princípios fundamentais introduzidos por Beauchemin, que revolucionaram a forma como a engenharia de dados é praticada:

1. **Pipelines como Código (Pipelines as Code):** Ao definir os pipelines de dados utilizando Python programaticamente, os fluxos se tornam dinâmicos, extensíveis e testáveis. Isso permite que equipes apliquem boas práticas de engenharia de software (como versionamento via Git, testes automatizados e integração contínua) sobre a infraestrutura de dados, superando as limitações de ferramentas legadas baseadas em XML ou interfaces gráficas rígidas de arrastar e soltar (drag-and-drop).
2. **Engenharia de Dados Funcional (Functional Data Engineering):** O Airflow incentiva pipelines projetados sob o paradigma funcional. Isso se apoia em dois pilares:
   * **Idempotência (Reproducibilidade):** A execução de uma tarefa múltiplas vezes com os mesmos parâmetros de entrada deve produzir exatamente o mesmo resultado final, sem gerar efeitos colaterais indesejados.
   * **Capacidade de Recomputação:** Como bugs ocorrem e regras de negócios mudam, a estrutura do pipeline deve permitir a fácil reexecução (*backfill*) de dados históricos sem comprometer o estado atual.
3. **Separação de Responsabilidades (Separation of Concerns):** O Airflow atua estritamente como um **orquestrador**, e não como um motor de processamento pesado. Ele é responsável pelo agendamento, dependências e monitoramento das tarefas, delegando o processamento computacional massivo para engines externas especializadas (como Apache Spark, Snowflake, ou bancos de dados SQL).
4. **Extensibilidade e Modularidade:** Através do conceito de **Providers** (provedores), o Airflow se conecta nativamente com praticamente qualquer tecnologia moderna de nuvem, banco de dados ou serviço SaaS, permitindo a construção de ecossistemas altamente integrados.

***

## 9.3.2 Arquitetura de Componentes

A arquitetura do Apache Airflow é distribuída e desacoplada, facilitando a escalabilidade horizontal em produção. O diagrama abaixo ilustra como seus principais componentes se comunicam:

```mermaid
graph TD
    Scheduler["Scheduler (Orquestração e Agendamento)"] -->|Consulta e Atualiza Estado| DB[(Metadata Database)]
    WebServer["WebServer (Interface WebUI)"] -->|Consulta Estado| DB
    Scheduler -->|Envia Tarefas| Exec["Executor (Mecanismo de Fila)"]
    Exec -->|Executa Tasks| Workers["Workers (Executores de Processos)"]
    Workers -->|Atualiza Status| DB
    DagDir["DAG Directory (Código Python)"] -.->|Leitura de Pipelines| Scheduler
    DagDir -.->|Leitura de Pipelines| WebServer
    DagDir -.->|Leitura de Pipelines| Workers

    style Scheduler fill:#d1ecf1,stroke:#0c5460,stroke-width:1px,color:#000
    style WebServer fill:#fff3cd,stroke:#856404,stroke-width:1px,color:#000
    style Workers fill:#d4edda,stroke:#155724,stroke-width:1px,color:#000
    style DB fill:#f8d7da,stroke:#721c24,stroke-width:1px,color:#000
    style Exec fill:#e2e3e5,stroke:#383d41,stroke-width:1px,color:#000
    style DagDir fill:#e8dbfc,stroke:#5c24b2,stroke-width:1px,color:#000
```

Os principais [componentes do Airflow](https://airflow.apache.org/docs/apache-airflow/stable/concepts/overview.html) são:

* **Scheduler**: O cérebro do sistema. Monitora continuamente as DAGs e suas tarefas, determinando quais tarefas estão prontas para execução com base nas suas dependências e disparando-as para o Executor.
* **WebServer**: A interface gráfica (Web UI) do usuário. Permite inspecionar o status das execuções, disparar DAGs manualmente, visualizar logs detalhados de cada execução e gerenciar conexões com sistemas externos.
* **Metadata Database**: O banco de dados relacional central (geralmente PostgreSQL ou MySQL) que armazena todo o estado do Airflow, incluindo o histórico de execuções das tarefas, usuários, conexões e variáveis.
* **Executor**: O mecanismo que define *onde* e *como* as tarefas serão executadas. Em ambientes de desenvolvimento, usa-se o `SequentialExecutor` ou `LocalExecutor`. Em produção, utilizam-se executores distribuídos baseados em filas (como `CeleryExecutor` ou `KubernetesExecutor`).
* **Workers**: Os nós de computação reais que executam o código definido nas tarefas. Em setups simples, as tarefas rodam no mesmo processo do Scheduler, mas em ambientes de produção de alta escala, os workers são pods isolados no Kubernetes ou servidores Celery dedicados.
* **DAG Directory**: Uma pasta compartilhada no sistema de arquivos contendo os códigos Python das DAGs, lida periodicamente pelo Scheduler e pelo WebServer.

### Blocos de Construção de um Workflow

Para desenvolver fluxos no Airflow, é preciso compreender seus conceitos fundamentais:

* **DAG (Directed Acyclic Graph):** O grafo acíclico direcionado que representa o fluxo completo. Ele agrupa as tarefas e define a ordem de execução e suas dependências.
* **Operators (Operadores):** A definição abstrata de uma tarefa. O Airflow possui centenas de operadores nativos ou fornecidos por provedores comunitários, divididos em três categorias gerais:
  * *Action Operators:* Executam uma ação direta (ex: `PythonOperator` para rodar uma função Python, `BashOperator` para comandos bash).
  * *Transfer Operators:* Movem dados de um sistema para outro (ex: `S3ToRedshiftOperator`).
  * *Sensors (Sensores):* Um tipo especial de operador que bloqueia a execução até que um critério seja satisfeito (ex: `FileSensor` aguarda um arquivo no disco, `SqlSensor` espera um registro no banco de dados).
* **Tasks (Tarefas):** A unidade básica de execução em um pipeline. Um operador instanciado dentro de uma DAG torna-se uma Task.
* **TaskInstance (Instância de Tarefa):** Representa a execução individual de uma Task em um ponto específico do tempo (associada a uma data de execução chamada `logical_date` ou `execution_date`).
* **Hooks (Ganchos):** A interface padronizada do Airflow para interagir com serviços externos. Os Hooks lidam com a autenticação e conexões de rede seguras de forma transparente. Os operadores usam hooks sob o capô para se comunicarem com bancos de dados, clusters Spark, APIs, etc.
* **XCom (Cross-Communication):** O mecanismo padrão do Airflow para trocar metadados simples (como strings ou pequenas listas) entre tarefas distintas do mesmo pipeline, permitindo que tarefas subsequentes tomem decisões com base em saídas anteriores.

## 9.3.3 Preparando container com Airflow

Para preparar as configurações do Airflow crie uma pasta `serving-batch/airflow`.

```docker
FROM apache/airflow:slim-3.2.2-python3.10 AS airflow

RUN python3 -m pip install --upgrade pip wheel setuptools && \
    python3 -m pip cache purge

COPY requirements.txt /tmp/
RUN python3 -m pip install -r /tmp/requirements.txt && \
    python3 -m pip cache purge

ENV AIRFLOW_HOME=/opt/airflow

ENTRYPOINT "airflow" "standalone"
```

Vamos usar o Airflow standalone para fins didáticos. Para **implementações** em **produção**, usar a [documentação oficial](https://airflow.apache.org/docs/apache-airflow/stable/production-deployment.html).

O arquivo de `docker-compose.yaml` nos ajuda a gerir os containers. Neste exemplo vamos criar um volume para as DAGs e os modelos. Copie o modelo joblib na pasta `model` do Airflow.

```yaml
services:
    airflow:
        build: .
        ports:
        - 8080:8080
        volumes:
        - ./dags:/opt/airflow/dags
        - ./model:/model
```

O arquivo de `requirements.txt` segue as mesmas versões das **bibliotecas** que usamos no Spark. Repare que também **incluímos** o provider de Spark do Airflow e o pacote de cliente do Spark Connect (`pyspark-client`), que serão utilizados para integrar o Airflow com o Spark.

```python
pandas
nltk==3.9.4
scikit-learn
apache-airflow-providers-apache-spark==6.1.0
pyspark-client==4.1.2
joblib==1.5.3
pyarrow==24.0.0
numpy
```

Inicialize o Airflow com o `docker-compose` e verifique a WebUI na url [localhost:8080](http://localhost:8080). Com a WebUI é possível acompanhar e operar as **execuções** das DAGs. Explore a [documentação da WebUI](https://airflow.apache.org/docs/apache-airflow/stable/ui.html) para o aprofundamento.

O exemplo a seguir traz a **utilização** do [Bash Operator](https://airflow.apache.org/docs/apache-airflow/stable/tutorial.html), operador do Airflow que executa um comando Bash nos Workers do Airflow.

```python
from datetime import datetime, timedelta
from textwrap import dedent
import joblib

# The DAG object; we'll need this to instantiate a DAG
from airflow import DAG

# Operators; we need this to operate!
from airflow.decorators import dag, task
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator # Correção: PythonOperator está em airflow.operators.python

from airflow.utils.task_group import TaskGroup

from airflow.providers.apache.spark.hooks.spark_connect import SparkConnectHook

with DAG(
    'bash-operator',
    # These args will get passed on to each operator
    # You can override them on a per-task basis during operator initialization
    default_args={
        'depends_on_past': False,
        'email': ['airflow@example.com'],
        'email_on_failure': False,
        'email_on_retry': False,
        'retries': 1,
        'retry_delay': timedelta(minutes=5),
    },
    description='A simple tutorial DAG',
    schedule=None,
    start_date=datetime(2021, 1, 1),
    catchup=False,
    tags=['example'],
) as dag:

    with TaskGroup(group_id='group_for') as tg_for:
        for i in range(1,5):
            t1 = BashOperator(
                task_id=f'print_date_t_for{i}',
                bash_command='date',
            )

    with TaskGroup(group_id='group_sleep') as tg1:
        t2 = BashOperator(
            task_id='sleep_t2',
            depends_on_past=False,
            bash_command='sleep 5',
            retries=3,
        )

        t3 = BashOperator(
            task_id='print_date_t3',
            bash_command='date',
        )

    t4 = BashOperator(
        task_id='print_date_t4',
        bash_command='date',
    )

    tg_for >> [t2, t3] >> t4
```

Faça teste com a dependência entre as Tasks e os [TaskGroups](https://airflow.apache.org/docs/apache-airflow/stable/concepts/dags.html?highlight=taskgroup#taskgroups) e veja o resultado na WebUI do Airflow.

O Scheduler do Airflow tem vários recursos de periodicidade e triggers de execução.

Explore a [documentação](https://airflow.apache.org/docs/apache-airflow/stable/_api/airflow/models/dag/index.html) de **parametrização** das DAGs, altere o pipeline e veja o resultado na WebUI.

Vamos explorar o [Python Operator](https://airflow.apache.org/docs/apache-airflow/stable/howto/operator/python.html).

Inclua a seguinte função na sua DAG:

```python
@task(task_id="print_the_context_t5")
def print_context(ds=None, **kwargs):
    """Print the Airflow context and ds variable from the context."""
    print(kwargs)
    print(ds)
    return 'Whatever you return gets printed in the logs'
```

O decorator `@task` indica que esta função é uma task do Airflow.

Inclua ela no fim como `tarefa5` e veja o resultado:

```python
t5 = print_context()
```

```python
tg_for >> [t2, t3] >> t4 >> t5
```

Repare que a **função** é executada e o contexto da DAG é apresentado no log das tasks.

Agora que já sabemos como orquestrar um **comando** Python em uma DAG, vamos simular a inferência do modelo em uma Task Python.

```python
@task(task_id="inferencia_t6")
def predict(ds=None, **kwargs):
    input_message = ["Figura Transformers Prime War Deluxe - E9687 - Hasbro",  "Senhor dos aneis", "Senhor dos anéis"]
    pipe = joblib.load("/model/classificador-produtos.joblib")
    final_prediction = pipe.predict(input_message)
    print("Predicted values:")
    print(",".join(final_prediction))
    return 'sucesso'
```

Inclua essa task no final da sua DAG como Task6 e veja que é possível executar a inferência dentro de uma Tarefa Python no Airflow. Explore o código para ler de arquivos externos, executar a inferência e gravar o resultado em um arquivo de saída.

## 9.3.4 Data-aware scheduling

É **possível** controlar **dependência** entre DAGs via [Task Sensor](https://airflow.apache.org/docs/apache-airflow/stable/howto/operator/external_task_sensor.html) ou desde a versão 2.4, quando o conceito de Dataset foi criado, é possível ter o mesmo resultado com [Data-aware scheduling](https://airflow.apache.org/docs/apache-airflow/stable/authoring-and-scheduling/datasets.html).

Com este novo conceito uma DAG pode depender de um Dataset para ser iniciada e as tasks dentro de uma DAG podem gerar versões de Datasets. Assim é possível ter **dependência** entre DAGs e **uma** visão parcial de Data Lineage na interface de Dataset.

Vamos gerar duas novas DAGs de teste para visualizarmos esta **dependência**:

```python
from datetime import datetime, timedelta
from airflow.sdk import Asset
from airflow.operators.bash import BashOperator
from airflow import DAG

with DAG(
    'DAG1',
    schedule=None,
    start_date=datetime(2021, 1, 1),
    catchup=False # Adicionado para evitar execuções retroativas
):
    t1 = BashOperator(
                task_id='GenerateDataset1', # Removido f-string desnecessário
                outlets=[Asset("Dataset1")],
                bash_command='date',
            )

with DAG(
    'DAG2',
    schedule=[Asset("Dataset1")],
    start_date=datetime(2021, 1, 1),
    catchup=False # Adicionado para evitar execuções retroativas
):
    t1 = BashOperator(
                task_id='GenerateDataset3', # Removido f-string desnecessário
                outlets=[Asset("Dataset2")],
                bash_command='date',
            )

```
