Apache Airflow DAGs में कैप्चा कैसे हल करें (Python)

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

Airflow में कैप्चा वाला step एक सामान्य Python task है। आप सॉल्वर को call करते हैं, token वापस मिलता है, उसी task में उसका इस्तेमाल कर लेते हैं। न कोई operator इंस्टॉल करना है और न कोई plugin लिखना है। लोग असल में जिस चीज़ पर अटकते हैं वह है जगह: आपके workers वहीं चलते हैं जहाँ scheduler ने उन्हें रखा है, और आपके laptop के loopback पते से बँधा सॉल्वर किसी दूसरे host पर चल रहे container की पहुँच में नहीं होता। पहले यह हिस्सा ठीक कर लें, उसके बाद DAG पंद्रह लाइनों का है।

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

  • Airflow 2.x या 3.x, और उसी image या virtualenv में CapSkip SDK इंस्टॉल किया हुआ जिसे आपके workers इस्तेमाल करते हैं।
  • ऐसा target जो कोई challenge लौटाता हो। किसी widget के पीछे बैठा कोई page, न कि कोई test key जो हमेशा पास हो जाती है।
  • CapSkip ऐसी मशीन पर Server mode में चालू हो जहाँ तक आपके workers पहुँच सकें, या Local mode में अगर worker और सॉल्वर एक ही मशीन हैं। दोनों यहाँ बताए गए हैं: कनेक्शन सेटिंग्स.
# Install into the worker environment, not just the scheduler.
pip install capskip

सॉल्वर को कहाँ रहना चाहिए

पूरी समस्या यही है, इसलिए यह सबसे पहले। कोई task किसी worker process के अंदर चलता है, और ज़्यादातर असली deployments में वह worker ऐसे host पर बैठा container होता है जो उन सब चीज़ों से अलग है जिन्हें आप हाथ से चलाते हैं। उस container के अंदर का loopback पता ख़ुद वही container है, इसलिए SDK को उस पर लगाने से कुछ नहीं मिलता।

मोडकिस पर सुनता हैकब इस्तेमाल करें
लोकल127.0.0.1, केवल उसी डिवाइस परसॉल्वर वाली उसी मशीन पर अकेला एक worker
सर्वरआपका नेटवर्क पता या पब्लिक IPContainers, workers का बेड़ा, कोई VPS या कोई managed Airflow

लगभग हर Airflow deployment का जवाब Server mode है। आप app में listen address बदल देते हैं, SDK का host उसी मशीन पर कर देते हैं, और बेड़े का हर worker एक ही सॉल्वर साझा करता है। जब call करने वाले आपके अपने network के बाहर बैठे हों तो एक static public IP रखने की सलाह दी जाती है। इससे यह नहीं बदलता कि product क्या है: यह अब भी आपका अपना hardware है और अब भी unmetered है, इसलिए इसे loopback पते से हटाने पर सिर्फ़ यह बदलता है कि वह कहाँ चलता है, और कुछ नहीं। CapSkip एक Windows application है, इसलिए व्यवहार में यह एक Windows मशीन होती है जिसे पूरा बेड़ा call करता है।

Step 1: host को एक बार configure करें

पते को DAG file में hardcode न करें। SDK environment से CAPSKIP_HOST, CAPSKIP_PORT और CAPSKIP_API_KEY पढ़ता है, जो Airflow के लिए सबसे साफ़-सुथरा तरीक़ा है क्योंकि worker environment variables सेट करने का रास्ता आपके पास पहले से है। अगर आप इसे metadata database में रखना चाहें तो कोई Airflow Variable भी काम कर जाता है।

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

client को task के अंदर बनाएँ, module scope पर नहीं। Airflow हर DAG file को बार-बार parse करता रहता है, और import के समय बनी हर चीज़ हर parse पर बनती है, worker के साथ-साथ scheduler में भी।

Step 2: हल को एक ही task के रूप में लिखें

सबसे छोटा रास्ता TaskFlow decorators हैं। Airflow 3 इन्हें SDK module से import करता है, Airflow 2 decorators module से, और body दोनों ही सूरतों में एक जैसी रहती है।

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

पूरा integration इतना ही है। सॉल्वर की call synchronous है, यह आपके लिए एक चौथाई सेकंड से शुरू होने वाले backoff के साथ poll करती है, और एक dictionary लौटाती है जिसके code field में token होता है।

Step 3: token को tasks के बीच पास न करें

यह वह ग़लती है जिसका नाम लेना ज़रूरी है, क्योंकि Airflow इसे स्वाभाविक महसूस कराता है। TaskFlow के returns XComs बन जाते हैं, इसलिए token लौटाने वाला एक solve task और उसे इस्तेमाल करने वाला एक downstream task अच्छा design लगता है। यह एक दौड़ है। कोई reCAPTCHA token लगभग दो मिनट तक वैध रहता है, और दो Airflow tasks के बीच का अंतराल scheduler का फ़ैसला है जो आपके नियंत्रण में नहीं। इसमें queue की देरी, कोई भरा हुआ pool, या किसी worker का restart जोड़ दीजिए, और token रास्ते में ही expire हो जाता है।

यह विफलता रुक-रुक कर होती है, जो सबसे बुरी क़िस्म है: DAG testing में काम करता है और production में कुछ प्रतिशत बार फेल हो जाता है, form अस्वीकार होता है और सॉल्वर की तरफ़ से कोई error नहीं आता। reCAPTCHA token कितनी देर वैध रहता है, इस पर लिखी पोस्ट समय-सीमाओं को कवर करती है। किसी DAG में यह नियम एक ही लाइन में सिमट जाता है: हल और submit एक ही task के अंदर करें, और दोबारा हल करने का काम retry पर छोड़ दें।

पूरा 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()

वहाँ pool argument असल काम कर रहा है। कोई pool यह सीमित करता है कि पूरे deployment में एक साथ कितने task instances चलें, ताकि दो सौ runs का कोई backfill एक ही मशीन पर दो सौ एक साथ हल शुरू न कर दे। UI में captcha नाम का एक pool बनाएँ, उसे उतने parallel हल दे दें जितने आप चाहते हैं, और उसका नाम लेने वाला हर task उसी सीमा के पीछे क़तार में लग जाता है।

ऐसे retries जो नुक़सान नहीं, मदद करें

task पर retries सेट करें और पूरे हल तथा submit को दोहराने दें। चूँकि token task की body के अंदर लाया जाता है, इसलिए हर retry को अपने आप ताज़ा token मिलता है, और यही वह व्यवहार है जो आप चाहते हैं तथा यही वजह है कि दोनों steps साथ रहते हैं।

retry को थोड़ी देरी दें। जिस site ने अभी-अभी आपको challenge दिया है, उसी पर तुरंत की गई retry को दोबारा challenge मिलने की संभावना रहती है, और किसी scheduled pipeline में दो-चार मिनट की कोई क़ीमत नहीं। दो retries आमतौर पर काफ़ी हैं: तीसरी विफलता आम तौर पर या तो ग़लत sitekey होती है या ऐसा सॉल्वर जिस तक पहुँचा नहीं जा सकता, और इनमें से कोई भी इंतज़ार करने से ठीक नहीं होता।

आम errors

आप जो देखते हैंकारणफिक्स
हर task पर NetworkExceptionworker सॉल्वर तक नहीं पहुँच पा रहाServer mode पर जाएँ और worker पर CAPSKIP_HOST सेट करें
scheduler पर चलता है, worker पर फेल हो जाता हैworker image में SDK मौजूद नहीं हैइसे वहाँ इंस्टॉल करें जहाँ tasks चलते हैं, सिर्फ़ वहाँ नहीं जहाँ DAGs parse होते हैं
load के दौरान TimeoutExceptionमशीन जितने संभाल सकती है उससे ज़्यादा एक साथ चलने वाले हलtask को किसी pool में डालें और slot count पर सीमा लगाएँ
Form ऐसे token को अस्वीकार करता है जो ठीक हल हुआ थाtoken दो tasks के बीच पुराना पड़ गयाsubmit को हल करने वाले task में ले आएँ
ERROR_PAGEURLAPI तक कोई relative URL पहुँचाabsolute URL भेजें, scheme समेत

हर code और उसे ट्रिगर करने वाली वजह, सब कुछ यहाँ दिया गया है: CapSkip API डॉक्यूमेंटेशन.

FAQ

मेरे workers containers में चलते हैं। क्या वे सॉल्वर तक पहुँच सकते हैं?

हाँ, Server mode में। सॉल्वर loopback पते के बजाय किसी network पते पर सुनता है, और workers उसे किसी भी दूसरी internal service की तरह API के ज़रिए call करते हैं। host और port को worker environment variables के रूप में सेट करें ताकि DAG file में कोई पता रहे ही नहीं। container और laptop के बीच code में कुछ नहीं बदलता।

उस managed Airflow का क्या जिसे मैं ख़ुद नहीं चलाता?

जवाब वही है, बस एक अतिरिक्त शर्त के साथ। कोई hosted scheduler आपके network के बाहर चलता है, इसलिए सॉल्वर को ऐसा पता चाहिए जिस तक वह route कर सके, और एक static public IP ही उसे स्थिर बनाता है। key validation चालू करें और हर environment को उसकी अपनी key दें, ताकि कोई value लीक हो जाए तो बाक़ी को छुए बिना उसे रद्द किया जा सके।

क्या observability के लिए हल को अपना अलग task होना चाहिए?

यह लुभावना लगता है और इसकी क़ीमत आपकी reliability है। अलग task का मतलब है कि token एक XCom के रूप में सफ़र करता है और scheduler जब तक तय करता है कि आगे क्या चलेगा तब तक पुराना पड़ता रहता है, और tokens ठीक इसी तरह रास्ते में expire होते हैं। दोनों को साथ रखें और observability इसके बजाय logs तथा task duration से लें, जो बिना उस दौड़ के वही बात दिखा देते हैं।

क्या कई DAGs एक ही सॉल्वर instance साझा कर सकते हैं?

हाँ, और आमतौर पर यही सेटअप होता है। Server mode में चल रहा एक instance हर उस worker को सेवा देता है जो उस तक पहुँच सके, और बाँटने के लिए प्रति हल कोई लागत होती ही नहीं। DAGs के बीच कुल concurrency बाँधने के लिए किसी Airflow pool का इस्तेमाल करें, क्योंकि जो सीमा आपके लिए मायने रखती है वह यह है कि एक साथ कितने हल चल रहे हैं, यह नहीं कि कितनी pipelines को हल चाहिए।

संक्षेप में

सॉल्वर को ऐसी मशीन पर रखें जहाँ तक workers पहुँच सकें, host को environment में सेट करें, और हल तथा submit को एक ही task में रखें, किसी pool के पीछे दो retries के साथ। बाक़ी सब एक सामान्य DAG है। इसके crawl वाले पहलू के लिए देखें: वेब स्क्रैपिंग के लिए कैप्चा सॉल्वर पेज, और client surface के लिए देखें: Python कैप्चा सॉल्वर पेज। शेड्यूल किया गया काम ही वह जगह है जहाँ pricing model सामने आता है। हर घंटे चलने वाला ऐसा DAG जो हर run पर हल करता है, किसी metered vendor की तरफ़ से एक असली बिल बन जाता है, और यहीं असली फ़र्क़ है: कैप्चा सॉल्वर का पहले से आपके अपने hardware पर चलना उतना ही खर्चीला रहता है, चाहे DAG दिन में एक बार चले या मिनट में एक बार।