Captcha in einer Dagster-Pipeline lösen (Ops und Assets)

dagster captcha - How to Solve CAPTCHA in a Dagster Pipeline (Ops and Assets)

Ein Captcha-Schritt in Dagster ist ein op mit einer Retry-Policy darauf, und die Pipeline drumherum umfasst etwa dreißig Zeilen. Der Teil, der spezifisch für Dagster ist und vor dem Sie kein anderer Orchestrator warnen wird, ist das, was mit dem Wert passiert, den der op zurückgibt. Dagster persistiert die Outputs von ops und Assets über einen IO manager, und der voreingestellte schreibt sie per pickle auf die Festplatte. Ein Captcha-Token ist ein einmalig verwendbares Credential mit zwei Minuten Lebensdauer, und das ist der letzte Ort, an dem es landen sollte.

Was Sie brauchen

  • Dagster 1.9 oder neuer und Python 3.10 oder neuer, dazu das CapSkip Python-SDK.
  • Dagster Open Source oder ein Dagster+-Deployment. Der Code ist in beiden Fällen derselbe.
  • Die Seiten-URL des geschützten Formulars und dessen sitekey.
  • CapSkip im Local mode, wenn Code und Löser auf derselben Maschine liegen, sonst im Server mode. Beides finden Sie unter Verbindungseinstellungen.
# pip install capskip
pip install dagster dagster-webserver capskip

# Run the UI locally while you build the job.
dagster dev

Wo Ihr Code tatsächlich läuft

Klären Sie das, bevor Sie einen Host-String wählen, denn davon hängt das gesamte Netzwerk-Setup ab. Dagster Open Source läuft vollständig auf Ihren eigenen Maschinen, hier gibt es also nichts zu besprechen. Dagster+ gibt es in zwei Ausprägungen, und für diesen Zweck sind sie das genaue Gegenteil voneinander.

Hybrid ist die Variante, die sich so verhält, wie man es sich wünscht. Sie betreiben einen Agent in Ihrer eigenen Infrastruktur, und er verbindet sich nach außen mit der Control Plane. Dagster+ hat keinen Zugang in Ihr Netzwerk hinein, sieht Ihren Code nicht und fasst Ihre Daten nicht an, ein aus der Cloud-UI gestarteter Lauf wird also trotzdem auf Hardware ausgeführt, die Ihnen gehört. Setzen Sie den Agent auf dieselbe Maschine wie CapSkip, und der op ruft 127.0.0.1:8080 genauso auf wie ein Skript auf Ihrem eigenen Rechner. Kein Tunnel und keine öffentliche Adresse.

Serverless ist die Ausnahme. Dort wird Ihr Code in der Umgebung von Dagster ausgeführt und nicht in Ihrer, und Loopback bedeutet nichts Brauchbares mehr: Es zeigt auf den Container von Dagster, wo nichts lauscht. In dieser Konfiguration braucht der Löser den Server mode und eine Adresse, die der Lauf erreichen kann. Das ist keine Verschlechterung: Es ist derselbe Löser auf derselben Hardware mit einer anderen Bind-Adresse, und pro Lösung fallen weiterhin keine Kosten an.

Schritt 1: Den Löser in eine Resource verpacken

Die Dagster-Idiomatik für alles Externe ist eine Resource, kein Client auf Modulebene. Leiten Sie von ConfigurableResource ab und deklarieren Sie die Verbindungsfelder: Die UI bekommt dafür einen Launchpad-Eintrag, die Konfiguration wird vor dem Start des Laufs validiert, und Tests können das Ganze gegen ein Fake austauschen. Lesen Sie den Key aus der Umgebung, statt ihn fest zu verdrahten.

# pip install capskip
import dagster as dg
from capskip import CapSkip

class CapSkipResource(dg.ConfigurableResource):
    host: str = "127.0.0.1"
    port: int = 8080
    api_key: str = "capskip"

    def solve_recaptcha(self, sitekey: str, page_url: str) -> str:
        solver = CapSkip(host=self.host, port=self.port, apiKey=self.api_key)
        return solver.recaptcha(sitekey=sitekey, url=page_url)["code"]

Eine Methode deckt reCAPTCHA v2, Invisible, Enterprise und v3 ab, denn die Varianten sind Keyword-Optionen statt eigener Aufrufe: invisible auf 1, enterprise auf 1 oder version auf v3 mit einer Action. Turnstile und GeeTest haben eigene Methoden mit demselben Aufbau. Die vollständige Parameterliste finden Sie in der CapSkip-API-Dokumentation.

Schritt 2: Ein op, mit einer Retry-Policy

Ein Löser-Aufruf ist ein Netzwerkaufruf, und Netzwerkaufrufe schlagen fehl. Die RetryPolicy von Dagster ist deklarativ und kommt an den Decorator: eine Maximalzahl, eine Basisverzögerung, eine Backoff-Kurve und ein Jitter-Modus. Exponentielles Backoff mit Plus-Minus-Jitter ist der sinnvolle Standard, weil es eine Reihe gleichzeitiger Fehlschläge auseinanderzieht, statt sie alle auf einmal zurückzubringen.

# The policy lives on the decorator, not in the body.
@dg.op(
    retry_policy=dg.RetryPolicy(
        max_retries=3,
        delay=5,
        backoff=dg.Backoff.EXPONENTIAL,
        jitter=dg.Jitter.PLUS_MINUS,
    )
)
def solve_and_submit(context, sitekey: str, capskip: CapSkipResource) -> int:
    token = capskip.solve_recaptcha(sitekey, PAGE_URL)
    context.log.info("Solved, submitting immediately.")
    return post_form(PAGE_URL, token)

Beachten Sie, was dieser op zurückgibt: einen Statuscode, kein Token. Genau darum geht es im nächsten Abschnitt.

Schritt 3: Das Token nie eine op-Grenze überschreiten lassen

Das ist die Dagster-spezifische Falle, und man tappt leicht hinein, denn das naheliegende Design besteht aus einem op, der löst, und einem zweiten, der absendet. Dagster reicht Werte zwischen ops nicht im Arbeitsspeicher weiter. Jeder Output läuft durch einen IO manager, und voreingestellt ist der für das Dateisystem, der Outputs als Pickle-Dateien auf der lokalen Platte ablegt. Ein Token, das ein op zurückgibt, wird also in eine Datei geschrieben, vom nächsten op wieder eingelesen und bleibt danach dort liegen.

Das ist gleich doppelt falsch. Es schreibt ein gültiges Credential auf die Platte, wo es niemand aufräumt. Und es legt einen Umweg über den Speicher zwischen das Lösen und das Absenden, also genau die Verzögerung, die sich ein Ablauf nach zwei Minuten nicht leisten kann. Bei einer erneuten Ausführung wird es schlimmer: Dagster kann den gespeicherten Output eines früheren Laufs laden, statt neu zu rechnen, und dann bekommt das Formular ein Token, das gestern abgelaufen ist.

Die Lösung ist kein ausgeklügelter IO manager. Sie besteht darin, Lösen und Absenden im selben op zu halten, damit das Token in einer lokalen Variablen lebt und nie zu einem Output wird. Alles davor und danach kann als eigene ops oder Assets bestehen bleiben. Nur das Paar, das sich das Token teilen muss, wird zusammengelegt.

# The token is a local variable. It is never an op output,
# so no IO manager ever sees it and nothing is written to disk.
@dg.op(retry_policy=dg.RetryPolicy(max_retries=3, delay=5))
def submit_protected_form(sitekey: str, capskip: CapSkipResource) -> int:
    token = capskip.solve_recaptcha(sitekey, PAGE_URL)
    return post_form(PAGE_URL, token)

Wenn Sie mit Assets statt mit ops arbeiten, gilt dieselbe Regel mit einer Ergänzung: Ein materialisiertes Asset ist ein dauerhafter Eintrag mit einer Historie in der UI, und ein einmalig verwendbares Geheimnis hat dort nichts zu suchen. Modellieren Sie das Lösen als op innerhalb eines graph-backed Asset und lassen Sie das Asset das Ergebnis der Übermittlung materialisieren.

Vollständiges lauffähiges Beispiel

Ein vollständiger Job. Die Seite holen, den Sitekey daraus ziehen, dann in einem op lösen und absenden. Zwei ops, eine Resource, ein Definitions-Objekt.

# pip install dagster capskip httpx
import re
import httpx
import dagster as dg
from capskip import CapSkip

PAGE_URL = "https://example.com/page-with-recaptcha"

class CapSkipResource(dg.ConfigurableResource):
    host: str = "127.0.0.1"
    port: int = 8080
    api_key: str = "capskip"

    def solve_recaptcha(self, sitekey: str, page_url: str) -> str:
        solver = CapSkip(host=self.host, port=self.port, apiKey=self.api_key)
        return solver.recaptcha(sitekey=sitekey, url=page_url)["code"]

@dg.op(retry_policy=dg.RetryPolicy(max_retries=2, delay=3))
def read_sitekey() -> str:
    html = httpx.get(PAGE_URL, timeout=30).text
    match = re.search(r'data-sitekey=["\']([^"\']+)', html)
    if not match:
        raise dg.Failure("No data-sitekey on the page.")
    return match.group(1)

@dg.op(
    pool="capskip",
    retry_policy=dg.RetryPolicy(
        max_retries=3, delay=5, backoff=dg.Backoff.EXPONENTIAL
    ),
)
def submit_protected_form(sitekey: str, capskip: CapSkipResource) -> int:
    token = capskip.solve_recaptcha(sitekey, PAGE_URL)
    reply = httpx.post(
        PAGE_URL,
        data={"g-recaptcha-response": token},
        timeout=30,
    )
    return reply.status_code

@dg.job
def captcha_protected_submit():
    submit_protected_form(read_sitekey())

defs = dg.Definitions(
    jobs=[captcha_protected_submit],
    resources={
        "capskip": CapSkipResource(api_key=dg.EnvVar("CAPSKIP_API_KEY")),
    },
)

EnvVar verschiebt die Auflösung auf die Laufzeit, statt den Wert in die Definition einzubacken, der Key steht also nie in Ihrem Repository und nie in einem serialisierten Snapshot.

Fehler, die nie durchgehen werden, nicht wiederholen

Eine deklarative Policy wiederholt alles, was drei Versuche und eine Minute Backoff an ein fehlerhaftes Argument verschwendet, das jedes Mal identisch scheitert. Die Notluke von Dagster besteht darin, den Retry selbst auszulösen: Fangen Sie die Ausnahmen ab, bei denen sich ein weiterer Versuch lohnt, fordern Sie einen Retry ausdrücklich an, und lassen Sie den Rest durchschlagen und den Lauf sofort scheitern.

# Retry the transient ones. Fail fast on the rest.
from capskip import NetworkException, TimeoutException

@dg.op
def solve_with_judgement(sitekey: str, capskip: CapSkipResource) -> int:
    try:
        token = capskip.solve_recaptcha(sitekey, PAGE_URL)
    except (NetworkException, TimeoutException) as err:
        raise dg.RetryRequested(max_retries=3, seconds_to_wait=10) from err
    return post_form(PAGE_URL, token)

NetworkException bedeutet, dass CapSkip nicht läuft oder der Host falsch ist, und TimeoutException bedeutet, dass das Lösen länger gedauert hat als recaptchaTimeout, das standardmäßig bei 300 Sekunden liegt. Beide verdienen einen weiteren Versuch. ValidationException bedeutet ein fehlerhaftes Argument und ApiException, dass die API die Anfrage abgelehnt hat, und keines von beiden wird beim zweiten Mal besser. Alle vier leiten sich von einer gemeinsamen Basis namens CapSkipError ab, Sie können also diese eine abfangen, wenn Sie lieber alles an einer Stelle behandeln.

Einen Batch mit einem Concurrency-Pool drosseln

Fächern Sie einen partitionierten Job über zweihundert URLs auf, und Dagster wird bereitwillig versuchen, sie alle auszuführen, was mehr Last ist, als Sie auf einen einzelnen Löser richten wollen. Concurrency-Pools sind die Steuerung dafür. Versehen Sie den op mit einem Pool-Namen, setzen Sie das Limit einmal, und alles über der Grenze wandert in die Warteschlange, statt sich aufzutürmen.

# Six solves in flight at a time, across every run.
dagster instance concurrency set capskip 6

Das Argument pool steht im Beispiel oben bereits am op. Limits gelten über alle Läufe hinweg und nicht pro Lauf, und genau das wollen Sie, wenn drei Schedules auf dieselbe Maschine zeigen. Wenn Sie das Auffächern lieber innerhalb eines Prozesses erledigen statt mit einem op pro URL: Das Python-SDK bringt einen echten asyncio-Client mit, und diesen Ansatz finden Sie in der Anleitung zum parallelen Lösen von Captchas.

Den Löser auf einer anderen Maschine betreiben

Container für Code Locations, Kubernetes-Agents und Serverless-Läufe verlagern Ihren op allesamt weg von Ihrem eigenen Rechner. Am Code ändert sich nichts außer dem Host, und die Resource liest ihn ohnehin aus der Konfiguration.

CapSkip hat zwei Verbindungsmodi. Local bindet an 127.0.0.1 und antwortet nur diesem einen Gerät. Server bindet an Ihre Netzwerk-IP oder öffentliche IP, sodass ein Container, eine VM oder ein Serverless-Lauf dieselbe Windows-Maschine über die API erreicht. Eine statische öffentliche IP hält die Adresse stabil. Es bleibt Ihre Hardware und wird so oder so nicht nach Verbrauch abgerechnet, die Kosten eines vollen Tages ändern sich mit dem Modus also nicht.

# Same resource, same call. Only the host moves.
resources={
    "capskip": CapSkipResource(
        host=dg.EnvVar("CAPSKIP_HOST"),
        api_key=dg.EnvVar("CAPSKIP_API_KEY"),
    ),
}

Aktivieren Sie die Key-Validierung, sobald der Löser auf einer Netzwerkadresse lauscht, und geben Sie jeder Code Location einen eigenen Key, damit sich einer widerrufen lässt, ohne die anderen anzufassen. Beide Modi werden Schritt für Schritt erklärt unter CapSkip-Einrichtungsanleitung.

Häufige Fehler und was sie bedeuten

Was Sie sehenUrsacheBeheben
Das Formular lehnt ein Token ab, das korrekt aussiehtEs lief zwischen zwei ops durch einen IO manager und kam veraltet anIm selben op lösen und absenden, damit das Token eine lokale Variable bleibt
Eine erneute Ausführung schickt ein Token aus einem alten Lauf abDagster hat den zuvor gespeicherten Output geladen, statt neu zu rechnenGleiche Lösung. Ein Token darf nie Output eines op oder eines Asset sein
Unter DAGSTER_HOME/storage taucht eine Pickle-Datei aufDer voreingestellte IO manager für das Dateisystem hat Ihr Token auf die Platte geschriebenGleiche Lösung, danach die Datei löschen. Behandeln Sie sie als geleaktes Credential
NetworkException auf Dagster+ ServerlessDer Lauf wurde in der Umgebung von Dagster ausgeführt, nicht in IhrerDen Löser auf Server mode umstellen oder einen Hybrid-Agent verwenden
NetworkException auf Ihrem eigenen AgentCapSkip läuft nicht, oder der Host ist falschDie App starten oder den Host der Resource auf die Server-Adresse richten
Drei Wiederholungen für einen falschen sitekey verbranntJeder Versuch scheitert auf dieselbe deterministische WeiseRetryRequested nur bei NetworkException und TimeoutException auslösen
Der Löser ist während eines Backfills überlastetPartitionen laufen standardmäßig parallelDen op in einen Pool legen und ein Limit darauf setzen
ValidationExceptionFehlender oder fehlerhafter Sitekey oder fehlerhafte Seiten-URLBeide vor dem Aufruf loggen und prüfen, ob der Sitekey der aktive ist

FAQ

Sollte das Lösen ein Asset sein?

Nein. Ein Asset ist ein dauerhaftes Objekt mit einer Materialisierungshistorie, und ein Token, das nach zwei Minuten abläuft, ist das genaue Gegenteil davon. Modellieren Sie das, was Sie tatsächlich erzeugt haben, als Asset, sei es der übermittelte Datensatz oder die Seite, die Sie endlich lesen durften, und lassen Sie das Lösen ein op darin bleiben.

Kann ein Dagster+-Lauf wirklich 127.0.0.1 erreichen?

Bei Hybrid ja, denn der Agent in Ihrer eigenen Infrastruktur führt den Code aus. Loopback bedeutet dort die Maschine des Agents, ein Löser auf dieser Maschine antwortet also ganz normal. Bei Serverless läuft der Lauf in der Umgebung von Dagster, und Sie brauchen den Server mode mit einer erreichbaren Adresse.

Worin unterscheidet sich das davon, es in Airflow zu tun?

Vor allem darin, wie sich Werte bewegen. Airflow schiebt kleine Werte durch XCom, und Sie legen dort einfach kein Token hinein. Dagster persistiert standardmäßig jeden Output über einen IO manager, derselbe Fehler schreibt also eine Pickle-Datei statt einer Datenbankzeile. Die Seite des Lösers ist identisch, und die Airflow-Variante ist beschrieben in der Airflow-Captcha-Anleitung.

Blockiert ein langes Lösen den Rest des Laufs?

Der op ist blockiert, während er pollt, aber der Multiprocess-Executor von Dagster führt weiter andere ops aus, die nicht von ihm abhängen. Ein Pool-Limit begrenzt, wie viele Lösungen gleichzeitig stattfinden, ohne sonst etwas zu stoppen. Auf Ihrer eigenen Hardware kostet ein langsames Lösen nur Wartezeit.

Die Kurzfassung

Verpacken Sie den Löser in eine ConfigurableResource und lesen Sie den Key aus EnvVar. Geben Sie dem op eine RetryPolicy mit exponentiellem Backoff und lösen Sie RetryRequested selbst aus, wenn Sie Wiederholungen überspringen wollen, die nicht gelingen können. Vor allem aber halten Sie Lösen und Absenden im selben op, denn Dagster persistiert die Outputs von ops standardmäßig, und ein Token ist nichts, was man persistieren sollte.

Die Python-Seite von alldem finden Sie auf der Python-Captcha-Solver-Seite. Die Details zu diesem Captcha-Typ finden Sie auf die reCAPTCHA-v2-Solver-Seite. Dieselben drei Aufrufe gibt es auch in Node.js, PHP und C#. Aufgeführt sind sie hier: Die Seite zu den SDKs fürs Captcha-Lösen.

Eines sollten Sie wissen, bevor Sie das für einen großen Backfill einplanen. CapSkip ist ein Captcha-Löser der auf Hardware läuft, die Ihnen bereits gehört, sodass ein Job, der fünfzigtausend Zeilen löst, genau so viel kostet wie ein Job, der fünfzig löst.