Como criar um Schedule de Spark Job em Microsoft Fabric: passo a passo
Automatizar a execução de Spark Jobs no Microsoft Fabric é útil para pipelines ETL, preparação de dados e tarefas recorrentes. Este tutorial mostra como criar um Spark Job, configurá‑lo com parâmetros, agendá‑lo e gerir erros para garantir execuções fiáveis. Vamos cobrir o porquê de cada opção e dar exemplos práticos com números plausíveis (memória, cores, retries) para que possas aplicar diretamente.
Pré‑requisitos
- Conta com acesso a um workspace no Microsoft Fabric e permissões de contributor. Sem esta permissão não consegues criar Jobs nem agendar triggers.
- Um Notebook Spark funcional no workspace (PySpark ou Spark SQL). Se já tens um notebook que lê de OneLake/Lakehouse e escreve parquet, estás quase pronto.
- Conhecimentos básicos de PySpark / Spark SQL e da UI do Fabric. Saber ler logs e ajustar memory/cores ajuda a evitar OOM.
- Recomendação prática: testa com um dataset amostra (10k–100k linhas) antes de ir para produção com milhões de linhas.
Passo 1: Preparar um Notebook com parâmetros
Explica o porquê: usar parâmetros permite reutilizar o mesmo Notebook para diferentes inputs, datas ou ambientes (dev/prod). Torna também o debug mais simples — alteras apenas os argumentos no Job UI, não o código.
# Exemplo PySpark no Notebook
from pyspark.sql import SparkSession
from datetime import datetime
# widgets para parametrização no Fabric
dbutils.widgets.text("input_path", "/lakehouse/default/mydata")
input_path = dbutils.widgets.get("input_path")
dbutils.widgets.text("output_root", "/lakehouse/default/output")
output_root = dbutils.widgets.get("output_root")
spark = SparkSession.builder.getOrCreate()
df = spark.read.format("parquet").load(input_path)
# pequena transformação
df2 = df.filter("value IS NOT NULL")
out_path = f"{output_root}/{datetime.now().strftime('%Y%m%d_%H%M%S')}"
df2.write.mode("overwrite").parquet(out_path)
print(f"Wrote to {out_path}")
Exemplo concreto: num ambiente de testes define input_path = /lakehouse/dev/sample (≈50k linhas) e output_root = /lakehouse/dev/out. Num ambiente de produção usa caminhos com partições por data.
Passo 2: Criar um Spark Job no workspace
Explicação simples: um Spark Job é um recurso gerido que executa um Notebook com uma pool de execução. No Fabric UI, abre a secção Jobs / Spark e escolhe "Create new job". Indica o Notebook e a Spark pool (ex.: pool com 2 workers, cada um com 4 vCPU e 8 GB RAM).
Exemplo prático: para um dataset médio (1–10M linhas) começa com driverMemory=4g, executorCores=2 e 4 executors; depois ajusta conforme a utilização real. Regista a configuração inicial para comparar custos.
Passo 3: Definir parâmetros e configurações do Job
Porquê: parametrizar o Job permite alterar input_path ou outras opções sem editar o Notebook. No formulário do Job, adiciona os widgets/args correspondentes — por exemplo input_path e output_root. Aqui também defines configs do Spark (driverMemory, executorMemory, executorCores) e tags para billing.
# Exemplo de parâmetros no Job UI
input_path = /lakehouse/default/mydata
output_root = /lakehouse/default/output
-- Spark configs --
driverMemory = 4g
executorMemory = 8g
executorCores = 2
numExecutors = 4
Nota: documenta as escolhas (ex.: "executorMemory 8g para 2 vCPU por executor"), porque isso afeta custos e performance. Se vires GC excessivo, aumenta memória ou reduz partições do shuffle.
Passo 4: Configurar o agendamento (schedule)
Explicação: define quando o Job corre automaticamente — diário, horário ou cron. No Job, escolhe "Schedule" e configura um trigger recorrente. Para evitar sobreposição, ativa "Max concurrent runs = 1" ou define políticas de retries e backoff.
# Exemplos de opções comuns no UI
Schedule: Recurring daily at 02:00
Timezone: Europe/Lisbon
Retry policy: 2 retries, backoff 5 minutes (exp. backoff opcional)
Max concurrent runs: 1
Exemplo concreto: agenda uma execução diária às 02:00 para jobs ETL que processam dados do dia anterior. Para pipelines hora a hora usa cron (ex.: "0 * * * *" para no início de cada hora). Se o tempo esperado por run for 30–45 min, evita um trigger a cada 15 minutos.
Passo 5: Adicionar notificações e estratégias de falha
Porquê: receber alertas e ter retries melhora a fiabilidade operacional. No Job, configura e‑mail/webhook em Notifications e define a política de retries. Regista também o estado do job em OneLake para auditoria e reconciliação.
# Boas práticas
- Enable email on failure para a equipa de operações
- Set 2 retries com exponential backoff (ex.: 5m, 15m)
- Write job status to /lakehouse/default/job_status as parquet/json
Exemplo de escrita de estado no Notebook (simples): escreve um ficheiro JSON com status, start_time, end_time e rows_processed para permitir dashboards de monitorização.
Passo 6: Testar manualmente antes do agendamento
Explicação: executa o Job manualmente com parâmetros de teste para confirmar que o Notebook e as ligações a OneLake/Lakehouse funcionam. Observa o tempo de execução (por exemplo 12 min na primeira corrida, 8 min em rerun) e ajusta recursos conforme necessário.
Verificar o resultado
Confirma que o Job correu e produziu output:
- Na secção Jobs vê o histórico de runs e status (Success/Failed). Verifica tempos de start/end e duration (úteis para SLAs).
- Verifica os logs para mensagens, warnings e exceptions; copia stack traces relevantes para o sistema de incidentes.
- Confirma que os ficheiros/parquet foram escritos no caminho especificado em OneLake/Lakehouse e valida tamanho/partições (ex.: 3 ficheiros parquet, 120 MB no total).
- Verifica notificações/e‑mail de falha se configuradas e o registo em /lakehouse/default/job_status.
Conclusão
Ao criar um Spark Job agendado no Microsoft Fabric com parâmetros, agendamento e notificações, automatizas tarefas ETL repetitivas de forma robusta. Próximos passos: integrar o Job numa pipeline mais ampla, adicionar testes unitários no Notebook e usar métricas de execução para optimizar custos. Dica prática: começa por agendar em horários fora de pico, monitoriza as primeiras 7–14 execuções e ajusta recursos e retries conforme o padrão de falhas e latências observadas.