Cara Memecahkan CAPTCHA di DAG Apache Airflow (Python)

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

Langkah captcha di Airflow adalah task Python biasa. Anda memanggil solver, Anda menerima token, lalu Anda memakainya di task yang sama. Tidak ada operator yang perlu diinstal dan tidak ada plugin yang perlu ditulis. Yang sebenarnya menjegal orang adalah lokasi: worker Anda berjalan di mana pun scheduler menempatkannya, dan solver yang terikat pada alamat loopback laptop Anda tidak dapat dijangkau dari container di host lain. Benahi bagian itu lebih dulu, setelah itu DAG-nya hanya lima belas baris.

Apa yang Anda butuhkan

  • Airflow 2.x atau 3.x, dengan SDK CapSkip terinstal di image atau virtualenv yang sama dengan yang dipakai worker Anda.
  • Target yang benar-benar mengembalikan challenge. Halaman yang dilindungi widget, bukan kunci uji yang selalu lolos.
  • CapSkip berjalan dalam mode Server di mesin yang dapat dijangkau worker Anda, atau dalam mode Local jika worker dan solver berada di mesin yang sama. Keduanya dijelaskan di pengaturan koneksi.
# Install into the worker environment, not just the scheduler.
pip install capskip

Di mana solver harus ditempatkan

Inilah inti seluruh masalahnya, jadi bagian ini didahulukan. Sebuah task berjalan di dalam proses worker, dan pada sebagian besar deployment nyata worker itu adalah container di host yang berbeda dari apa pun yang Anda kelola secara manual. Alamat loopback di dalam container itu adalah container itu sendiri, jadi mengarahkan SDK ke sana tidak menemukan apa-apa.

ModeMendengarkan diGunakan saat
Lokal127.0.0.1, hanya perangkat ituSatu worker tunggal di mesin yang sama dengan solver
ServerAlamat jaringan atau IP publik AndaContainer, armada worker, sebuah VPS, atau Airflow terkelola

Mode Server adalah jawaban untuk hampir setiap deployment Airflow. Anda mengubah alamat listen di aplikasinya, mengarahkan host SDK ke mesin tersebut, dan setiap worker dalam armada berbagi satu solver. IP publik statis disarankan ketika pemanggilnya berada di luar jaringan Anda sendiri. Ini tidak mengubah hakikat produknya: perangkat kerasnya tetap milik Anda dan tetap tanpa kuota, jadi memindahkannya dari alamat loopback hanya memindahkan tempat ia berjalan, tidak lebih. CapSkip adalah aplikasi Windows, jadi dalam praktiknya itu berarti satu mesin Windows yang dipanggil oleh seluruh armada.

Langkah 1: konfigurasikan host satu kali

Jangan menuliskan alamatnya secara permanen di file DAG. SDK membaca CAPSKIP_HOST, CAPSKIP_PORT, dan CAPSKIP_API_KEY dari environment, dan itulah yang paling rapi untuk Airflow karena Anda sudah punya cara menyetel variabel lingkungan worker. Airflow Variable juga bisa dipakai jika Anda lebih suka menyimpannya di database metadata.

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

Bangun client di dalam task, bukan di lingkup modul. Airflow mengurai setiap file DAG secara berulang, dan apa pun yang dibuat pada saat import akan dibuat pada setiap penguraian, di scheduler maupun di worker.

Langkah 2: tulis pemecahannya sebagai satu task

Decorator TaskFlow adalah jalan masuk paling singkat. Airflow 3 mengimpornya dari modul SDK, Airflow 2 dari modul decorators, dan isi fungsinya identik pada keduanya.

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

Itulah keseluruhan integrasinya. Panggilan solver bersifat sinkron, ia melakukan polling untuk Anda dengan backoff yang dimulai dari seperempat detik, dan ia mengembalikan sebuah dictionary yang kolom code-nya berisi token.

Langkah 3: jangan mengoper token antar task

Inilah kesalahan yang layak disebut namanya, karena Airflow membuatnya terasa wajar. Nilai kembalian TaskFlow menjadi XCom, jadi satu task pemecahan yang mengembalikan token dan satu task berikutnya yang memakainya tampak seperti desain yang bagus. Padahal itu sebuah balapan. Token reCAPTCHA tetap valid kira-kira dua menit, sedangkan jeda antara dua task Airflow adalah keputusan scheduler yang tidak Anda kendalikan. Tambahkan penundaan antrean, pool yang penuh, atau worker yang restart, dan token itu kedaluwarsa di tengah jalan.

Kegagalannya bersifat sesekali, dan itu jenis yang paling buruk: DAG berhasil saat pengujian lalu gagal beberapa persen dari waktu di produksi, dengan formulir yang ditolak dan tanpa error dari solver. Artikel tentang berapa lama token reCAPTCHA tetap valid membahas rincian waktunya. Di dalam DAG aturannya menyusut jadi satu baris: pecahkan dan kirim di dalam task yang sama, lalu biarkan percobaan ulang memecahkannya lagi.

DAG lengkapnya

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

Argumen pool di situ benar-benar bekerja. Sebuah pool membatasi berapa banyak instance task yang berjalan sekaligus di seluruh deployment, sehingga backfill dua ratus run tidak membuka dua ratus pemecahan serentak terhadap satu mesin. Buat pool bernama captcha di UI, beri jumlah pemecahan paralel yang Anda inginkan, dan setiap task yang menyebut nama itu akan mengantre di balik batas tersebut.

Percobaan ulang yang menolong, bukan merugikan

Setel retries pada task dan biarkan seluruh proses pemecahan dan pengiriman diulang. Karena token diambil di dalam badan task, percobaan ulang otomatis mendapat token baru, dan itulah perilaku yang Anda inginkan sekaligus alasan kedua langkah itu tetap menyatu.

Beri jeda pada percobaan ulang. Percobaan ulang seketika terhadap situs yang baru saja menantang Anda cenderung ditantang lagi, dan menunggu beberapa menit tidak merugikan apa pun dalam pipeline terjadwal. Dua kali percobaan ulang biasanya cukup: kegagalan ketiga umumnya berarti sitekey yang salah atau solver yang tidak terjangkau, dan keduanya tidak akan beres hanya dengan menunggu.

Kesalahan umum

Apa yang Anda lihatPenyebabPerbaiki
NetworkException di setiap taskWorker tidak dapat menjangkau solverBeralihlah ke mode Server dan setel CAPSKIP_HOST di worker
Berhasil di scheduler, gagal di workerSDK tidak ada di image workerInstal di tempat task berjalan, bukan hanya di tempat DAG diurai
TimeoutException saat beban tinggiPemecahan bersamaan lebih banyak daripada yang sanggup ditangani mesinMasukkan task ke dalam pool dan batasi jumlah slotnya
Formulir menolak token yang sebenarnya terpecahkan dengan baikToken menua di antara dua taskPindahkan pengiriman ke dalam task pemecahan
ERROR_PAGEURLURL relatif sampai ke APIKirim URL absolutnya, termasuk skemanya

Daftar lengkap kode beserta pemicunya ada di dokumentasi API CapSkip.

FAQ

Worker saya berjalan di dalam container. Bisakah mereka menjangkau solver?

Bisa, dalam mode Server. Solver mendengarkan di alamat jaringan alih-alih alamat loopback, dan worker memanggilnya lewat API seperti layanan internal lainnya. Setel host dan port sebagai variabel lingkungan worker supaya file DAG tidak memuat alamat apa pun. Tidak ada bagian kode yang berbeda antara container dan laptop.

Bagaimana dengan Airflow terkelola yang tidak saya kelola sendiri?

Jawabannya sama, dengan satu syarat tambahan. Scheduler yang dihosting berjalan di luar jaringan Anda, jadi solver membutuhkan alamat yang bisa dirutekan ke sana, dan IP publik statis itulah yang membuatnya stabil. Aktifkan validasi kunci dan beri setiap lingkungan kuncinya sendiri, supaya nilai yang bocor bisa dicabut tanpa mengusik yang lain.

Haruskah pemecahan dijadikan task tersendiri demi observabilitas?

Itu menggoda, tetapi mengorbankan keandalan. Task terpisah berarti token berjalan sebagai XCom dan menua sementara scheduler memutuskan apa yang berjalan berikutnya, dan persis begitulah token kedaluwarsa di tengah jalan. Satukan keduanya, lalu dapatkan observabilitas dari log dan durasi task, yang menunjukkan hal yang sama tanpa balapan itu.

Bisakah beberapa DAG berbagi satu instance solver?

Bisa, dan itulah penyiapan yang lazim. Satu instance dalam mode Server melayani setiap worker yang dapat menjangkaunya, dan tidak ada biaya per pemecahan yang perlu dibagi. Gunakan pool Airflow untuk membatasi total konkurensi lintas DAG, karena batas yang Anda pedulikan adalah berapa banyak pemecahan yang berjalan sekaligus, bukan berapa banyak pipeline yang kebetulan membutuhkannya.

Versi singkatnya

Tempatkan solver di mesin yang dapat dijangkau worker, setel host-nya di environment, dan pertahankan pemecahan serta pengiriman dalam satu task dengan beberapa percobaan ulang di balik sebuah pool. Selebihnya hanyalah DAG biasa. Untuk sisi crawl-nya lihat halaman pemecah CAPTCHA untuk web scraping, dan untuk permukaan client-nya lihat halaman pemecah CAPTCHA Python. Pekerjaan terjadwal adalah tempat model harga benar-benar terasa. DAG per jam yang memecahkan pada setiap run berarti tagihan nyata dari vendor berbasis kuota, dan di situlah perbedaannya: pemecah captcha berjalan di perangkat keras yang sudah Anda miliki memakan biaya yang sama entah DAG dipicu sekali sehari atau sekali semenit.