> 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-4-testando-tudo-junto.md).

# 9.4 Testando tudo junto

Para integrar o Apache Airflow com o Apache Spark de forma moderna e simplificada, utilizaremos o **Spark Connect**. Esse protocolo gRPC nativo permite que o Airflow atue como um cliente fino (*thin client*), comunicando-se diretamente com o driver do Spark remoto, eliminando completamente gateways intermediários complexos (como o antigo Apache Livy) e reduzindo a sobrecarga de JVM/Java no lado do Airflow.

***

## 9.4.1 Preparando o Spark Connect Server

Utilizaremos o mesmo ecossistema Docker criado para o Spark na seção [9.2](/pratica-devops-com-docker-para-machine-learning/id-9-model-serving-batch/9-2-spark.md). Criaremos uma pasta `serving-batch/spark-connect` para organizar o Dockerfile do servidor. O contêiner será configurado para subir o servidor do Spark Connect no momento da inicialização do Docker em primeiro plano.

```docker
FROM python:3.10.11-slim

ENV SPARK_VERSION=4.1.2
ENV HADOOP_VERSION=3
ENV SPARK_VERSION_STRING=spark-${SPARK_VERSION}-bin-hadoop${HADOOP_VERSION}
ENV SPARK_HOME=/spark

# Instalação do Spark
ADD https://dlcdn.apache.org/spark/spark-${SPARK_VERSION}/${SPARK_VERSION_STRING}.tgz /
RUN tar xzf ${SPARK_VERSION_STRING}.tgz \
    && mv ${SPARK_VERSION_STRING} ${SPARK_HOME} \
    && rm ${SPARK_VERSION_STRING}.tgz

ENV OPENJDK_VERSION=17

# Instalação do Java e dependências do SO
RUN apt-get -y update && \
    apt-get install --no-install-recommends -y \
    openjdk-${OPENJDK_VERSION}-jdk \
    procps zip libssl-dev libkrb5-dev libffi-dev libxml2-dev libxslt1-dev python-dev build-essential unzip && \
    apt-get clean && rm -rf /var/lib/apt/lists/*

# Atualização de pip e instalação de dependências Python
COPY requirements.txt /tmp/ 
RUN python3 -m pip install --upgrade pip wheel setuptools && \
    python3 -m pip install -r /tmp/requirements.txt && \
    python3 -m pip cache purge

ENV PATH=/usr/local/openjdk-${OPENJDK_VERSION}/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin:${SPARK_HOME}/bin
ENV PYSPARK_PYTHON=/usr/local/bin/python3
ENV SPARK_CONF_DIR=${SPARK_HOME}/conf

# Porta gRPC padrão do Spark Connect Server
EXPOSE 15002

# Para manter o contêiner Docker ativo, iniciamos a classe Java do Spark Connect 
# diretamente em primeiro plano (foreground) usando o spark-class
ENTRYPOINT ["/spark/bin/spark-class", "org.apache.spark.sql.connect.service.SparkConnectServer", "--packages", "org.apache.spark:spark-connect_2.13:4.1.2"]
```

***

## 9.4.2 Docker Compose de Integração

Agora, na pasta raiz `serving-batch/`, vamos criar o `docker-compose.yaml` para orquestrar o Airflow e o Spark Connect Server. Compartilharemos os volumes locais de dados e modelos para que o Spark consiga ler o modelo treinado e os dados do arquivo de produtos.

```yaml
services:
    airflow:
        build: ./airflow
        ports:
        - "8080:8080" # Porta do Airflow Webserver
        volumes:
        - ./dags:/opt/airflow/dags
        - ./model:/model
        - ./data:/data
    spark-connect:
        build: ./spark-connect
        ports:
        - "15002:15002" # Porta do Spark Connect gRPC
        - "4040:4040"   # Porta do Spark Web UI (opcional para análise)
        volumes:
        - ./model:/model
        - ./data:/data
```

Certifique-se de ter colocado o modelo serializado `classificador-produtos.joblib` na pasta `./model` e a tabela `produtos.csv` na pasta `./data` antes de subir o ambiente.

***

## 9.4.3 Escrevendo a DAG de Inferência com Spark Connect

Com o Spark Connect, a comunicação é gRPC e baseada em planos de dados. O Airflow 3.x fornece suporte nativo a este fluxo através do decorador `@task.pyspark` (que utiliza o `PySparkOperator` e `SparkConnectHook` sob o capô, injetando a `SparkSession` conectada automaticamente no método decorado).

> \[!IMPORTANT] **Broadcast em Spark Connect:** O protocolo do Spark Connect não oferece suporte direto ao objeto `SparkContext.broadcast()`. Em vez disso, utilizamos o carregamento do modelo binário diretamente em cache dentro do escopo da UDF vetorizada no Worker do Python. Para otimizar a performance e carregar o modelo em memória uma única vez por processo executor, aplicamos o padrão de inicialização preguiçosa (*Lazy Initialization*) utilizando uma variável global.

Crie o arquivo da DAG na pasta `./dags/inference_dag.py`:

```python
from datetime import datetime, timedelta
from airflow import DAG
from airflow.decorators import task

@task.pyspark(conn_id="spark_connect_default", task_id="batch_inference_task")
def run_batch_inference(spark):
    from pyspark.sql.functions import monotonically_increasing_id, pandas_udf
    import pandas as pd
    import joblib

    from typing import Iterator

    # 1 & 2. Definição da UDF do Pandas com Lazy Loading via Iterator (carrega uma vez por partição)
    @pandas_udf("string")
    def predict(iterator: Iterator[pd.Series]) -> Iterator[pd.Series]:
        import joblib
        model = joblib.load("/model/classificador-produtos.joblib")
        for series in iterator:
            yield pd.Series(model.predict(series.tolist()))

    # 3. Carregando dados do volume do cluster
    df_input = spark.read.option("delimiter", ';') \
                         .option("header", "true") \
                         .csv("/data/produtos.csv")

    # 4. Filtrando dados nulos e criando ID único para as transações
    df_prepared = df_input.filter("descricao is not null") \
                          .select("descricao") \
                          .withColumn("id", monotonically_increasing_id())

    # 5. Execução da UDF para classificar as descrições dos produtos em lote
    df_predicted = df_prepared.dropna().withColumn("predict", predict("descricao"))    

    # 6. Salvando o resultado das predições finais
    df_predicted.limit(10).write.mode('overwrite').parquet("/data/predicoes_output")

with DAG(
    'spark-connect-inference-dag',
    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='Inferência em lote de ML com Spark Connect no Spark 4.1.2',
    schedule=None,
    start_date=datetime(2026, 1, 1),
    catchup=False,
    tags=['mlops', 'pyspark', 'connect'],
) as dag:

    run_batch_inference()
```

### Configurando a Conexão no Airflow UI

Antes de executar a DAG, você deve registrar a conexão de destino no console de administração do Airflow:

1. Acesse o painel administrativo do Airflow em <http://localhost:8080>.
2. Vá em **Admin > Connections** e clique em **Add a new record**.
3. Configure os seguintes parâmetros:
   * **Connection ID:** `spark_connect_default`
   * **Connection Type:** `Apache Spark Connect`
   * **Host:** `spark-connect`
   * **Port:** `15002`
4. Salve a conexão.

Agora, basta ativar a DAG `spark-connect-inference-dag` e dispará-la manualmente para verificar as predições escritas na pasta local `./data/predicoes_output`.

***

## 9.4.4 Kubernetes

No Kubernetes, a orquestração do Spark Connect é facilitada por operadores nativos como o **Apache Spark Kubernetes Operator** (oficial da ASF) ou o **Stackable Operator**. Estes operadores permitem subir o driver do Spark Connect como um pod de longa duração exposto por serviços internos gRPC. Os pipelines do Airflow rodam como pods efêmeros e submetem jobs leves ao servidor do Spark Connect por meio da rede interna do Kubernetes, aproveitando o auto-escalonamento de pods executores gerenciado diretamente pelo Kubernetes.

***

## 9.4.5 Cloud

Em ambientes corporativos de produção rodando em nuvem pública (AWS, GCP ou Azure), as empresas geralmente optam por não gerenciar a infraestrutura do Spark e Spark Connect diretamente:

* **Databricks:** Fornece suporte nativo completo ao Spark Connect através do gRPC em seus clusters gerenciados, além da *Jobs API* para disparo via Airflow.
* **AWS EMR & Google Cloud Dataproc:** Disponibilizam conectores nativos e endpoints gRPC seguros integrados ao ecossistema do AWS MWAA e Cloud Composer, simplificando pipelines corporativos de larga escala.

***

## 9.4.6 Gestão de Dependências de Python no Apache Spark

Um dos maiores desafios operacionais em pipelines de dados distribuídos com PySpark é garantir que as dependências e pacotes Python (como `pandas`, `scikit-learn`, `joblib` ou outras bibliotecas) estejam disponíveis e sejam idênticos em todas as máquinas executoras (*worker nodes*) do cluster.

Se uma biblioteca importada dentro de uma UDF (User Defined Function) não estiver instalada nos workers, a tarefa falhará em tempo de execução. Para resolver isso, existem três abordagens principais utilizadas na indústria:

### 1. Imagens Docker Customizadas (Recomendado para Kubernetes)

Quando o Spark roda sobre Kubernetes ou em containers Docker, a forma mais segura e robusta é criar uma imagem Docker base contendo todas as dependências pré-instaladas (tanto no driver quanto nos workers).

* **Vantagem:** Garante 100% de consistência do ambiente de execução e tempo de inicialização rápido para as tarefas.
* **Como é feito:** Define-se um arquivo `requirements.txt` e executa-se o `pip install` durante o build da imagem Docker do Spark (como fizemos no Dockerfile do Spark Connect deste laboratório).

### 2. Ambientes Virtuais Compactados (`venv-pack` ou `conda-pack`)

Em clusters tradicionais executados diretamente sobre máquinas virtuais (como AWS EMR ou ambientes baseados em YARN), instalar pacotes globalmente em cada nó pode ser inviável. Nesse cenário, empacotamos um ambiente virtual Python inteiro em um arquivo compactado `.tar.gz` e o enviamos dinamicamente para o Spark:

1. **Criação do ambiente virtual local:**

   ```bash
   python3 -m venv meu_env
   source meu_env/bin/activate
   pip install pandas scikit-learn joblib venv-pack
   ```
2. **Empacotamento do ambiente virtual:**

   ```bash
   venv-pack -o meu_env.tar.gz
   ```
3. **Distribuição para o Spark:** Ao disparar o job, passamos a configuração `spark.archives` (ou o parâmetro `--archives` no `spark-submit`), especificando o caminho do arquivo e um alias para onde ele deve ser extraído nos workers:

   ```python
   spark.conf.set("spark.archives", "meu_env.tar.gz#environment")
   spark.conf.set("spark.pyspark.python", "./environment/bin/python")
   ```

### 3. Distribuição de Módulos Puros (`--py-files`)

Para dependências simples contendo apenas código Python proprietário (arquivos `.py` avulsos ou empacotados em um arquivo `.zip` ou `.egg`), que não necessitam de compilação C externa:

* **Utilização:** Podemos adicionar esses arquivos diretamente à SparkSession activa para que sejam transmitidos sob demanda para o driver e workers:

  ```python
  spark.sparkContext.addPyFile("meu_modulo_auxiliar.py")
  ```

### 4. Gestão de Dependências em Serviços de Cloud Gerenciados

> \[!ATENÇÃO] **Isolamento de Internet em Produção:** Em ambientes produtivos reais, é uma prática recomendada (e frequentemente obrigatória por segurança corporativa) **não depender da internet** para a instalação dinâmica de pacotes em tempo de execução.
>
> * **Segurança:** Clusters de produção rodam em subredes privadas sem acesso direto à internet externa para mitigar vazamentos de dados.
> * **Confiabilidade:** Se o repositório PyPI estiver fora do ar ou houver instabilidade de rede, seu pipeline de dados crítico falhará.
> * **Reprodutibilidade:** Instalações dinâmicas podem sofrer com *dependency drift* (mudanças silenciosas em subdependências não fixadas), quebrando códigos estáveis.
> * **Custo e Performance:** Baixar pacotes sob demanda em tempo de execução desperdiça tempo e poder computacional de clusters caros. Sempre prefira empacotar previamente as bibliotecas via Imagens Docker ou ambientes virtuais (`venv-pack`) e armazená-los em storages privados (como S3, GCS ou Artifact Registry interno).

Em ambientes corporativos reais, a escolha de como gerenciar dependências depende fortemente de qual plataforma gerenciada de Spark está sendo utilizada:

* **Databricks:**
  * **Notebook-Scoped Libraries:** Para testes dinâmicos e desenvolvimento ágil, é possível rodar o comando `%pip install <pacote>` em uma célula de notebook. Isso instala a biblioteca de forma isolada na sessão atual do notebook e a propaga automaticamente e de forma transparente para todos os workers do cluster sem interferir com outros usuários ou notebooks.
  * **Cluster-Scoped Libraries:** Para produção, as dependências são registradas de forma centralizada através da aba "Libraries" do cluster (ou via APIs/Terraform), sendo instaladas permanentemente em todos os nós.
* **AWS EMR (Elastic MapReduce):**
  * **Bootstrap Actions:** Durante a criação do cluster, um script Bash personalizado (como `sudo pip3 install pandas`) é especificado para instalar dependências em cada nó durante a inicialização do cluster.
  * **Amazon S3 & venv-pack:** Para deploys de produção isolados, é comum salvar o ambiente virtual `.tar.gz` empacotado no S3 e apontá-lo nos jobs por meio da configuração `--archives s3://meu-bucket/meu_env.tar.gz#environment`.
* **Google Cloud Dataproc:**
  * **Initialization Actions:** Equivalente às bootstrap actions do EMR, executa scripts no startup do cluster para instalar pacotes via pip/conda em todos os nós do cluster.
  * **Dataproc Serverless & Custom Containers:** Para execuções serverless ou Kubernetes (GKE), a prática recomendada é apontar uma imagem Docker personalizada construída previamente com todas as dependências pré-instaladas.

No ecossistema moderno do **Spark Connect**, a gestão de dependências foi otimizada para ser isolada por sessão (*session-based isolation*), permitindo que diferentes usuários ou pipelines enviem seus próprios pacotes `.zip` ou arquivos `.py` dinamicamente sem afetar o estado global de execução das outras aplicações do cluster.
