So lösen Sie Captchas in Apache-Airflow-DAGs (Python)

airflow captcha - How to Solve CAPTCHA in Apache Airflow DAGs (Python)

Ein Captcha-Schritt in Airflow ist ein gewöhnlicher Python-Task. Sie rufen den Löser auf, Sie bekommen einen Token zurück, Sie verwenden ihn im selben Task. Es gibt keinen Operator zu installieren und kein Plugin zu schreiben. Woran es tatsächlich scheitert, ist der Ort: Ihre Worker laufen dort, wohin der Scheduler sie gesetzt hat, und ein Löser, der an die Loopback-Adresse Ihres Laptops gebunden ist, ist aus einem Container auf einem anderen Host nicht erreichbar. Bringen Sie diesen Teil zuerst in Ordnung, dann ist der DAG fünfzehn Zeilen lang.

Was Sie brauchen

  • Airflow 2.x oder 3.x, mit dem CapSkip-SDK, installiert im selben Image oder virtualenv, das Ihre Worker verwenden.
  • Ein Ziel, das eine Challenge zurückgibt. Eine Seite hinter einem Widget, kein Testschlüssel, der immer durchgeht.
  • CapSkip läuft im Server-Modus auf einer Maschine, die Ihre Worker erreichen können, oder im Local-Modus, wenn Worker und Löser derselbe Rechner sind. Beides finden Sie unter Verbindungseinstellungen.
# Install into the worker environment, not just the scheduler.
pip install capskip

Wo der Löser stehen muss

Das ist das ganze Problem, deshalb steht es zuerst. Ein Task läuft in einem Worker-Prozess, und in den meisten echten Deployments ist dieser Worker ein Container auf einem anderen Host als alles, was Sie von Hand verwalten. Die Loopback-Adresse in diesem Container ist der Container selbst; wenn Sie das SDK darauf ausrichten, findet es nichts.

ModusLauscht aufSinnvoll, wenn
Lokal127.0.0.1, nur dieses GerätEin einzelner Worker auf derselben Maschine wie der Löser
ServerIhre Netzwerkadresse oder öffentliche IPContainer, eine Worker-Flotte, ein VPS oder ein verwaltetes Airflow

Der Server-Modus ist die Antwort für fast jedes Airflow-Deployment. Sie stellen die Lauschadresse in der App um, richten den SDK-Host auf diese Maschine aus, und jeder Worker der Flotte teilt sich einen Löser. Eine statische öffentliche IP wird empfohlen, wenn die Aufrufer außerhalb Ihres eigenen Netzwerks sitzen. Am Produkt selbst ändert das nichts: Es bleibt Ihre Hardware und bleibt ohne Verbrauchsabrechnung, der Weg weg von der Loopback-Adresse verschiebt also nur, wo es läuft, und sonst nichts. CapSkip ist eine Windows-Anwendung, in der Praxis ist das also ein Windows-Rechner, den die Flotte aufruft.

Schritt 1: den Host einmal konfigurieren

Schreiben Sie die Adresse nicht fest in die DAG-Datei. Das SDK liest CAPSKIP_HOST, CAPSKIP_PORT und CAPSKIP_API_KEY aus der Umgebung, was am saubersten zu Airflow passt, weil Sie ohnehin schon einen Weg haben, Umgebungsvariablen für Worker zu setzen. Eine Airflow Variable funktioniert ebenfalls, wenn Sie den Wert lieber in der Metadaten-Datenbank halten.

# pip install capskip
import os
from capskip import CapSkip


def get_solver():
    """One place that knows where the solver lives."""
    return CapSkip(
        host=os.environ.get("CAPSKIP_HOST", "127.0.0.1"),
        port=int(os.environ.get("CAPSKIP_PORT", "8080")),
        recaptchaTimeout=300,
    )

Bauen Sie den Client innerhalb des Tasks auf, nicht auf Modulebene. Airflow parst jede DAG-Datei in einer Schleife, und alles, was zur Importzeit erzeugt wird, entsteht bei jedem Parsen neu, im Scheduler ebenso wie im Worker.

Schritt 2: das Lösen als einen Task schreiben

Die TaskFlow-Dekoratoren sind der kürzeste Einstieg. Airflow 3 importiert sie aus dem SDK-Modul, Airflow 2 aus dem decorators-Modul, und der Rumpf ist in beiden Fällen identisch.

# Airflow 3.x. On 2.x use: from airflow.decorators import dag, task
from airflow.sdk import dag, task

PAGE = "https://example.com/page-with-recaptcha"
SITEKEY = "YOUR_SITEKEY"


@task(retries=2)
def fetch_protected_page():
    solver = get_solver()

    # Solve and use in the same task. The token is short lived.
    token = solver.recaptcha(sitekey=SITEKEY, url=PAGE)["code"]

    return post_form(PAGE, token)

Das ist die gesamte Integration. Der Löser-Aufruf ist synchron, er pollt für Sie mit einem Backoff, der bei einer Viertelsekunde beginnt, und er gibt ein Dictionary zurück, dessen Feld code den Token enthält.

Schritt 3: den Token nicht zwischen Tasks weiterreichen

Das ist der Fehler, den man benennen muss, weil Airflow ihn natürlich wirken lässt. Rückgaben von TaskFlow werden zu XComs, ein Lösungs-Task, der einen Token zurückgibt, und ein nachgelagerter Task, der ihn konsumiert, sehen also nach gutem Design aus. Es ist ein Wettlauf. Ein reCAPTCHA-Token bleibt rund zwei Minuten gültig, und der Abstand zwischen zwei Airflow-Tasks ist eine Entscheidung des Schedulers, die Sie nicht steuern. Kommen eine Warteschlangenverzögerung, ein voller Pool oder ein Worker-Neustart hinzu, läuft der Token unterwegs ab.

Der Fehler tritt sporadisch auf, und das ist die schlimmste Sorte: Der DAG funktioniert im Test und scheitert in der Produktion in ein paar Prozent der Fälle, mit einem abgewiesenen Formular und ohne Fehler vom Löser. Der Beitrag zur Gültigkeitsdauer eines reCAPTCHA-Tokens behandelt die Zeitfenster. In einem DAG schrumpft die Regel auf eine Zeile: lösen und absenden im selben Task, und den Retry neu lösen lassen.

Der vollständige DAG

# pip install capskip
import os
from datetime import datetime, timedelta

import requests
from airflow.sdk import dag, task
from capskip import CapSkip

PAGE = "https://example.com/page-with-recaptcha"
SITEKEY = "YOUR_SITEKEY"


@dag(
    schedule="@hourly",
    start_date=datetime(2026, 1, 1),
    catchup=False,
    tags=["scraping"],
)
def protected_source():

    @task(retries=2, retry_delay=timedelta(minutes=2), pool="captcha")
    def scrape():
        solver = CapSkip(
            host=os.environ.get("CAPSKIP_HOST", "127.0.0.1"),
            port=int(os.environ.get("CAPSKIP_PORT", "8080")),
        )
        token = solver.recaptcha(sitekey=SITEKEY, url=PAGE)["code"]

        # Same task, so the token is seconds old when it is used.
        r = requests.post(
            PAGE,
            data={"g-recaptcha-response": token},
            timeout=60,
        )
        r.raise_for_status()
        return len(r.text)

    scrape()


protected_source()

Das Argument pool leistet dort echte Arbeit. Ein Pool begrenzt, wie viele Task-Instanzen im gesamten Deployment gleichzeitig laufen, sodass ein Backfill von zweihundert Läufen nicht zweihundert gleichzeitige Lösungen gegen eine Maschine öffnet. Legen Sie in der UI einen Pool namens captcha an, geben Sie ihm die Anzahl paralleler Lösungen, die Sie wollen, und jeder Task, der ihn nennt, reiht sich hinter diesem Limit ein.

Retries, die helfen statt zu schaden

Setzen Sie retries am Task und lassen Sie Lösen und Absenden komplett wiederholen. Weil der Token im Rumpf des Tasks geholt wird, bekommt ein Retry automatisch einen frischen, und genau das ist das gewünschte Verhalten und der Grund, warum die beiden Schritte zusammenbleiben.

Geben Sie dem Retry eine Verzögerung. Ein sofortiger Retry gegen eine Website, die Ihnen gerade eine Challenge gestellt hat, bekommt meist gleich die nächste, und ein paar Minuten kosten in einer geplanten Pipeline nichts. Zwei Retries reichen in der Regel: Ein dritter Fehlschlag liegt normalerweise an einem falschen Sitekey oder an einem Löser, der nicht erreichbar ist, und beides behebt Warten nicht.

Häufige Fehler

Was Sie sehenUrsacheBeheben
NetworkException bei jedem TaskDer Worker erreicht den Löser nichtWechseln Sie in den Server-Modus und setzen Sie CAPSKIP_HOST auf dem Worker
Funktioniert im Scheduler, scheitert im WorkerDas SDK fehlt im Worker-ImageInstallieren Sie es dort, wo Tasks laufen, nicht nur dort, wo DAGs geparst werden
TimeoutException unter LastMehr gleichzeitige Lösungen, als die Maschine schafftStecken Sie den Task in einen Pool und begrenzen Sie die Slot-Anzahl
Das Formular weist einen Token ab, der sauber gelöst wurdeDer Token ist zwischen zwei Tasks gealtertVerlagern Sie das Absenden in den lösenden Task
ERROR_PAGEURLEine relative URL hat die API erreichtSenden Sie die absolute URL, inklusive Schema

Die vollständige Liste der Codes und ihrer Auslöser finden Sie unter der CapSkip-API-Dokumentation.

FAQ

Meine Worker laufen in Containern. Können sie den Löser erreichen?

Ja, im Server-Modus. Der Löser lauscht auf einer Netzwerkadresse statt auf der Loopback-Adresse, und die Worker rufen ihn über die API auf wie jeden anderen internen Dienst. Setzen Sie Host und Port als Umgebungsvariablen der Worker, damit in der DAG-Datei keine Adresse steht. Am Code ändert sich zwischen einem Container und einem Laptop nichts.

Was ist mit verwaltetem Airflow, das ich nicht selbst administriere?

Dieselbe Antwort, mit einer zusätzlichen Anforderung. Ein gehosteter Scheduler läuft außerhalb Ihres Netzwerks, der Löser braucht also eine Adresse, zu der er routen kann, und eine statische öffentliche IP macht das stabil. Schalten Sie die Schlüsselvalidierung ein und geben Sie jeder Umgebung ihren eigenen Schlüssel, damit ein durchgesickerter Wert widerrufen werden kann, ohne die anderen anzufassen.

Sollte das Lösen für die Observability ein eigener Task sein?

Es ist verlockend und kostet Sie Zuverlässigkeit. Ein eigener Task bedeutet, dass der Token als XCom reist und altert, während der Scheduler entscheidet, was als Nächstes läuft, und genau so laufen Token unterwegs ab. Halten Sie beides zusammen und holen Sie sich die Observability stattdessen aus Logs und Task-Dauer, die dasselbe zeigen, nur ohne den Wettlauf.

Können sich mehrere DAGs eine Löser-Instanz teilen?

Ja, und das ist das übliche Setup. Eine Instanz im Server-Modus bedient jeden Worker, der sie erreicht, und es gibt keine Kosten pro Lösung, die aufzuteilen wären. Nutzen Sie einen Airflow-Pool, um die Gesamtparallelität über alle DAGs hinweg zu begrenzen, denn das Limit, auf das es Ihnen ankommt, ist, wie viele Lösungen gleichzeitig laufen, nicht, wie viele Pipelines gerade eine wollen.

Die Kurzfassung

Stellen Sie den Löser auf eine Maschine, die die Worker erreichen, setzen Sie den Host in der Umgebung und halten Sie Lösen und Absenden in einem Task, mit ein paar Retries hinter einem Pool. Alles andere ist ein normaler DAG. Für die Crawl-Seite davon siehe die Seite zum Captcha-Solver für Web Scraping, und für die Client-Schnittstelle siehe der Python-Captcha-Solver-Seite. Bei geplanter Arbeit zeigt sich das Preismodell. Ein stündlicher DAG, der bei jedem Lauf löst, ist bei einem verbrauchsabhängigen Anbieter eine echte Rechnung, und genau das ist hier der Unterschied: Captcha-Löser läuft auf Hardware, die Ihnen bereits gehört, und kostet gleich viel, ob der DAG einmal am Tag oder einmal pro Minute ausgelöst wird.