De volta ao blog

Como Construir um Pipeline de Web Scraping Agendado Com Airflow e Scrapeless

Daniel Kim
Daniel Kim

Lead Scraping Automation Engineer

10-Sep-2026

TL;DR:

  • Um pipeline de scraping no Airflow consiste em três tarefas — buscar, extrair, armazenar — conectadas pela passagem de valores de retorno, e a execução completa aqui terminou em 25,2 segundos com state=success.
  • O Airflow 3 renomeou o parâmetro de agendamento. schedule_interval= gera TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval', então a maioria dos tutoriais escritos para o Airflow 2 para na definição do DAG.
  • Instale o Airflow com seu arquivo de restrições. Sem um, pip install apache-airflow==3.1.2 gastou mais de 13 minutos na resolução de dependências aqui; com o arquivo de restrições, ele foi concluído.
  • Você não precisa do Docker para executar um DAG. airflow dags test executa todo o gráfico no processo contra o SQLite.
  • Valores passados entre as tarefas são serializados no banco de dados de metadados — o HTML da tarefa de busca caiu lá como 53.889 bytes contra 2.154 para os registros analisados.
  • Executar o mesmo DAG duas vezes deixou 20 linhas em vez de 40, porque a gravação é um upsert baseado no título do produto.
  • Comece a agendar contra páginas renderizadas no plano gratuito da Scrapeless.

Um scraper que roda quando você lembra de executá-lo é um script. Transformá-lo em algo que produz uma tabela datada toda manhã significa decidir o que uma segunda execução faz com as linhas de ontem e onde as credenciais estão. O Airflow lida com o agendamento e a conexão das tarefas; essas duas decisões continuam sendo suas.

O pipeline abaixo coleta um catálogo público de livros em um agendamento diário: uma busca renderizada através da API Universal Scraping, uma análise em registros e um upsert no SQLite. Cada número vem da execução no Airflow 3.1.2.

Pipeline de Olho

Estágio Tarefa Entrada Saída Resultado verificado
Buscar fetch URL da categoria HTML renderizado 50.403 caracteres
Extrair extract string HTML lista de registros 20 registros
Armazenar store lista de registros contagem de linhas 20 linhas no SQLite

O fluxo é buscar → extrair → armazenar, e cada estágio entrega seu valor de retorno ao próximo. Manter os limites nesses três pontos permite que você re-execute um estágio sem repetir os outros, e mantém a lógica de análise testável sem uma chamada de rede.

Instale o Airflow Sem a Longa Resolução

O Airflow depende de um grande grafo fixado, e instalá-lo sem o arquivo de restrições correspondente deixa o pip fazendo retrocesso através de combinações de versões. Um simples pip install apache-airflow==3.1.2 aqui ainda estava resolvendo após 13 minutos; a mesma instalação com o arquivo de restrições terminou normalmente.

bash Copy
python3 -m venv .venv
./.venv/bin/pip install "apache-airflow==3.1.2" \
  --constraint "https://raw.githubusercontent.com/apache/airflow/constraints-3.1.2/constraints-3.12.txt"
./.venv/bin/pip install "parsel==1.10.0" "requests==2.32.5"

A URL de restrições codifica tanto a versão do Airflow quanto a versão do Python, e a referência de instalação do projeto considera isso parte da instalação em vez de uma otimização. Aponte AIRFLOW_HOME em algum lugar explícito e crie o banco de dados de metadados:

bash Copy
export AIRFLOW_HOME="$PWD/home"
./.venv/bin/airflow db migrate
text Copy
Airflow database tables created
DB: sqlite:////root/verify-airflow/home/airflow.db
Database migrating done!

O SQLite é o back-end padrão e é suficiente para um pipeline de máquina única. O agendador precisa de um banco de dados concorrente para executar tarefas em paralelo, mas nada abaixo requer isso.

Estágio 1: Buscar uma Página Renderizada

O primeiro estágio é aquele que decide se o restante do pipeline vê dados ou não. Um cliente HTTP simples retorna o que o servidor envia antes que o JavaScript seja executado, então a tarefa de busca chama a API Universal Scraping, que retorna o documento renderizado.

python Copy
@task
def fetch() -> str:
    payload = json.dumps({
        "actor": "unlocker.webunlocker",
        "input": {"url": CATEGORY_URL, "js_render": True, "headless": False},
    }).encode()
    request = urllib.request.Request(
        "https://api.scrapeless.com/api/v2/unlocker/request",
        data=payload,
        headers={"Content-Type": "application/json",
                 "x-api-token": os.environ["SCRAPELESS_API_KEY"]},
    )
    with urllib.request.urlopen(request, timeout=180) as response:
        body = json.loads(response.read())
    html = body["data"]
    print(f"fetched {len(html)} chars")
    return html
text Copy
fetched 50403 chars

O envelope de resposta é {"code": ..., "data": ...}, onde data transporta o HTML como uma string. Como a API retorna JSON, o documento chega já decodificado como UTF-8 — o £ em cada preço sobrevive sem nenhum tratamento de codificação na tarefa.

O token da API vem do ambiente em vez do arquivo do DAG. O Airflow lê o ambiente do processo quando executa uma tarefa, então exportá-lo no shell que roda o pipeline mantém a credencial fora do repositório. As Variáveis e Conexões do Airflow são a resposta mais completa uma vez que mais de um DAG precise disso.

Estágio 2: Extrair Registros

O estágio de extração pega uma string e retorna registros. Ele não toca na rede, o que significa que pode ser exercitado contra uma fixture salva sempre que a marcação mudar.

python Copy
@task
def extract(html: str) -> list[dict]:
    sel = Selector(text=html)
    rows = []
    for card in sel.css("article.product_pod"):
        rows.append({
            "title": card.css("h3 a::attr(title)").get(),
            "price": card.css("p.price_color::text").get(),
            "rating": (card.css("p.star-rating::attr(class)").get() or "").replace("star-rating", "").strip(),
            "in_stock": "In stock" in (card.css("p.instock.availability::text").getall() or [""])[-1],
        })
    print(f"extracted {len(rows)} records")
    return rows
text Copy
extracted 20 records

Selecionar cada cartão primeiro e depois consultar dentro dele mantém os campos do mesmo produto juntos. Selecionar todos os títulos e todos os preços como duas listas planas e zipá-las produz pares silenciosamente não correspondidos no momento em que um cartão está faltando um campo.

Estágio 3: Armazenar com um Upsert

Um pipeline agendado grava as mesmas linhas novamente amanhã, então a gravação deve ser definida em termos do que torna um registro único em vez de anexar cegamente.

python Copy
@task
def store(rows: list[dict]) -> int:
    conn = sqlite3.connect(DB_PATH)
    conn.execute("""CREATE TABLE IF NOT EXISTS books (
        title TEXT PRIMARY KEY, price TEXT, rating TEXT,
        in_stock INTEGER, scraped_at TEXT)""")
    now = datetime.utcnow().isoformat(timespec="seconds")
    conn.executemany(
        "INSERT OR REPLACE INTO books VALUES (?,?,?,?,?)",
        [(r["title"], r["price"], r["rating"], int(r["in_stock"]), now) for r in rows],
    )
    conn.commit()
    count = conn.execute("SELECT COUNT(*) FROM books").fetchone()[0]
    conn.close()
    print(f"stored; table now holds {count} rows")
    return count
text Copy
stored; table now holds 20 rows

title é a chave primária, e o comportamento INSERT OR REPLACE do SQLite exclui a linha conflitante antes de inserir a nova. A tabela armazenada após a execução:

text Copy
title                                  price     rating stock
A Murder in Time                       £16.64    One    1
A Study in Scarlet (Sherlock Holmes    £16.73    Two    1
A Time of Torment (Charlie Parker #1   £48.35    Five   1
Boar Island (Anna Pigeon #19)          £59.48    Three  1
Delivering the Truth (Quaker Midwife   £20.89    Four   1
Hide Away (Eve Duncan #20)             £11.84    One    1

Conectando os Estágios em um DAG

A declaração do DAG contém o cronograma; a última linha declara a cadeia de dependência passando valores de retorno.

python Copy
@dag(
    dag_id="books_pipeline",
    schedule="0 6 * * *",
    start_date=pendulum.datetime(2026, 9, 1, tz="UTC"),
    catchup=False,
    tags=["scraping"],
)
def books_pipeline():
    ...
    store(extract(fetch()))


books_pipeline()

schedule leva uma expressão cron padrão, então 0 6 * * * é 06:00 diariamente no fuso horário do DAG. Configurar start_date com uma zona explícita é mais importante do que parece: um DAG fixado em uma zona que observa o horário de verão muda seu real tempo de execução duas vezes por ano, e os deslocamentos vêm de o banco de dados de fusos horários da IANA. Fixar em UTC mantém o intervalo fixo. catchup=False é importante para scrapers especificamente: com ele definido para True, implantar um DAG cujo start_date é um mês atrás coloca em espera uma execução para cada intervalo perdido, e um scraper não pode recuperar uma página como estava três semanas atrás de qualquer forma.

Escrever store(extract(fetch())) é o grafo de dependência completo. O Airflow lê a aninhamento de chamadas e constrói as arestas.

Execute Sem Docker

A maioria dos tutoriais do Airflow começa com um arquivo Docker Compose e um contêiner Postgres. Para desenvolver um DAG, airflow dags test executa todo o grafo em processo, sem iniciar o agendador, o servidor web ou um contêiner.

bash Copy
export AIRFLOW_HOME="$PWD/home"
export SCRAPELESS_API_KEY="your-api-key"
./.venv/bin/airflow dags test books_pipeline
text Copy
fetched 50403 chars
extracted 20 records
stored; table now holds 20 rows
DagRun Finished: dag_id=books_pipeline, run_duration=25.211182, state=success

Essa execução é o loop de feedback mais rápido disponível enquanto a lógica de análise ainda está em movimento. O agendador é o que você inicia uma vez que o DAG está fazendo o que você deseja.

Agendando contra páginas que renderizam do lado do cliente? O plano gratuito Scrapeless cobre solicitações suficientes para que um DAG diário funcione de ponta a ponta.

Executar Duas Vezes Não Deve Dobrar Suas Linhas

Um pipeline agendado é executado sem supervisão, então o teste interessante é a segunda execução, em vez da primeira. Executando o mesmo DAG novamente:

text Copy
stored; table now holds 20 rows

Vinte, não quarenta. O upsert chaveado em title substituiu cada linha e atualizou seu scraped_at, o que torna o DAG seguro para ser acionado manualmente enquanto depura. Um INSERT simples teria dobrado a tabela e deixado nenhuma maneira de dizer qual cópia estava atual.

Se a história importa — rastreando como um preço se move — a chave se torna o par de título e data de execução, em vez do título sozinho, e a linha de ontem permanece.

A Fronteira do XCom

Os valores de retorno transitam entre as tarefas através do XCom, e os valores do XCom são serializados no banco de dados de metadados. Os três estágios aqui armazenados:

text Copy
task_id      serialized bytes
fetch        53889
extract      2154
store        2

O HTML é 25 vezes o tamanho dos registros analisados, e fica no banco de dados do Airflow, em vez de na memória. Isso é aceitável nesse tamanho e deixa de ser aceitável à medida que as páginas aumentam ou as execuções se tornam mais frequentes.

A abordagem que escala é analisar cedo e passar pequenas quantidades: fazer a tarefa de busca salvar o documento bruto no armazenamento de objetos e retornar uma chave, e então permitir que a extração o leia de volta. As fronteiras de estágio permanecem idênticas e o banco de dados de metadados continua armazenando valores do tamanho de registros, em vez de documentos.

O Código do Airflow 2 Não Funcionará Aqui Sem Alterações

Duas mudanças chamam a atenção de quem segue um tutorial mais antigo, e apenas uma delas se anuncia.

python Copy
from airflow.decorators import dag, task   # deprecated
from airflow.sdk import dag, task          # Airflow 3

O antigo caminho de importação ainda se resolve e emite DeprecatedImportWarning: The airflow.decorators.dag attribute is deprecated. Please use 'airflow.sdk.dag'. O parâmetro de agendamento é uma falha crítica:

text Copy
TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'

schedule_interval= se tornou schedule=. Como quase todos os tutoriais publicados sobre scraping do Airflow são anteriores à versão 3, esse TypeError é a primeira coisa que muitos leitores encontram — e isso aparece no momento da análise do DAG, antes que qualquer código de scraping seja executado.

Conclusão

O pipeline consiste em três funções e um decorador: fetch retorna HTML, extract retorna registros, store retorna uma contagem, e store(extract(fetch())) é o grafo. O que o torna um pipeline em vez de um script é o upsert que sobrevive a uma segunda execução, o cronograma expresso como cron, e a fronteira do estágio que permite reanalisar sem reobter.
Comece com airflow dags test e SQLite, mantenha os documentos fora do XCom uma vez que eles cresçam, e fixe a instalação com o arquivo de restrições para que a primeira coisa que você enfrente seja a marcação em vez do resolvedor de dependências. Para o armazenamento da mesma questão, nosso pipeline DuckDB e Parquet cobre a saída em coluna, e preços lista o custo de um cronograma diário em solicitações.

Pronto para colocar um scraper em um cronograma? Comece com o plano gratuito da Scrapeless e aponte a fase de busca para seu próprio alvo.

FAQ

Q: Eu preciso do Docker para rodar o Airflow para scraping?

Não. airflow dags test <dag_id> executa todo o gráfico em processo contra o banco de dados SQLite padrão, que é como a execução de 25,2 segundos acima foi realizada. Docker e Postgres tornam-se úteis quando você deseja que o agendador funcione sem supervisão com tarefas sendo executadas em paralelo, pois SQLite não suporta as gravações concorrentes que isso exige.

Q: Onde a chave da API deve estar em um DAG do Airflow?

Não no arquivo do DAG. Lê-la do ambiente do processo mantém fora do controle de versão, e Variáveis ou Conexões do Airflow são o melhor lugar uma vez que vários DAGs precisam da mesma credencial — ambos são armazenados no banco de dados de metadados e referenciados pelo nome em vez de colados no código.

Q: Como eu paro um scraper agendado de criar linhas duplicadas?

Defina o que torna um registro único e escreva contra essa chave. Aqui title é a chave primária e a gravação é INSERT OR REPLACE, então a segunda execução deixou 20 linhas em vez de 40. Se você quer histórico em vez do estado atual, amplie a chave para incluir a data da execução para que o instantâneo de cada dia seja sua própria linha.

Q: É seguro passar HTML raspado entre tarefas do Airflow?

Em tamanhos pequenos, sim — o documento de 50.403 caracteres acima foi serializado para 53.889 bytes no banco de dados de metadados e movido entre tarefas sem problema. Isso deixa de ser uma boa ideia à medida que os documentos ou a frequência de execução crescem, pois cada valor reside nesse banco de dados. Escrever o documento em armazenamento de objetos e passar uma chave mantém os mesmos limites de tarefas com valores XCom do tamanho de registro.

Q: Por que o código do meu tutorial do Airflow gera um TypeError no schedule_interval?

Porque foi escrito para o Airflow 2. A versão 3 renomeou o parâmetro para schedule, e passar o nome antigo gera TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval' enquanto o DAG está sendo analisado. A mudança correspondente é a importação: airflow.decorators ainda funciona, mas emite um aviso, e airflow.sdk é o caminho atual.

Q: O catchup deve ser ativado para um DAG de scraping?

Geralmente não. Com catchup=True, um DAG cujo start_date está no passado enfileira uma execução por intervalo perdido assim que é implantado. Esse comportamento se adequa à reprocessamento de um conjunto de dados datado, mas um scraper lê o que a página mostra agora, então essas execuções coletariam os dados de hoje e os rotulariam com datas históricas.

Na Scorretless, acessamos apenas dados disponíveis ao público, enquanto cumprem estritamente as leis, regulamentos e políticas de privacidade do site aplicáveis. O conteúdo deste blog é apenas para fins de demonstração e não envolve atividades ilegais ou infratoras. Não temos garantias e negamos toda a responsabilidade pelo uso de informações deste blog ou links de terceiros. Antes de se envolver em qualquer atividade de raspagem, consulte seu consultor jurídico e revise os termos de serviço do site de destino ou obtenha as permissões necessárias.

Artigos mais populares

Catálogo