Bir Dagster Pipeline'ında CAPTCHA Nasıl Çözülür (Op'lar ve Asset'ler)

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

Dagster'da bir captcha adımı, üzerinde bir retry policy bulunan tek bir op'tur ve etrafındaki pipeline yaklaşık otuz satırdır. Dagster'a özgü olan ve başka hiçbir orkestratörün sizi uyarmayacağı kısım, op'un döndürdüğü değere ne olduğudur. Dagster, op ve asset çıktılarını bir IO manager üzerinden kalıcı hale getirir ve varsayılan olan onları diske pickle eder. Bir CAPTCHA token'ı iki dakikalık ömrü olan tek kullanımlık bir kimlik bilgisidir, dolayısıyla orası onun varabileceği en son yerdir.

Neye ihtiyacınız var

  • Dagster 1.9 veya daha yenisi ve Python 3.10 veya daha yenisi, ayrıca CapSkip Python SDK'sı.
  • Dagster açık kaynak ya da bir Dagster+ dağıtımı. Kod her iki durumda da aynıdır.
  • Korumalı formun sayfa URL’si ve sitekey'i.
  • Kod ile çözücü aynı makineyi paylaşıyorsa Local mode, paylaşmıyorsa Server mode ile çalışan CapSkip. Her ikisi de şurada anlatılır: bağlantı ayarları.
# pip install capskip
pip install dagster dagster-webserver capskip

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

Kodunuz aslında nerede çalışıyor

Bir host dizesi seçmeden önce buna cevap verin, çünkü tüm ağ kurulumunu bu belirler. Dagster açık kaynak tamamen kendi makinelerinizde çalışır, dolayısıyla burada tartışılacak bir şey yok. Dagster+'ın iki biçimi vardır ve bu amaç açısından birbirlerinin tam zıddıdırlar.

Hybrid, umduğunuz gibi davranan seçenektir. Kendi altyapınızın içinde bir agent çalıştırırsınız ve o, dışarıya doğru kontrol düzlemine bağlanır. Dagster+'ın ağınıza girişi yoktur, kodunuzu görmez ve verilerinize dokunmaz; bu yüzden bulut arayüzünden başlatılan bir çalıştırma yine de sizin sahip olduğunuz donanımda yürütülür. Agent'ı CapSkip ile aynı makineye koyun; op, tıpkı masanızdaki bir script gibi 127.0.0.1:8080 adresini çağırır. Tünel yok, genel adres yok.

Serverless istisnadır. Orada kodunuz sizin ortamınızda değil Dagster’ın ortamında yürütülür ve loopback işe yarar bir anlam taşımayı bırakır: hiçbir şeyin dinlemediği Dagster’ın konteynerini işaret eder. Bu yapılandırmada çözücünün Server mode'a ve çalıştırmanın erişebileceği bir adrese ihtiyacı vardır. Bu bir gerileme değildir: aynı donanımdaki aynı çözücüdür, yalnızca bind adresi farklıdır ve çözüm başına yine hiçbir ücret çıkarmaz.

Adım 1: çözücüyü bir resource içine sarın

Dagster’ın dış dünyaya ait her şey için kullandığı yaklaşım, modül düzeyinde bir istemci değil bir resource'tur. ConfigurableResource'tan türetin, bağlantı alanlarını tanımlayın; böylece arayüzde bunlar için bir launchpad girdisi oluşur, yapılandırma çalıştırma başlamadan önce doğrulanır ve testler tüm bu yapıyı sahte bir nesneyle değiştirebilir. Anahtarı koda gömmek yerine ortamdan okuyun.

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

Tek bir metot reCAPTCHA v2, Invisible, Enterprise ve v3'ü kapsar, çünkü varyantlar ayrı çağrılar değil anahtar kelime seçenekleridir: invisible 1 olarak ayarlanır, enterprise 1 olarak ayarlanır veya version bir action ile birlikte v3 olarak ayarlanır. Turnstile ve GeeTest'in kendi metotları ve aynı yapısı vardır. Tam parametre listesi şurada: CapSkip API dokümantasyonu.

Adım 2: retry policy'si olan tek bir op

Çözücü çağrısı bir ağ çağrısıdır ve ağ çağrıları başarısız olur. Dagster’ın RetryPolicy yapısı bildirimseldir ve dekoratörün üzerine gelir: bir maksimum sayı, bir temel gecikme, bir geri çekilme eğrisi ve bir jitter modu. Artı eksi jitter'lı üstel geri çekilme makul varsayılandır, çünkü eşzamanlı hata yığınını hepsini bir anda geri getirmek yerine zamana yayar.

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

O op'un ne döndürdüğüne dikkat edin: bir token değil, bir durum kodu. Bir sonraki bölümün tüm amacı budur.

Adım 3: token'ın bir op sınırını geçmesine asla izin vermeyin

Dagster'a özgü tuzak budur ve içine düşmek kolaydır, çünkü akla ilk gelen tasarım çözen bir op ile gönderen ikinci bir op'tur. Dagster op'lar arasında değerleri bellekte aktarmaz. Her çıktı bir IO manager üzerinden geçer ve varsayılan olan, çıktıları yerel diskte pickle dosyaları olarak saklayan dosya sistemi yöneticisidir. Dolayısıyla bir op'tan dönen bir token bir dosyaya yazılır, sonraki op tarafından geri okunur ve ardından orada öylece kalır.

Bu iki yönden birden yanlıştır. Canlı bir kimlik bilgisini, hiçbir şeyin temizlemediği diske yazar. Ayrıca çözüm ile gönderim arasına bir depolama gidiş dönüşü koyar ki bu, iki dakikalık bir geçerlilik süresinin kaldıramayacağı gecikmenin ta kendisidir. Yeniden yürütmede işler daha da kötüleşir: Dagster yeniden hesaplamak yerine önceki bir çalıştırmanın saklanmış çıktısını yükleyebilir ve o zaman forma dün süresi dolmuş bir token verilir.

Çözüm akıllı bir IO manager değildir. Çözme ile gönderimi aynı op'un içinde tutmaktır; böylece token yerel bir değişkende yaşar ve hiçbir zaman bir çıktıya dönüşmez. Yukarı ve aşağı akıştaki her şey ayrı op'lar veya asset'ler olarak kalabilir. Yalnızca token'ı paylaşmak zorunda olan ikili birleştirilir.

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

Op'lar yerine asset'lerle çalışıyorsanız aynı kural bir ek notla geçerlidir: materyalize edilmiş bir asset, arayüzde geçmişi olan kalıcı bir kayıttır ve tek kullanımlık bir sırrın böyle bir şey olmasının hiçbir gereği yoktur. Çözmeyi, graph tabanlı bir asset'in içindeki bir op olarak modelleyin ve gönderimin sonucunu asset materyalize etsin.

Tam çalışan örnek

Eksiksiz bir job. Sayfayı getirin, sitekey'i içinden çekin, ardından tek bir op içinde çözün ve gönderin. İki op, bir resource, bir Definitions nesnesi.

# 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, değeri tanımın içine gömmek yerine aramayı çalışma zamanına erteler; böylece anahtar ne deponuzda ne de serileştirilmiş bir anlık görüntüde bulunur.

Asla başarılı olmayacak hataları yeniden denemeyin

Bildirimsel bir politika her şeyi yeniden dener; bu da her seferinde aynı şekilde başarısız olacak hatalı biçimlendirilmiş bir argüman için üç denemeyi ve bir dakikalık geri çekilmeyi boşa harcar. Dagster’ın kaçış yolu, yeniden denemeyi kendinizin fırlatmasıdır: bir şans daha vermeye değer istisnaları yakalayın, açıkça yeniden deneme isteyin ve geri kalanların yayılmasına ve çalıştırmayı hemen başarısız kılmasına izin verin.

# 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, CapSkip'in çalışmadığı veya host'un yanlış olduğu anlamına gelir; TimeoutException ise çözümün, varsayılanı 300 saniye olan recaptchaTimeout değerini aştığı anlamına gelir. İkisi de bir deneme daha hak eder. ValidationException hatalı bir argüman, ApiException ise API'nin isteği reddettiği anlamına gelir ve ikisi de ikinci denemede düzelmez. Dördü de CapSkipError adlı ortak bir temel sınıftan türer, dolayısıyla her şeyi tek bir yerde ele almayı tercih ediyorsanız bunun yerine onu yakalayabilirsiniz.

Bir grubu concurrency pool ile sınırlamak

Bölümlenmiş bir job'ı iki yüz URL'ye dağıtın, Dagster hepsini seve seve çalıştırmaya kalkar; bu da tek bir çözücüye yöneltmek isteyeceğinizden fazla yüktür. Kontrol mekanizması concurrency pool'lardır. Op'u bir pool adıyla etiketleyin, sınırı bir kez ayarlayın; sınırın üzerindeki her şey üst üste yığılmak yerine kuyruğa girer.

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

Yukarıdaki örnekte pool argümanı zaten op'un üzerinde. Sınırlar çalıştırma başına değil, tüm çalıştırmalar genelinde geçerlidir; üç zamanlama aynı makineye yöneldiğinde tam olarak isteyeceğiniz şey budur. URL başına bir op yerine dağıtımı tek bir işlem içinde yapmayı tercih ederseniz, Python SDK'sı gerçek bir asyncio istemcisiyle gelir ve bu yaklaşım şurada ele alınıyor: CAPTCHA'ları paralel çözme kılavuzu.

Çözücüyü başka bir makinede çalıştırmak

Code location konteynerleri, Kubernetes agent'ları ve Serverless çalıştırmaları op'unuzu masanızdan uzaklaştırır. Kodda host dışında hiçbir şey değişmez ve resource onu zaten yapılandırmadan okur.

CapSkip'in iki bağlantı modu vardır. Local, 127.0.0.1'e bağlanır ve yalnızca o cihaza yanıt verir. Server ise ağ IP'nize veya genel IP'nize bağlanır, böylece bir konteyner, bir VM veya bir Serverless çalıştırması aynı Windows makinesine API üzerinden ulaşır. Statik bir genel IP adresi sabit tutar. Her iki durumda da donanım yine sizindir ve yine sayaçsızdır, bu yüzden yoğun bir günün maliyeti moda göre değişmez.

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

Çözücü bir ağ adresini dinlemeye başladığında anahtar doğrulamasını açın ve her code location'a kendi anahtarını verin ki biri diğerlerine dokunmadan iptal edilebilsin. Her iki mod da şurada adım adım anlatılıyor: CapSkip kurulum kılavuzu.

Sık görülen hatalar ve anlamları

GördüğünüzNedenDüzeltme
Form, doğru görünen bir token'ı reddediyorToken iki op arasında bir IO manager üzerinden geçti ve bayatlamış olarak ulaştıToken'ın yerel bir değişken olarak kalması için çözme ve gönderme işlemini tek bir op içinde yapın
Yeniden yürütme, eski bir çalıştırmadan gelen bir token gönderiyorDagster yeniden hesaplamak yerine daha önce saklanan çıktıyı yüklediAynı çözüm. Bir token asla bir op veya asset çıktısı olmamalıdır
DAGSTER_HOME/storage altında bir pickle dosyası beliriyorVarsayılan dosya sistemi IO manager'ı token'ınızı diske yazdıAynı çözüm, ardından dosyayı silin. Onu sızmış bir kimlik bilgisi olarak değerlendirin
Dagster+ Serverless üzerinde NetworkExceptionÇalıştırma sizin ortamınızda değil, Dagster’ın ortamında yürütüldüÇözücüyü Server mode'a geçirin veya bir Hybrid agent kullanın
Kendi agent'ınızda NetworkExceptionCapSkip çalışmıyor ya da ana makine adresi yanlışUygulamayı başlatın veya resource host'unu sunucu adresine yönlendirin
Hatalı bir sitekey'e üç yeniden deneme harcandıHer deneme aynı deterministik şekilde başarısız oluyorRetryRequested'ı yalnızca NetworkException ve TimeoutException için fırlatın
Bir backfill sırasında çözücü aşırı yükleniyorBölümler varsayılan olarak eşzamanlı çalışırOp'u bir pool'a koyun ve üzerine bir sınır ayarlayın
ValidationExceptionEksik veya hatalı biçimlendirilmiş bir sitekey ya da sayfa URL'siÇağrıdan önce ikisini de loglayın ve sitekey'in canlı olan olduğunu doğrulayın

FAQ

Çözme işlemi bir asset olmalı mı?

Hayır. Bir asset, materyalizasyon geçmişi olan kalıcı bir nesnedir ve iki dakikada sona eren bir token bunun tam tersidir. Asıl ürettiğiniz şeyi asset olarak modelleyin: bu, gönderilen kayıt ya da sonunda okumanıza izin verilen sayfa olabilir; çözme işlemini de onun içinde bir op olarak tutun.

Bir Dagster+ çalıştırması gerçekten 127.0.0.1 adresine ulaşabilir mi?

Hybrid'de evet, çünkü kodu yürüten şey kendi altyapınızdaki agent'tır. Oradaki loopback, agent'ın makinesi anlamına gelir, dolayısıyla o makinedeki bir çözücü normal şekilde yanıt verir. Serverless'ta ise çalıştırma Dagster’ın ortamında gerçekleşir ve erişilebilir bir adresle birlikte Server mode'a ihtiyacınız olur.

Bunu Airflow’da yapmaktan farkı ne?

Çoğunlukla değerlerin nasıl hareket ettiği konusunda. Airflow küçük değerleri XCom üzerinden aktarır ve siz oraya bir token koymaktan kaçınırsınız, o kadar. Dagster ise varsayılan olarak her çıktıyı bir IO manager üzerinden kalıcı hale getirir, dolayısıyla aynı hata bir veritabanı satırı yerine bir pickle dosyası yazar. Çözücü tarafı aynıdır ve Airflow sürümü şurada anlatılıyor: Airflow CAPTCHA kılavuzu.

Uzun bir çözüm çalıştırmanın geri kalanını bloke eder mi?

Op, sorgulama yaparken bloke olur, ancak Dagster’ın multiprocess executor'ı ona bağımlılığı olmayan diğer op'ları çalıştırmayı sürdürür. Bir pool sınırı, başka hiçbir şeyi durdurmadan aynı anda kaç çözümün yapılacağını kısıtlar. Kendi donanımınızda yavaş bir çözümün tek maliyeti duvar saati süresidir.

Kısa özet

Çözücüyü bir ConfigurableResource içine sarın ve anahtarı EnvVar'dan okuyun. Op'a üstel geri çekilmeli bir RetryPolicy verin ve başarılı olamayacak yeniden denemeleri atlamak istediğinizde RetryRequested'ı kendiniz fırlatın. Her şeyden önce, çözme ile gönderme işlemini aynı op içinde tutun, çünkü Dagster op çıktılarını varsayılan olarak kalıcı hale getirir ve bir token kalıcı hale getirilecek bir şey değildir.

Tüm bunların Python tarafı şurada ele alınıyor: Python CAPTCHA çözücü sayfası. O CAPTCHA türünün ayrıntıları şurada yer alıyor: reCAPTCHA v2 çözücü sayfası. Aynı üç çağrı Node.js, PHP ve C#'ta da mevcuttur ve şurada listelenmiştir: CAPTCHA çözme SDK sayfası.

Bunu büyük bir backfill için zamanlamadan önce bilmeniz gereken bir şey. CapSkip, bir captcha çözücü olup zaten sahip olduğunuz donanımda çalışır, bu yüzden elli bin satır çözen bir job, elli satır çözen bir job ile tam olarak aynı maliyettedir.