Cara Memecahkan CAPTCHA dalam Pipeline Dagster (Op dan Asset)

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

Langkah captcha di Dagster hanyalah satu op dengan retry policy padanya, dan pipeline di sekelilingnya sekitar tiga puluh baris. Bagian yang khas Dagster, dan yang tidak akan diperingatkan oleh orchestrator mana pun, adalah apa yang terjadi pada nilai yang dikembalikan op tersebut. Dagster menyimpan output op dan asset melalui sebuah IO manager, dan yang default mem-pickle-nya ke disk. Token CAPTCHA adalah kredensial sekali pakai dengan umur dua menit, jadi disk adalah tempat terakhir yang seharusnya menampungnya.

Apa yang Anda butuhkan

  • Dagster 1.9 atau lebih baru dan Python 3.10 atau lebih baru, ditambah SDK Python CapSkip.
  • Dagster open source, atau deployment Dagster+. Kodenya sama untuk keduanya.
  • URL halaman formulir yang terproteksi, beserta sitekey-nya.
  • CapSkip berjalan dalam Local mode ketika kode dan pemecah berbagi satu mesin, atau dalam Server mode ketika tidak. Keduanya dijelaskan di bagian pengaturan koneksi.
# pip install capskip
pip install dagster dagster-webserver capskip

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

Di mana kode Anda sebenarnya berjalan

Jawab ini sebelum Anda memilih string host, karena jawabannya menentukan seluruh penyiapan jaringan. Dagster open source berjalan sepenuhnya di mesin Anda sendiri, jadi tidak ada yang perlu dibahas di sini. Dagster+ punya dua bentuk dan keduanya berlawanan untuk keperluan ini.

Hybrid adalah bentuk yang berperilaku seperti yang Anda harapkan. Anda menjalankan sebuah agent di dalam infrastruktur Anda sendiri dan agent itu menghubungi control plane keluar. Dagster+ tidak punya jalur masuk ke jaringan Anda, tidak melihat kode Anda dan tidak menyentuh data Anda, sehingga run yang diluncurkan dari UI cloud tetap dieksekusi di perangkat keras milik Anda. Tempatkan agent di mesin yang sama dengan CapSkip dan op akan memanggil 127.0.0.1:8080 persis seperti skrip di komputer Anda. Tanpa tunnel dan tanpa alamat publik.

Serverless adalah pengecualiannya. Di sana kode Anda dieksekusi di lingkungan Dagster, bukan lingkungan Anda, dan loopback berhenti berarti apa pun yang berguna: loopback mengarah ke container Dagster, tempat tidak ada yang mendengarkan. Pada konfigurasi itu pemecah membutuhkan Server mode dan alamat yang dapat dijangkau oleh run. Itu bukan penurunan: pemecah yang sama di perangkat keras yang sama dengan alamat bind yang berbeda, dan tetap tidak menagih apa pun per pemecahan.

Langkah 1: bungkus pemecah dalam sebuah resource

Idiom Dagster untuk apa pun yang bersifat eksternal adalah resource, bukan client di level modul. Buat subclass dari ConfigurableResource, deklarasikan field koneksinya, dan UI akan mendapat entri launchpad untuk field itu, config divalidasi sebelum run dimulai, serta pengujian dapat menukar keseluruhannya dengan tiruan. Baca key dari environment alih-alih menuliskannya langsung di dalam kode.

# 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"]

Satu metode mencakup reCAPTCHA v2, Invisible, Enterprise dan v3, karena variannya berupa opsi keyword alih-alih pemanggilan terpisah: invisible diatur ke 1, enterprise diatur ke 1, atau version diatur ke v3 dengan sebuah action. Turnstile dan GeeTest punya metodenya sendiri dan bentuk yang sama. Daftar parameter lengkapnya ada di dokumentasi API CapSkip.

Langkah 2: satu op, dengan retry policy

Pemanggilan pemecah adalah pemanggilan jaringan, dan pemanggilan jaringan bisa gagal. RetryPolicy milik Dagster bersifat deklaratif dan dipasang pada decorator: jumlah maksimum, jeda dasar, kurva backoff dan mode jitter. Backoff eksponensial dengan jitter plus atau minus adalah default yang masuk akal, karena menyebarkan sekumpulan kegagalan yang terjadi bersamaan alih-alih memulangkan semuanya sekaligus.

# 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)

Perhatikan apa yang dikembalikan op itu: sebuah kode status, bukan token. Itulah inti dari bagian berikutnya.

Langkah 3: jangan pernah biarkan token melewati batas op

Inilah jebakan khas Dagster dan mudah sekali terperosok ke dalamnya, karena desain yang paling jelas adalah satu op yang memecahkan dan op kedua yang mengirim. Dagster tidak meneruskan nilai antar op di memori. Setiap output melewati sebuah IO manager, dan yang default adalah IO manager filesystem, yang menyimpan output sebagai file pickle di disk lokal. Karena itu token yang dikembalikan dari sebuah op ditulis ke sebuah file, dibaca kembali oleh op berikutnya, lalu tertinggal di sana setelahnya.

Itu salah dalam dua hal. Pertama, ia menulis kredensial aktif ke disk, tempat tidak ada yang membersihkannya. Kedua, ia menyisipkan satu perjalanan bolak-balik ke penyimpanan di antara pemecahan dan pengiriman, dan itu persis penundaan yang tidak sanggup ditanggung oleh masa berlaku dua menit. Pada eksekusi ulang keadaannya makin buruk: Dagster bisa memuat output tersimpan dari run sebelumnya alih-alih menghitung ulang, dan formulir pun menerima token yang sudah kedaluwarsa kemarin.

Perbaikannya bukan IO manager yang pintar. Perbaikannya adalah menjaga pemecahan dan pengiriman tetap berada di dalam op yang sama, sehingga token hidup di dalam sebuah variabel lokal dan tidak pernah menjadi output sama sekali. Semua yang ada di hulu dan di hilir boleh tetap berupa op atau asset terpisah. Hanya pasangan yang harus berbagi token yang digabungkan.

# 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)

Jika Anda bekerja dengan asset alih-alih op, aturan yang sama berlaku dengan satu catatan tambahan: asset yang dimaterialisasi adalah catatan permanen dengan riwayat di UI, dan rahasia sekali pakai tidak pantas menjadi catatan semacam itu. Modelkan pemecahan sebagai sebuah op di dalam asset yang didukung graph, dan biarkan asset menjadi materialisasi dari hasil pengirimannya.

Contoh lengkap yang berfungsi

Sebuah job lengkap. Ambil halamannya, tarik sitekey darinya, lalu pecahkan dan kirim dalam satu op. Dua op, satu resource, satu objek Definitions.

# 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 menunda pencarian nilainya hingga run time alih-alih menanamkan nilai itu ke dalam definisi, sehingga key tidak pernah berada di repositori Anda dan tidak pernah masuk ke snapshot terserialisasi.

Jangan mencoba ulang kegagalan yang tidak akan pernah berhasil

Kebijakan deklaratif mencoba ulang semuanya, sehingga membuang tiga percobaan dan satu menit backoff pada argumen salah bentuk yang akan gagal dengan cara yang sama setiap kali. Jalan keluar dari Dagster adalah memunculkan percobaan ulang itu sendiri: tangkap exception yang layak dicoba lagi, minta percobaan ulang secara eksplisit, dan biarkan sisanya merambat naik lalu menggagalkan run seketika.

# 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 berarti CapSkip tidak berjalan atau host-nya salah, dan TimeoutException berarti pemecahan melampaui recaptchaTimeout, yang secara default 300 detik. Keduanya layak dicoba lagi. ValidationException berarti argumen yang salah dan ApiException berarti API menolak permintaannya, dan keduanya tidak akan membaik pada percobaan kedua. Keempatnya diturunkan dari basis bersama bernama CapSkipError, jadi Anda bisa menangkap yang satu itu saja bila lebih suka menangani semuanya di satu tempat.

Membatasi laju sebuah batch dengan concurrency pool

Sebarkan sebuah job berpartisi ke dua ratus URL dan Dagster akan dengan senang hati mencoba menjalankan semuanya, dan itu beban yang lebih besar daripada yang Anda inginkan mengarah ke satu pemecah. Concurrency pool adalah pengendalinya. Beri tag nama pool pada op, tetapkan batasnya sekali, dan apa pun yang melebihi batas akan mengantre alih-alih menumpuk.

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

Argumen pool sudah ada pada op di contoh di atas. Batas berlaku di seluruh run, bukan per run, dan itulah yang Anda inginkan ketika tiga schedule mengarah ke mesin yang sama. Jika Anda lebih suka melakukan penyebaran di dalam satu proses alih-alih satu op per URL, SDK Python menyertakan client asyncio yang sesungguhnya dan pendekatan itu dibahas di panduan memecahkan CAPTCHA secara paralel.

Menjalankan solver di mesin lain

Container code location, agent Kubernetes dan run Serverless sama-sama memindahkan op Anda menjauh dari komputer Anda. Tidak ada yang berubah pada kode kecuali host-nya, dan resource sudah membacanya dari config.

CapSkip punya dua mode koneksi. Local mengikat ke 127.0.0.1 dan hanya melayani perangkat itu. Server mengikat ke IP jaringan atau IP publik Anda, sehingga sebuah container, VM atau run Serverless menjangkau mesin Windows yang sama melalui API. IP publik statis menjaga alamatnya tetap stabil. Perangkat kerasnya tetap milik Anda dan tetap tanpa batas pemakaian pada kedua mode, jadi biaya sebuah hari yang sibuk tidak berubah karena modenya.

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

Aktifkan validasi key begitu pemecah mendengarkan di alamat jaringan, dan beri setiap code location key-nya sendiri agar satu key bisa dicabut tanpa mengganggu yang lain. Kedua mode dibahas lengkap di panduan penyiapan CapSkip.

Kesalahan umum dan artinya

Apa yang Anda lihatPenyebabPerbaiki
Formulir menolak token yang tampak benarToken melewati sebuah IO manager di antara dua op dan tiba dalam keadaan basiPecahkan dan kirim di dalam satu op agar token tetap menjadi variabel lokal
Eksekusi ulang mengirimkan token dari run lamaDagster memuat output yang tersimpan sebelumnya alih-alih menghitung ulangPerbaikan yang sama. Token tidak boleh menjadi output op atau asset
Sebuah file pickle muncul di dalam DAGSTER_HOME/storageIO manager filesystem default menulis token Anda ke diskPerbaikan yang sama, lalu hapus file itu. Perlakukan sebagai kredensial yang bocor
NetworkException di Dagster+ ServerlessRun dieksekusi di lingkungan Dagster, bukan lingkungan AndaAlihkan pemecah ke Server mode, atau gunakan agent Hybrid
NetworkException di agent Anda sendiriCapSkip tidak berjalan, atau host-nya salahJalankan aplikasinya, atau arahkan host resource ke alamat server
Tiga retry terbuang pada sitekey yang salahSetiap percobaan gagal dengan cara deterministik yang samaMunculkan RetryRequested hanya untuk NetworkException dan TimeoutException
Pemecah kelebihan beban selama backfillPartisi berjalan bersamaan secara defaultMasukkan op ke dalam sebuah pool dan tetapkan batas padanya
ValidationExceptionSitekey atau URL halaman yang hilang atau salah bentukCatat keduanya sebelum pemanggilan dan pastikan sitekey-nya adalah yang aktif

FAQ

Haruskah pemecahan dijadikan sebuah asset?

Tidak. Asset adalah objek tahan lama dengan riwayat materialisasi, dan token yang kedaluwarsa dalam dua menit adalah kebalikannya. Modelkan hal yang benar-benar Anda hasilkan sebagai asset, entah itu catatan yang dikirimkan atau halaman yang akhirnya boleh Anda baca, dan biarkan pemecahan tetap berupa op di dalamnya.

Apakah sebuah run Dagster+ benar-benar bisa menjangkau 127.0.0.1?

Pada Hybrid, ya, karena agent di infrastruktur Anda sendirilah yang mengeksekusi kodenya. Loopback di sana berarti mesin milik agent, jadi pemecah di mesin itu menjawab dengan normal. Pada Serverless, run terjadi di lingkungan Dagster dan Anda memerlukan Server mode dengan alamat yang dapat dijangkau.

Apa bedanya dengan melakukannya di Airflow?

Terutama pada cara nilai berpindah. Airflow mendorong nilai kecil melalui XCom dan Anda cukup menghindari menaruh token di sana. Dagster secara default menyimpan setiap output melalui sebuah IO manager, jadi kesalahan yang sama menghasilkan file pickle alih-alih satu baris database. Sisi pemecahnya identik, dan versi Airflow-nya ditulis di panduan CAPTCHA Airflow.

Apakah pemecahan yang lama memblokir sisa run?

Op tersebut terblokir selama melakukan polling, tetapi executor multiprocess Dagster tetap menjalankan op lain yang tidak bergantung padanya. Batas pool membatasi berapa banyak pemecahan yang berlangsung sekaligus tanpa menghentikan hal lain. Di perangkat keras Anda sendiri, satu-satunya biaya dari pemecahan yang lambat hanyalah waktu yang berlalu.

Versi singkatnya

Bungkus pemecah dalam sebuah ConfigurableResource dan baca key-nya dari EnvVar. Beri op tersebut RetryPolicy dengan backoff eksponensial, dan munculkan sendiri RetryRequested ketika Anda ingin melewati percobaan ulang yang tidak mungkin berhasil. Yang terpenting, jaga pemecahan dan pengiriman tetap berada di op yang sama, karena Dagster secara default menyimpan output op dan token bukanlah sesuatu yang layak disimpan.

Sisi Python dari semua ini dibahas di halaman pemecah CAPTCHA Python. Detail tipe CAPTCHA tersebut ada di halaman pemecah reCAPTCHA v2. Tiga panggilan yang sama juga ada di Node.js, PHP, dan C#, dan semuanya tercantum di halaman SDK pemecahan CAPTCHA.

Satu hal yang perlu diketahui sebelum Anda menjadwalkan ini terhadap backfill besar. CapSkip adalah pemecah captcha yang berjalan di perangkat keras yang sudah Anda miliki, jadi job yang memecahkan lima puluh ribu baris berbiaya persis sama dengan job yang memecahkan lima puluh.