Dagster Pipeline में कैप्चा कैसे हल करें (Ops और Assets)

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

Dagster में कैप्चा step एक op है जिस पर एक retry policy लगी होती है, और उसके आसपास की पूरी pipeline लगभग तीस लाइनों की है। जो हिस्सा Dagster के लिए ख़ास है, और जिसके बारे में कोई दूसरा orchestrator आपको चेतावनी नहीं देगा, वह यह है कि op जो value लौटाता है उसका क्या होता है। Dagster, op और asset outputs को एक IO manager के ज़रिए persist करता है, और default वाला उन्हें pickle करके disk पर लिख देता है। कैप्चा token दो मिनट की उम्र वाला single use credential है, इसलिए उसके पहुँचने की सबसे आख़िरी जगह यही होनी चाहिए।

आपको क्या चाहिए

  • Dagster 1.9 या नया और Python 3.10 या नया, साथ में CapSkip Python SDK।
  • Dagster open source, या कोई Dagster+ deployment। दोनों ही सूरतों में code एक जैसा है।
  • सुरक्षित form का पेज URL, और उसकी sitekey।
  • Local mode में चलता CapSkip, जब code और solver एक ही machine साझा करते हैं, या Server mode में जब वे ऐसा नहीं करते। दोनों का वर्णन यहाँ है: कनेक्शन सेटिंग्स.
# pip install capskip
pip install dagster dagster-webserver capskip

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

आपका code असल में कहाँ चलता है

कोई host string चुनने से पहले इसका जवाब दें, क्योंकि यही पूरा network setup तय करता है। Dagster open source पूरी तरह आपकी अपनी machines पर चलता है, इसलिए यहाँ चर्चा करने को कुछ है ही नहीं। Dagster+ के दो रूप हैं और इस मामले में वे एक दूसरे के उलट हैं।

Hybrid वह रूप है जो आपकी उम्मीद के मुताबिक़ बर्ताव करता है। आप अपने ही infrastructure के अंदर एक agent चलाते हैं और वह बाहर control plane से जुड़ता है। Dagster+ के पास आपके network में कोई ingress नहीं है, वह न आपका code देखता है और न आपके data को छूता है, इसलिए cloud UI से शुरू किया गया run भी आपके अपने hardware पर ही चलता है। agent को CapSkip वाली machine पर ही रखें और op ठीक वैसे ही 127.0.0.1:8080 को कॉल करता है जैसे आपकी मेज़ पर रखी कोई script करती। न कोई tunnel और न कोई public address।

Serverless अपवाद है। वहाँ आपका code आपके नहीं, बल्कि Dagster के environment में चलता है, और loopback का कोई काम का मतलब नहीं रह जाता: वह Dagster के container की ओर इशारा करता है, जहाँ कुछ भी सुन नहीं रहा। उस configuration में solver को Server mode चाहिए और ऐसा address चाहिए जहाँ तक run पहुँच सके। यह कोई गिरावट नहीं है: यह उसी hardware पर वही solver है, बस अलग bind address के साथ, और इसमें अब भी प्रति हल कोई शुल्क नहीं लगता।

चरण 1: solver को एक resource में लपेटें

किसी भी बाहरी चीज़ के लिए Dagster का तरीक़ा एक resource है, module level client नहीं। ConfigurableResource को subclass करें, connection fields घोषित करें, और UI में उनके लिए एक launchpad entry मिल जाती है, run शुरू होने से पहले config validate हो जाता है, और tests पूरी चीज़ को किसी fake से बदल सकते हैं। key को hardcode करने के बजाय environment से पढ़ें।

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

एक ही method reCAPTCHA v2, Invisible, Enterprise और v3 को कवर करता है, क्योंकि variants अलग कॉल नहीं बल्कि keyword options हैं: invisible को 1, enterprise को 1, या version को v3 के साथ एक action। Turnstile और GeeTest के अपने methods हैं और आकार वही। parameters की पूरी सूची यहाँ है: CapSkip API डॉक्यूमेंटेशन.

चरण 2: एक op, एक retry policy के साथ

solver कॉल एक network कॉल है, और network कॉल विफल होते हैं। Dagster की RetryPolicy declarative है और decorator पर लगती है: एक अधिकतम count, एक base delay, एक backoff curve और एक jitter mode। plus या minus jitter के साथ exponential backoff समझदार default है, क्योंकि यह एक साथ हुई विफलताओं के झुंड को फैला देता है, उन सबको एक ही बार में वापस लाने के बजाय।

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

ध्यान दें कि वह op क्या लौटाता है: एक status code, token नहीं। अगले हिस्से का पूरा मतलब यही है।

चरण 3: token को कभी किसी op की सीमा पार न करने दें

यह Dagster वाला ख़ास जाल है और इसमें फँसना आसान है, क्योंकि सीधा सा design यही लगता है कि एक op हल करे और दूसरा submit करे। Dagster, ops के बीच values memory में नहीं भेजता। हर output किसी IO manager से होकर जाता है, और default वाला filesystem IO manager है, जो outputs को local disk पर pickle files के रूप में रखता है। इसलिए किसी op से लौटाया गया token एक file में लिखा जाता है, अगला op उसे वापस पढ़ता है, और उसके बाद वह वहीं पड़ा रह जाता है।

यह दो तरह से ग़लत है। यह एक ज़िंदा credential को disk पर लिख देता है, जहाँ उसे कोई साफ़ नहीं करता। और यह हल और submit के बीच एक storage round trip डाल देता है, यानी ठीक वही देरी जो दो मिनट की expiry बर्दाश्त नहीं कर सकती। re-execution पर बात और बिगड़ती है: Dagster दोबारा गणना करने के बजाय पिछले run का सहेजा हुआ output load कर सकता है, और तब form को ऐसा token थमा दिया जाता है जो कल ही expire हो चुका था।

इसका हल कोई चतुर IO manager नहीं है। हल यह है कि हल करना और submit करना, दोनों एक ही op के अंदर रखें, ताकि token एक local variable में रहे और कभी output बने ही नहीं। उससे पहले और उसके बाद की हर चीज़ अलग ops या assets के रूप में रह सकती है। सिर्फ़ वही जोड़ी मिलाई जाती है जिसे token साझा करना ही है।

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

अगर आप ops के बजाय assets में काम कर रहे हैं, तो वही नियम एक अतिरिक्त बात के साथ लागू होता है: materialise किया गया asset एक स्थायी रिकॉर्ड होता है जिसका UI में इतिहास रहता है, और किसी single use secret का ऐसा बनना बनता ही नहीं। हल को किसी graph backed asset के अंदर एक op के रूप में रखें, और asset को submission का नतीजा materialise करने दें।

पूरा चलने वाला उदाहरण

एक पूरा job। page fetch करें, उसमें से sitekey निकालें, फिर एक ही op में हल करें और submit करें। दो ops, एक resource, एक Definitions object।

# 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, value को definition में पका देने के बजाय lookup को run time तक टाल देता है, इसलिए key न कभी आपकी repository में होती है और न किसी serialised snapshot में।

उन विफलताओं को retry न करें जो कभी पास नहीं होंगी

एक declarative policy हर चीज़ को retry करती है, जिससे ऐसे ख़राब argument पर तीन प्रयास और एक मिनट का backoff बर्बाद होता है जो हर बार बिल्कुल वैसे ही विफल होगा। Dagster का रास्ता यह है कि retry आप खुद उठाएं: जिन exceptions पर एक और कोशिश बनती है उन्हें catch करें, साफ़ तौर पर retry माँगें, और बाक़ी को आगे बढ़ने दें ताकि run तुरंत विफल हो जाए।

# 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 चल नहीं रहा या host ग़लत है, और TimeoutException का मतलब है कि हल recaptchaTimeout से लंबा खिंच गया, जिसकी default 300 सेकंड है। दोनों एक और कोशिश के हक़दार हैं। ValidationException का मतलब है कोई ख़राब argument और ApiException का मतलब है कि API ने request ठुकरा दी, और दूसरी कोशिश में इनमें से कोई नहीं सुधरता। चारों एक साझा base से निकलते हैं जिसका नाम CapSkipError है, इसलिए अगर आप सब कुछ एक ही जगह संभालना चाहें तो उसी को catch कर सकते हैं।

किसी batch को concurrency pool से throttle करना

किसी partitioned job को दो सौ URLs पर फैलाइए और Dagster ख़ुशी ख़ुशी सबको चलाने की कोशिश करेगा, जो एक अकेले solver पर डालने लायक़ भार से कहीं ज़्यादा है। इसका नियंत्रण concurrency pools हैं। op को एक pool नाम से tag करें, limit एक बार सेट करें, और cap से ऊपर की हर चीज़ ढेर लगाने के बजाय क़तार में लग जाती है।

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

ऊपर दिए उदाहरण में op पर pool argument पहले से मौजूद है। Limits हर run पर अलग नहीं, बल्कि सभी runs पर मिलाकर लागू होती हैं, और जब तीन schedules एक ही machine की ओर इशारा कर रहे हों तो आपको यही चाहिए। अगर आप हर URL के लिए एक op के बजाय fan out को एक ही process के अंदर करना चाहें, तो Python SDK एक असली asyncio client के साथ आता है और वह तरीक़ा यहाँ समझाया गया है: समानांतर में कैप्चा हल करने की गाइड.

सॉल्वर को किसी दूसरी मशीन पर चलाना

Code location containers, Kubernetes agents और Serverless runs, ये सब आपके op को आपकी मेज़ से दूर ले जाते हैं। code में host के अलावा कुछ नहीं बदलता, और resource उसे पहले से ही config से पढ़ता है।

CapSkip के दो connection modes हैं। Local, 127.0.0.1 पर bind होता है और सिर्फ़ उसी device को जवाब देता है। Server, आपके network या public IP पर bind होता है, इसलिए कोई container, VM या Serverless run उसी Windows machine तक API के ज़रिए पहुँचता है। एक static public IP उस address को स्थिर रखता है। यह अब भी आपका अपना hardware है और दोनों में unmetered ही है, इसलिए किसी व्यस्त दिन की लागत mode के साथ नहीं बदलती।

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

जैसे ही solver किसी network address पर सुनने लगे, key validation चालू करें, और हर code location को उसकी अपनी key दें ताकि एक को बाक़ियों को छुए बिना revoke किया जा सके। दोनों modes को विस्तार से यहाँ समझाया गया है: CapSkip सेटअप गाइड.

आम errors और उनका मतलब

आप जो देखते हैंकारणफिक्स
form ऐसे token को ठुकरा देता है जो देखने में सही लगता हैवह दो ops के बीच एक IO manager से होकर गुज़रा और बासी पहुँचाएक ही op के अंदर हल करें और submit करें ताकि token एक local variable बना रहे
कोई re-execution किसी पुराने run का token submit कर देता हैDagster ने दोबारा गणना करने के बजाय पहले से सहेजा हुआ output load कर लियावही हल। token कभी किसी op या asset का output नहीं होना चाहिए
DAGSTER_HOME/storage के नीचे एक pickle file दिखाई देती हैdefault filesystem IO manager ने आपका token disk पर लिख दियावही हल, फिर उस file को हटा दें। इसे लीक हुआ credential मानें
Dagster+ Serverless पर NetworkExceptionrun, आपके नहीं बल्कि Dagster के environment में चलाsolver को Server mode पर लाएं, या कोई Hybrid agent इस्तेमाल करें
आपके अपने agent पर NetworkExceptionCapSkip चल नहीं रहा, या host ग़लत हैapp शुरू करें, या resource के host को server के address पर लगाएं
ख़राब sitekey पर तीन retries बर्बादहर प्रयास एक ही तय तरीक़े से fail होता हैRetryRequested सिर्फ़ NetworkException और TimeoutException के लिए उठाएं
किसी backfill के दौरान solver पर बोझ बढ़ जाता हैPartitions default रूप से एक साथ चलती हैंop को एक pool में डालें और उस पर एक limit सेट करें
ValidationExceptionग़ायब या ख़राब sitekey या page URLकॉल से पहले दोनों को log करें और जाँचें कि sitekey वही live वाली है

FAQ

क्या हल को एक asset होना चाहिए?

नहीं। asset एक टिकाऊ object है जिसका materialisation इतिहास होता है, और दो मिनट में expire होने वाला token उसका ठीक उलटा है। आपने असल में जो चीज़ बनाई है उसी को asset बनाएं, चाहे वह submit किया गया record हो या वह page जिसे आख़िरकार आपको पढ़ने दिया गया, और हल को उसके अंदर एक op ही रहने दें।

क्या कोई Dagster+ run सचमुच 127.0.0.1 तक पहुँच सकता है?

Hybrid पर, हाँ, क्योंकि code को चलाने वाला आपके अपने infrastructure में मौजूद agent ही है। वहाँ loopback का मतलब agent की machine है, इसलिए उस machine पर चल रहा solver सामान्य रूप से जवाब देता है। Serverless पर run, Dagster के environment में होता है और आपको Server mode के साथ एक पहुँच योग्य address चाहिए।

यह Airflow में करने से कैसे अलग है?

ज़्यादातर इसमें कि values कैसे चलती हैं। Airflow छोटी values को XCom के ज़रिए भेजता है और आप बस वहाँ token डालने से बचते हैं। Dagster default रूप से हर output को एक IO manager के ज़रिए persist करता है, इसलिए वही ग़लती किसी database row के बजाय एक pickle file लिख देती है। solver वाला हिस्सा एक जैसा ही है, और Airflow वाला रूप यहाँ लिखा गया है: Airflow कैप्चा गाइड.

क्या एक लंबा हल बाक़ी run को रोक देता है?

जब तक op poll करता है तब तक वह ब्लॉक रहता है, लेकिन Dagster का multiprocess executor उन दूसरे ops को चलाता रहता है जिनकी उस पर कोई निर्भरता नहीं है। एक pool limit तय कर देती है कि एक साथ कितने हल हों, और इससे बाक़ी कुछ नहीं रुकता। आपके अपने hardware पर एक धीमे हल की इकलौती लागत wall clock समय है।

संक्षेप में

solver को एक ConfigurableResource में लपेटें और key को EnvVar से पढ़ें। op को exponential backoff वाली एक RetryPolicy दें, और जब आप ऐसे retries छोड़ना चाहें जो सफल हो ही नहीं सकते तो RetryRequested खुद उठाएं। सबसे बढ़कर, हल और submit दोनों को एक ही op में रखें, क्योंकि Dagster default रूप से op outputs को persist करता है और token कोई persist करने की चीज़ नहीं है।

इस सबका Python वाला हिस्सा यहाँ कवर किया गया है: Python कैप्चा सॉल्वर पेज। उस कैप्चा प्रकार का ब्योरा यहाँ मिलता है: reCAPTCHA v2 सॉल्वर पेज। Node.js, PHP और C# में भी वही तीन calls मौजूद हैं, और वे यहाँ सूचीबद्ध हैं: CAPTCHA solving SDK पेज.

किसी बड़े backfill पर इसे schedule करने से पहले एक बात जान लें। CapSkip एक ऐसा कैप्चा सॉल्वर है जो आपके पास पहले से मौजूद hardware पर चलता है, इसलिए पचास हज़ार rows हल करने वाले job की क़ीमत ठीक उतनी ही है जितनी पचास हल करने वाले job की।