Cómo construir un pipeline de web scraping programado con Airflow y Scrapeless
Lead Scraping Automation Engineer
TL;DR:
- Un pipeline de raspado en Airflow son tres tareas: fetch, extract, store — conectadas al pasar valores de retorno, y la ejecución completa aquí terminó en 25.2 segundos con
state=success. - Airflow 3 renombró el parámetro de programación.
schedule_interval=elevaTypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval', por lo que la mayoría de los tutoriales escritos para Airflow 2 se detienen en la definición del DAG. - Instalar Airflow con su archivo de restricciones. Sin uno,
pip install apache-airflow==3.1.2gastó más de 13 minutos en la resolución de dependencias aquí; con el archivo de restricciones, se completó. - No necesitas Docker para ejecutar un DAG.
airflow dags testejecuta todo el gráfico en proceso contra SQLite. - Los valores pasados entre tareas se serializan en la base de datos de metadatos: el HTML de la tarea de fetch aterrizó allí como 53,889 bytes frente a 2,154 para los registros analizados.
- Ejecutar el mismo DAG dos veces dejó 20 filas en lugar de 40, porque la escritura es un upsert clave en el título del producto.
- Comienza a programar contra páginas renderizadas en el plan gratuito de Scrapeless.
Un raspador que se ejecuta cuando recuerdas ejecutarlo es un script. Transformarlo en algo que produzca una tabla fechada cada mañana significa decidir qué hace una segunda ejecución con las filas de ayer y dónde residen las credenciales. Airflow maneja la programación y la conexión de tareas; esas dos decisiones siguen siendo tuyas.
El pipeline a continuación recopila un catálogo de libros público en un horario diario: una recuperación renderizada a través de la API de Raspado Universal, un análisis en registros y un upsert en SQLite. Cada número proviene de ejecutarlo en Airflow 3.1.2.
Pipeline a Primera Vista
| Etapa | Tarea | Entrada | Salida | Resultado verificado |
|---|---|---|---|---|
| Fetch | fetch |
URL de categoría | HTML renderizado | 50,403 caracteres |
| Extract | extract |
cadena HTML | lista de registros | 20 registros |
| Store | store |
lista de registros | número de filas | 20 filas en SQLite |
El flujo es fetch → extract → store, y cada etapa entrega su valor de retorno a la siguiente. Mantener los límites en esos tres puntos es lo que te permite volver a ejecutar una etapa sin repetir las demás, y lo que mantiene la lógica de parsing comprobable sin una llamada a la red.
Instalar Airflow Sin la Larga Resolución
Airflow depende de un gran gráfico fijado, y instalarlo sin el archivo de restricciones correspondiente deja a pip retrocediendo a través de combinaciones de versiones. Un pip install apache-airflow==3.1.2 simple aquí aún estaba resolviendo después de 13 minutos; la misma instalación con el archivo de restricciones terminó 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"
La URL de restricciones codifica tanto la versión de Airflow como la versión de Python, y la referencia de instalación del proyecto lo trata como parte de la instalación en lugar de una optimización. Apunta AIRFLOW_HOME en algún lugar explícito y crea la base de datos de metadatos:
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!
SQLite es el backend predeterminado y es suficiente para un pipeline de una sola máquina. El programador necesita una base de datos concurrente para ejecutar tareas en paralelo, pero nada de lo siguiente requiere eso.
Etapa 1: Obtener una Página Renderizada
La primera etapa es la que decide si el resto del pipeline ve datos en absoluto. Un cliente HTTP simple devuelve lo que el servidor envía antes de que se ejecute JavaScript, así que la tarea de fetch llama a la API de Raspado Universal, que devuelve el 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
El sobre de respuesta es {"code": ..., "data": ...}, donde data lleva el HTML como una cadena. Debido a que la API devuelve JSON, el documento llega ya decodificado como UTF-8 — el £ en cada precio sobrevive sin ningún manejo de codificación en la tarea.
El token de la API proviene del entorno en lugar del archivo DAG. Airflow lee el entorno del proceso cuando ejecuta una tarea, así que exportarlo en el shell que ejecuta el pipeline mantiene la credencial fuera del repositorio. Las Variables y Conexiones de Airflow son la respuesta más completa una vez que más de un DAG las necesita.
Etapa 2: Extraer Registros
La etapa de extracción toma una cadena y devuelve registros. No toca ninguna red, lo que significa que se puede probar contra un fixture guardado cada vez que cambia el marcado.
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
Seleccionar cada tarjeta primero y luego consultar dentro de ella mantiene los campos del mismo producto juntos. Seleccionar todos los títulos y todos los precios como dos listas planas y combinarlas produce silenciosamente pares desajustados en el momento en que falta un campo en una tarjeta.
Etapa 3: Almacenar Con un Upsert
Un pipeline programado escribe las mismas filas nuevamente mañana, por lo que la escritura debe definirse en términos de lo que hace que un registro sea único en lugar de agregar ciegamente.
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 es la clave primaria, y el comportamiento INSERT OR REPLACE de SQLite elimina la fila en conflicto antes de insertar la nueva. La tabla almacenada después de la ejecución:
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 las Etapas en un DAG
La declaración del DAG lleva el horario; la última línea declara la cadena de dependencias mediante la entrega de 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 toma una expresión cron estándar, así que 0 6 * * * es 06:00 diario en la zona horaria del DAG. Configurar start_date con una zona explícita importa más de lo que parece: un DAG fijado a una zona que observa el horario de verano cambia su tiempo de ejecución real dos veces al año, y los desplazamientos provienen de la base de datos de zonas horarias de IANA. Fijar a UTC mantiene el intervalo fijo. catchup=False importa específicamente para los scrapers: con él configurado a True, desplegar un DAG cuyo start_date es un mes atrás pone en cola una ejecución por cada intervalo perdido, y un scraper no puede recuperar una página tal como se veía hace tres semanas de todos modos.
Escribir store(extract(fetch())) es todo el gráfico de dependencias. Airflow lee el anidamiento de llamadas y construye los bordes.
Ejecútalo Sin Docker
La mayoría de las guías de Airflow comienzan con un archivo Docker Compose y un contenedor de Postgres. Para desarrollar un DAG, airflow dags test ejecuta todo el gráfico en proceso, sin iniciar el programador, el servidor web o un contenedor.
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
Esa ejecución es el bucle de retroalimentación más rápido disponible mientras la lógica de análisis aún se está moviendo. El programador es lo que inicias una vez que el DAG está haciendo lo que deseas.
¿Programando contra páginas que se renderizan del lado del cliente? El plan gratuito de Scrapeless cubre suficientes solicitudes para ejecutar un DAG diario de principio a fin.
Ejecutar Dos Veces No Debe Duplicar Tus Filas
Un pipeline programado se ejecuta sin supervisión, así que la prueba interesante es la segunda ejecución en lugar de la primera. Ejecutando el mismo DAG de nuevo:
text
stored; table now holds 20 rows
Veinte, no cuarenta. El upsert basado en title reemplazó cada fila y actualizó su scraped_at, lo que hace que el DAG sea seguro para activar manualmente mientras se depura. Un INSERT simple habría duplicado la tabla y no habría dejado forma de saber qué copia era la actual.
Si la historia importa — rastrear cómo se mueve un precio — la clave se convierte en el par de título y fecha de ejecución en lugar de solo el título, y la fila de ayer se mantiene.
El Límite de XCom
Los valores de retorno se mueven entre tareas a través de XCom, y los valores de XCom se serializan en la base de datos de metadatos. Las tres etapas aquí almacenadas:
text
task_id serialized bytes
fetch 53889
extract 2154
store 2
El HTML es 25 veces el tamaño de los registros extraídos de él, y se encuentra en la base de datos de Airflow en lugar de en la memoria. Eso está bien a este tamaño y deja de estar bien a medida que las páginas se vuelven más grandes o las ejecuciones se vuelven más frecuentes.
La forma que escala es analizar temprano y pasar pequeño: hacer que la tarea de recuperación escriba el documento en bruto en almacenamiento de objetos y devuelva una clave, luego permitir que la extracción lo lea de nuevo. Los límites de las etapas se mantienen idénticos y la base de datos de metadatos sigue manteniendo valores del tamaño de registros en lugar de documentos.
El Código de Airflow 2 No Se Ejecutará Aquí Sin Cambios
Dos cambios atrapan a cualquiera que siga un tutorial más antiguo, y solo uno de ellos se anuncia a sí mismo.
python
from airflow.decorators import dag, task # deprecated
from airflow.sdk import dag, task # Airflow 3
La antigua ruta de importación aún se resuelve y emite DeprecatedImportWarning: The airflow.decorators.dag attribute is deprecated. Please use 'airflow.sdk.dag'. El parámetro de programación es un fallo definitivo:
text
TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'
schedule_interval= se convirtió en schedule=. Dado que casi todos los tutoriales de scraping de Airflow publicados son anteriores a la versión 3, ese TypeError es lo primero que muchos lectores encuentran — y aparece en el momento del análisis del DAG, antes de que se ejecute cualquier código de scraping.
Conclusión
El pipeline son tres funciones y un decorador: fetch devuelve HTML, extract devuelve registros, store devuelve un conteo, y store(extract(fetch())) es el gráfico. Lo que lo convierte en un pipeline en lugar de un script es el upsert que sobrevive a una segunda ejecución, el horario expresado como cron, y el límite de etapa que te permite re-analizar sin volver a recuperar.
airflow dags test y SQLite, mantén los documentos fuera de XCom una vez que crezcan, y fija la instalación con el archivo de restricciones para que la primera cosa con la que luches sea el marcado en lugar del resolvedor de dependencias. Para el almacenamiento de la misma problemática, nuestra tubería DuckDB y Parquet cubre la salida columnar, y los precios enumeran cuánto cuesta un horario diario en solicitudes.
¿Listo para programar un raspador? Comienza con el plan gratuito de Scrapeless y apunta la etapa de obtención a tu propio objetivo.
FAQ
P: ¿Necesito Docker para ejecutar Airflow para raspado?
No. airflow dags test <dag_id> ejecuta el gráfico completo en el proceso contra la base de datos SQLite predeterminada, que es cómo se ejecutó el tiempo de 25.2 segundos anterior. Docker y Postgres se vuelven útiles cuando quieres que el programador funcione desatendido con tareas ejecutándose en paralelo, porque SQLite no admite las escrituras concurrentes que se requieren.
P: ¿Dónde debería vivir la clave API en un DAG de Airflow?
No en el archivo DAG. Leerla desde el entorno del proceso la mantiene fuera del control de versiones, y las Variables o Conexiones de Airflow son el mejor hogar una vez que varios DAG necesitan la misma credencial; ambos se almacenan en la base de datos de metadatos y se hacen referencia por nombre en lugar de ser pegados en el código.
P: ¿Cómo detengo un raspador programado de crear filas duplicadas?
Define lo que hace que un registro sea único y escribe contra esa clave. Aquí title es la clave primaria y la escritura es INSERT OR REPLACE, por lo que la segunda ejecución dejó 20 filas en lugar de 40. Si quieres un historial en lugar del estado actual, amplía la clave para incluir la fecha de ejecución para que la instantánea de cada día sea su propia fila.
P: ¿Es seguro pasar HTML raspado entre tareas de Airflow?
En tamaños pequeños, sí; el documento de 50,403 caracteres anterior se serializó a 53,889 bytes en la base de datos de metadatos y se movió entre tareas sin problemas. Deja de ser una buena idea a medida que los documentos o la frecuencia de ejecución crecen, porque cada valor vive en esa base de datos. Escribir el documento en almacenamiento de objetos y pasar una clave mantiene los límites de tarea con valores XCom del tamaño del registro.
P: ¿Por qué mi código de tutorial de Airflow genera un TypeError en schedule_interval?
Porque fue escrito para Airflow 2. La versión 3 renombró el parámetro a schedule, y pasar el nombre antiguo genera TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval' mientras se está analizando el DAG. El cambio complementario es la importación: airflow.decorators aún funciona pero advierte, y airflow.sdk es la ruta actual.
P: ¿Debería habilitarse catchup para un DAG de raspado?
Normalmente no. Con catchup=True, un DAG cuya start_date está en el pasado pone en cola una ejecución por cada intervalo perdido tan pronto como se despliega. Ese comportamiento se adapta a reprocesar un conjunto de datos datado, pero un raspador lee lo que la página muestra ahora, por lo que esas ejecuciones recogerían los datos de hoy y los etiquetarían con fechas históricas.
En Scrapeless, solo accedemos a datos disponibles públicamente y cumplimos estrictamente con las leyes, regulaciones y políticas de privacidad del sitio web aplicables. El contenido de este blog es sólo para fines de demostración y no implica ninguna actividad ilegal o infractora. No ofrecemos garantías y renunciamos a toda responsabilidad por el uso de la información de este blog o enlaces de terceros. Antes de realizar cualquier actividad de scraping, consulte a su asesor legal y revise los términos de servicio del sitio web de destino u obtenga los permisos necesarios.



