Prefect 3 em Python: Guia Completo de Orquestração de Pipelines em 2026
Guia prático de Prefect 3 em Python para 2026: instalação, flows assíncronos, tarefas autônomas, deployments com work pools, pipeline ETL com Pandas e DuckDB, e comparativo com Airflow e Dagster.
O Prefect 3 é um framework open source de orquestração de workflows em Python que transforma qualquer função Python em uma tarefa monitorada, agendada e resiliente, com suporte nativo a execução assíncrona, tarefas autônomas e uma nova engine baseada em events introduzida na versão 3.0. Neste guia prático de 2026, você vai aprender como instalar o Prefect 3, escrever seu primeiro flow, criar deployments com work pools, orquestrar pipelines de ETL de forma assíncrona e comparar a ferramenta com Airflow e Dagster, tudo com exemplos reproduzíveis alinhados à versão 3.4 estável.
Prefect 3 (release estável 3.0 em setembro de 2024, atual 3.4.x em 2026) reescreveu a engine em código async-first, permitindo executar milhares de tarefas por segundo sem overhead de subprocess.
Você define pipelines com apenas dois decorators, @flow e @task, e ganha retry, cache, logging estruturado e observabilidade sem configuração adicional.
A nova arquitetura de work pools e workers substituiu os agents do Prefect 2, permitindo deploy híbrido em Docker, Kubernetes, ECS ou processo local.
Tarefas autônomas (autonomous tasks) permitem invocar qualquer @task como job de fila sem precisar de um flow envolvente.
Comparado ao Airflow, o Prefect 3 é Python-nativo (sem DAGs declarativas), suporta parâmetros dinâmicos em runtime e roda tanto self-hosted quanto na Prefect Cloud com plano gratuito.
Integração direta com Pydantic 2, Pandas, Polars, DuckDB e ferramentas de qualidade como Great Expectations facilita a construção de pipelines de dados sólidos.
O que é o Prefect 3 e o que mudou desde a versão 2
O Prefect é um sistema de orquestração de workflows criado pela Prefect Technologies que trata pipelines como código Python comum, sem a necessidade de escrever DAGs declarativos. A versão 3, lançada como estável em setembro de 2024 e evoluída até o release 3.4 em 2026, é a maior reescrita desde o projeto original. Segundo as notas oficiais da release 3.0, o novo engine roda transações de tarefas nativamente em asyncio, elimina a arquitetura de agents em favor de workers específicos por infraestrutura e introduz o conceito de tarefas autônomas que podem ser disparadas como jobs de fila.
Na prática, isso significa três mudanças que impactam desenvolvedores todos os dias. Primeiro, você não precisa mais registrar blocos de storage separados para cada deployment; o novo prefect.yaml descreve tudo em um único arquivo declarativo. Segundo, o overhead por task caiu de aproximadamente 200 ms para menos de 20 ms, permitindo pipelines com dezenas de milhares de tarefas leves. Terceiro, a integração com eventos e automações tornou trivial acionar flows a partir de webhooks, mensagens de fila ou mudanças em arquivos S3.
Se você já usava Prefect 2, o guia de migração é curto. A maior parte do código continua funcionando após atualizar o prefect para a série 3.x, mas Deployment.build_from_flow foi removido e substituído por flow.deploy(). (Eu subi um projeto pequeno em uma tarde; um monolito de flows antigos levou dois dias, principalmente para revisar blocos de storage.)
Como instalar o Prefect 3 em Python
Então, vamos ao que interessa. A instalação é feita via pip ou uv e requer Python 3.9 ou superior. Recomendo criar um ambiente virtual isolado para evitar conflitos com outras ferramentas de orquestração já instaladas no sistema:
# criar e ativar um ambiente virtual
python -m venv .venv
source .venv/bin/activate # Linux/macOS
# .venv\Scripts\activate # Windows
# instalar a última versão estável do Prefect 3
pip install "prefect>=3.4,<4.0"
# verificar a instalação
prefect version
A saída deve indicar algo semelhante a Version: 3.4.x. Para o backend, você tem duas opções: rodar o servidor local (Prefect Server) ou usar a Prefect Cloud gratuita. Para experimentos e desenvolvimento, o servidor local costuma ser suficiente:
# em um terminal separado, inicie o servidor
prefect server start
# em outro terminal, aponte o CLI para o servidor local
prefect config set PREFECT_API_URL="http://127.0.0.1:4200/api"
Se preferir a Prefect Cloud, execute prefect cloud login e escolha seu workspace. Em ambos os casos, a UI web disponível em http://127.0.0.1:4200 (self-hosted) mostra em tempo real todos os flows, tasks, deployments e logs. Se você já leu nosso guia sobre Pydantic 2 para pipelines de dados, vai notar que o Prefect 3 usa a mesma versão do Pydantic para validar parâmetros de flows.
Seu primeiro flow: tasks, retry e cache
No Prefect, um flow é a unidade de trabalho orquestrada e pode conter uma ou várias tasks. A diferença principal entre task e flow é que tasks têm cache automático, retries granulares e são registradas individualmente na UI, enquanto flows organizam a lógica de negócio e podem chamar outros flows como subflows. O exemplo abaixo baixa dados da API pública do IBGE, transforma-os com Pandas e imprime um resumo:
from datetime import timedelta
import httpx
import pandas as pd
from prefect import flow, task
from prefect.cache_policies import INPUTS
@task(retries=3, retry_delay_seconds=[1, 5, 15], cache_policy=INPUTS,
cache_expiration=timedelta(hours=6))
def baixar_populacao(uf: str) -> dict:
"""Baixa estimativas populacionais do IBGE para uma UF."""
url = f"https://servicodados.ibge.gov.br/api/v1/localidades/estados/{uf}/municipios"
resp = httpx.get(url, timeout=30.0)
resp.raise_for_status()
return resp.json()
@task
def resumir(municipios: list[dict]) -> pd.DataFrame:
df = pd.DataFrame(municipios)
df["nome_micro"] = df["microrregiao"].apply(lambda x: x["nome"])
return df.groupby("nome_micro").size().reset_index(name="qtd_municipios")
@flow(name="ibge-populacao", log_prints=True)
def pipeline_populacao(ufs: list[str] = ["SP", "RJ", "MG"]):
for uf in ufs:
dados = baixar_populacao(uf)
resumo = resumir(dados)
print(f"{uf}: {len(dados)} municípios em {len(resumo)} microrregiões")
if __name__ == "__main__":
pipeline_populacao()
Execute com python pipeline.py. Ao abrir a UI, você verá o run do flow com cada task listada, seus tempos de execução e o log estruturado. Se a chamada HTTP falhar, o Prefect tenta novamente 3 vezes com backoff exponencial (1s, 5s, 15s). O parâmetro cache_policy=INPUTS garante que, se você rodar duas vezes com o mesmo uf, o resultado é servido do cache durante 6 horas, sem repetir a chamada de rede.
Execução assíncrona e tarefas autônomas
O Prefect 3 é async-first: qualquer @task ou @flow pode ser declarado como coroutine e será executado pela engine sem subprocess overhead. Isso é particularmente útil para pipelines I/O-bound, como raspagem de APIs em paralelo:
import asyncio
import httpx
from prefect import flow, task
@task(retries=2)
async def fetch(url: str) -> int:
async with httpx.AsyncClient() as client:
r = await client.get(url, timeout=15.0)
return len(r.content)
@flow
async def crawl(urls: list[str]) -> dict[str, int]:
resultados = await asyncio.gather(*(fetch(u) for u in urls))
return dict(zip(urls, resultados))
if __name__ == "__main__":
asyncio.run(crawl([
"https://python.org",
"https://numpy.org",
"https://pandas.pydata.org",
]))
Já as tarefas autônomas, novidade da versão 3, permitem executar uma task sem envolvê-la em um flow. Elas ficam em uma fila e são processadas por um task worker, atuando como jobs de background (perfeito para arquiteturas orientadas a eventos):
from prefect import task
from prefect.task_worker import serve
@task(log_prints=True)
def enviar_email(destino: str, assunto: str):
print(f"enviando email para {destino}: {assunto}")
if __name__ == "__main__":
# Em um processo dedicado, sirva a task como worker
serve(enviar_email)
Em outro processo, chame enviar_email.delay("[email protected]", "Bem-vindo!") e o worker processa a mensagem. Honestamente, essa arquitetura substitui Celery ou RQ em muitos casos, com a vantagem de já vir com a UI e observabilidade do Prefect. Se você trabalha com validação de dados de entrada, vale combinar as tarefas autônomas com a nossa abordagem de testes de dados com Great Expectations para validar payloads antes de processá-los.
Deployments, work pools e workers
Um deployment é a forma de agendar um flow para rodar em uma infraestrutura definida. Ele empacota o código, os parâmetros padrão, o cronograma e o work pool onde será executado. A partir do Prefect 3, o modelo é sempre pull-based: um worker roda em sua infra, faz polling do work pool e executa flows quando eles aparecem.
O comando abaixo cria um deployment agendado para rodar diariamente às 6h da manhã em um work pool Docker:
# criar o work pool (só uma vez)
prefect work-pool create --type docker meu-pool-docker
# no código Python
if __name__ == "__main__":
pipeline_populacao.deploy(
name="ibge-diario",
work_pool_name="meu-pool-docker",
cron="0 6 * * *",
image="python:3.12-slim",
parameters={"ufs": ["SP", "RJ", "MG", "BA", "RS"]},
)
Em produção, você inicia o worker em um servidor, container ou pod Kubernetes:
prefect worker start --pool meu-pool-docker
Os tipos de work pool disponíveis em 2026 incluem process, docker, kubernetes, ecs, cloud-run, vertex-ai e azure-container-instance. A tabela a seguir compara quando usar cada um:
Work pool
Melhor para
Isolamento
Escala automática
process
Desenvolvimento local, POCs
Baixo
Não
docker
Deploy em VM única com múltiplos flows
Container por run
Não
kubernetes
Ambientes corporativos com K8s existente
Pod por run
Sim (via HPA/KEDA)
ecs / cloud-run
Serverless em AWS ou GCP sem gerenciar cluster
Task/Container por run
Sim, nativa
Prefect 3 vs Airflow vs Dagster: qual escolher?
A escolha entre Prefect, Apache Airflow e Dagster depende do estilo do time e do estágio do produto. Airflow é o padrão de mercado há uma década, com ecossistema enorme de operators, mas exige aprender DAGs declarativas e tem overhead significativo por task. Dagster foca em software-defined assets e é excelente para modelagem de linhagem, mas é bem mais opinativo. Prefect 3 fica no meio: Python-nativo como Dagster, porém com menos abstrações, e mais leve que Airflow.
Critério
Prefect 3
Airflow 2.x
Dagster 1.x
Modelo de definição
Funções Python + decorators
DAG declarativa
Assets e ops
Suporte async nativo
Sim (engine async-first)
Parcial (2.7+)
Parcial
Parâmetros dinâmicos em runtime
Sim, sem restrições
Limitado (TaskGroups)
Sim (partitions)
Overhead por task
~20 ms
~500 ms
~100 ms
Plano gerenciado gratuito
Prefect Cloud Free
Não (MWAA, Astronomer pagos)
Dagster+ Trial
Curva de aprendizagem
Baixa
Média
Média-alta
Para times pequenos, startups e pipelines ML, Prefect é geralmente a escolha mais produtiva. Para empresas com centenas de DAGs legadas, migrar para Airflow gerenciado (MWAA, Astronomer) costuma fazer sentido. E para times que enxergam pipelines como grafos de assets versionados, Dagster acaba sendo a opção mais completa. No meu último projeto de dados, começamos com Airflow porque já existia na infra; migrei um subconjunto para Prefect 3 e o feedback do time foi imediato, principalmente pela facilidade de testar flows como funções Python normais.
Pipeline ETL completo com Pandas e DuckDB
Vamos juntar tudo em um pipeline de ETL prático: extrair dados de uma API, transformar com Pandas, carregar em DuckDB e validar antes do commit final. Esse padrão é o que costumo usar em projetos reais de análise de dados com DuckDB.
from datetime import datetime, timedelta
import duckdb
import httpx
import pandas as pd
from prefect import flow, task
from prefect.cache_policies import INPUTS
@task(retries=3, retry_delay_seconds=[2, 10, 30])
def extrair_cotacoes(moeda: str = "USD-BRL") -> pd.DataFrame:
url = f"https://economia.awesomeapi.com.br/json/daily/{moeda}/30"
dados = httpx.get(url, timeout=15.0).json()
df = pd.DataFrame(dados)
df["timestamp"] = pd.to_datetime(df["timestamp"].astype(int), unit="s")
return df[["timestamp", "bid", "ask", "high", "low"]].astype(
{"bid": float, "ask": float, "high": float, "low": float}
)
@task
def transformar(df: pd.DataFrame) -> pd.DataFrame:
df["spread"] = df["ask"] - df["bid"]
df["variacao_dia"] = df["high"] - df["low"]
df["media_movel_7d"] = df["bid"].rolling(window=7, min_periods=1).mean()
return df
@task
def validar(df: pd.DataFrame) -> pd.DataFrame:
assert not df.empty, "DataFrame vazio: extração falhou"
assert (df["spread"] >= 0).all(), "Spread negativo detectado"
assert df["timestamp"].is_monotonic_decreasing, "Timestamps fora de ordem"
return df
@task
def carregar(df: pd.DataFrame, banco: str = "cotacoes.duckdb") -> int:
con = duckdb.connect(banco)
con.execute("""
CREATE TABLE IF NOT EXISTS cotacoes (
timestamp TIMESTAMP, bid DOUBLE, ask DOUBLE,
high DOUBLE, low DOUBLE, spread DOUBLE,
variacao_dia DOUBLE, media_movel_7d DOUBLE
)
""")
con.execute("INSERT INTO cotacoes SELECT * FROM df")
total = con.execute("SELECT COUNT(*) FROM cotacoes").fetchone()[0]
con.close()
return total
@flow(name="etl-cotacoes", log_prints=True)
def etl_cotacoes(moeda: str = "USD-BRL"):
bruto = extrair_cotacoes(moeda)
tratado = transformar(bruto)
validado = validar(tratado)
total = carregar(validado)
print(f"Pipeline concluído. Total de linhas no DuckDB: {total}")
if __name__ == "__main__":
etl_cotacoes()
Este pipeline segue o padrão canônico de ETL (extract, transform, validate, load) com retries automáticos apenas na etapa que fala com a rede, cache opcional para reprocessamento seguro e uma etapa explícita de validação que falha rápido se qualquer invariante for violada. A UI do Prefect mostra o grafo de dependências, tempos de execução por task e permite reexecutar apenas a etapa que falhou, poupando chamadas à API. Eu bati de frente com esse cenário há alguns meses: um limite de requisições da API travou o pipeline no meio, e sem retries granulares por task, teríamos que rodar tudo de novo.
Como monitorar flows e configurar alertas
A observabilidade nativa do Prefect 3 cobre três camadas: logs por task, events emitidos por qualquer mudança de estado e automations que reagem a esses eventos. Você pode, por exemplo, criar uma automação que envia mensagem no Slack sempre que um flow entrar no estado Failed mais de duas vezes em 15 minutos. Isso é feito pela UI em Automations → Add Automation ou via API.
Para métricas mais avançadas, o Prefect exporta um endpoint Prometheus a partir da versão 3.2. Basta habilitar em prefect.yaml ou pela variável de ambiente PREFECT_API_ENABLE_METRICS=true e apontar seu Prometheus para /api/metrics. Painéis Grafana prontos estão disponíveis no repositório prefect-community/grafana-dashboards. Nas melhores implementações que vi, os times combinam esse dashboard com alertas no Slack e um Runbook curto por flow crítico. Assim, o on-call sabe exatamente o que fazer quando algo quebra às 3h da manhã.
Boas práticas para rodar Prefect 3 em produção
Depois de operar Prefect em vários times, consolidei algumas regras que evitam problemas recorrentes. Primeiro, separe flows de bibliotecas: mantenha a lógica de negócio em pacotes Python testáveis com pytest e importe-os nos flows. Isso permite testar 90% do código sem precisar de um servidor Prefect rodando. Segundo, use tags e descrições em todos os deployments; a UI cresce rápido e sem tags fica quase impossível filtrar entre desenvolvimento, staging e produção.
Terceiro, defina timeouts explícitos em cada task com @task(timeout_seconds=300): sem isso, uma chamada HTTP travada pode segurar um worker por horas. Quarto, versione a imagem Docker usada pelos workers com hash SHA em vez de latest. Atualizações silenciosas de dependências já causaram muitos incidentes por aqui. Quinto, habilite o result_storage em flows críticos para que resultados de tasks sejam persistidos em S3 ou GCS e possam ser reutilizados sem reprocessamento. Por fim, revise os limites de concurrency do work pool para não sobrecarregar APIs externas: o Prefect suporta concurrency limits globais e por tag.
Perguntas frequentes
Qual a diferença entre Prefect e Airflow?
O Prefect 3 usa Python puro com decorators (@flow e @task) e permite parâmetros dinâmicos em runtime, enquanto o Airflow exige que você defina DAGs declarativos e tem restrições ao passar dados entre tasks. Prefect tem overhead de ~20 ms por task; o Airflow gira em torno de 500 ms.
O Prefect é gratuito?
Sim. O Prefect Core é open source (licença Apache 2.0) e você pode rodar o servidor auto-hospedado sem custos. A Prefect Cloud oferece um plano Free com até 20.000 tasks/mês e planos pagos para times maiores.
Como executo Prefect em produção?
Crie um work pool do tipo docker, kubernetes, ecs ou cloud-run, faça deploy do flow com flow.deploy() e inicie um worker naquela infraestrutura com prefect worker start. O worker faz polling do pool e executa os flows agendados.
O Prefect suporta processamento assíncrono?
Sim. A engine do Prefect 3 é async-first e você pode declarar @task async def ... e @flow async def .... Isso torna trivial paralelizar chamadas I/O com asyncio.gather, sem overhead de subprocess ou thread pool.
Preciso migrar do Prefect 2 para o Prefect 3?
A versão 2 continua recebendo correções críticas até final de 2026, mas todo o desenvolvimento novo acontece na série 3.x. A migração costuma ser rápida (poucas horas para projetos médios) e traz ganhos claros de performance e experiência de deploy.
Aprenda a usar Narwhals para escrever código Python compatível com Pandas, Polars, PyArrow, Modin, cuDF e Dask sem duplicação. Guia 2026 com exemplos práticos, benchmarks e três padrões de produção que uso no dia a dia.
Marimo é o notebook Python reativo que substitui o Jupyter: instale, use SQL com DuckDB, publique no navegador via WASM e versione tudo no Git com arquivos .py puros.
Guia prático de Ibis 12.0 em Python: escreva DataFrames uma vez e execute em DuckDB, BigQuery, Snowflake ou Polars trocando apenas o objeto de conexão.