Como Construir um Pipeline de Web Scraping Agendado Com Airflow e Scrapeless
Lead Scraping Automation Engineer
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=geraTypeError: 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.2gastou 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 testexecuta 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
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
export AIRFLOW_HOME="$PWD/home"
./.venv/bin/airflow db migrate
text
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
@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
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
@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
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
@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
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
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
@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
export AIRFLOW_HOME="$PWD/home"
export SCRAPELESS_API_KEY="your-api-key"
./.venv/bin/airflow dags test books_pipeline
text
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
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
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
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
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.



