Как решать капчу в DAG Apache Airflow (Python)

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

Шаг с капчей в Airflow представляет собой обычную задачу на Python. Вы вызываете решатель, получаете токен и используете его в той же задаче. Никакого оператора ставить не нужно, и плагин писать тоже. По-настоящему спотыкаются люди на другом, на расположении: ваши воркеры работают там, куда их поместил планировщик, а решатель, привязанный к адресу loopback на вашем ноутбуке, недоступен из контейнера на другом хосте. Сначала разберитесь с этим, и тогда DAG уложится в пятнадцать строк.

Что понадобится

  • Airflow 2.x или 3.x, с CapSkip SDK, установленным в том же образе или virtualenv, который используют ваши воркеры.
  • Цель, которая действительно отдаёт проверку. Страница за виджетом, а не тестовый ключ, который проходит всегда.
  • Запущенный CapSkip в режиме Server на машине, доступной вашим воркерам, или в режиме Local, если воркер и решатель находятся на одной машине. Оба режима описаны в разделе Настройки подключения.
# Install into the worker environment, not just the scheduler.
pip install capskip

Где должен находиться решатель

В этом и состоит вся проблема, поэтому она идёт первой. Задача выполняется внутри процесса-воркера, и в большинстве реальных развёртываний этот воркер представляет собой контейнер на другом хосте, чем всё, что вы администрируете руками. Адрес loopback внутри такого контейнера указывает на сам контейнер, поэтому SDK, направленный туда, ничего не находит.

РежимПрослушиваетКогда использовать
Локально127.0.0.1, только это устройствоОдин воркер на той же машине, что и решатель
СерверВаш сетевой адрес или публичный IPКонтейнеры, парк воркеров, VPS или управляемый Airflow

Режим Server подходит почти для любого развёртывания Airflow. Вы переключаете адрес прослушивания в приложении, направляете хост SDK на эту машину, и весь парк воркеров пользуется одним решателем. Статический публичный IP рекомендуется, когда вызывающие стороны находятся за пределами вашей сети. Сути продукта это не меняет: железо по-прежнему ваше и по-прежнему без тарификации, так что уход с адреса loopback меняет только место запуска. CapSkip представляет собой приложение для Windows, поэтому на практике речь идёт об одной Windows-машине, в которую стучится весь парк.

Шаг 1: настройте хост один раз

Не прописывайте адрес прямо в файле DAG. SDK читает CAPSKIP_HOST, CAPSKIP_PORT и CAPSKIP_API_KEY из окружения, и для Airflow это ложится лучше всего, потому что способ задавать переменные окружения воркерам у вас уже есть. 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,
    )

Создавайте клиент внутри задачи, а не на уровне модуля. Airflow разбирает каждый файл DAG по кругу, и всё, что создаётся во время импорта, создаётся при каждом разборе, причём и в планировщике, и в воркере.

Шаг 2: опишите решение одной задачей

Декораторы TaskFlow дают самый короткий путь. Airflow 3 импортирует их из модуля SDK, Airflow 2 из модуля decorators, а тело в обоих случаях одинаковое.

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

Это и есть вся интеграция. Вызов решателя синхронный, он сам опрашивает результат с задержкой, которая начинается с четверти секунды, и возвращает словарь, в поле code которого лежит токен.

Шаг 3: не передавайте токен между задачами

Вот ошибка, которую стоит назвать вслух, потому что Airflow делает её естественной. Возвраты TaskFlow превращаются в XCom, поэтому задача решения, возвращающая токен, и следующая задача, которая его потребляет, выглядят как хороший дизайн. На деле это гонка. Токен reCAPTCHA остаётся действительным примерно две минуты, а промежуток между двумя задачами Airflow решает планировщик, и вы им не управляете. Добавьте задержку в очереди, заполненный пул или перезапуск воркера, и токен протухнет по дороге.

Сбой получается плавающим, а это худший вид: DAG работает при тестировании и падает в нескольких процентах случаев на проде, с отклонённой формой и без единой ошибки от решателя. Отдельный пост про срок жизни токена reCAPTCHA разбирает тайминги. В DAG правило сводится к одной строке: решайте и отправляйте внутри одной задачи, а повтор пусть получает токен заново.

Полный 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 делает там реальную работу. Пул ограничивает, сколько экземпляров задачи выполняется одновременно во всём развёртывании, поэтому бэкфилл из двухсот запусков не откроет двести одновременных решений против одной машины. Создайте в интерфейсе пул с именем captcha, задайте ему нужное число параллельных решений, и каждая задача, которая его называет, встанет в очередь за этим лимитом.

Повторы, которые помогают, а не вредят

Задайте retries на задаче и позвольте всему решению вместе с отправкой повториться целиком. Поскольку токен запрашивается внутри тела задачи, повтор автоматически получает свежий, и именно такого поведения вы хотите, и именно поэтому два шага держатся вместе.

Дайте повтору задержку. Немедленный повтор к сайту, который только что выдал вам проверку, обычно получает проверку снова, а пара минут в запланированном пайплайне ничего не стоит. Двух повторов обычно хватает: третий сбой, как правило, означает неверный sitekey или недоступный решатель, а это ожиданием не лечится.

Частые ошибки

Что вы видитеПричинаИсправить
NetworkException на каждой задачеВоркер не может достучаться до решателяПерейдите в режим Server и задайте CAPSKIP_HOST на воркере
Работает в планировщике, падает на воркереВ образе воркера нет SDKСтавьте его там, где выполняются задачи, а не только там, где разбираются DAG
TimeoutException под нагрузкойОдновременных решений больше, чем тянет машинаПоместите задачу в пул и ограничьте число слотов
Форма отклоняет токен, который решился нормальноТокен состарился между двумя задачамиПеренесите отправку в задачу решения
ERROR_PAGEURLВ API пришёл относительный URLОтправляйте абсолютный URL, вместе со схемой

Полный список кодов и того, что их вызывает, находится в документации CapSkip API.

FAQ

Мои воркеры работают в контейнерах. Смогут ли они достучаться до решателя?

Да, в режиме Server. Решатель слушает сетевой адрес вместо адреса loopback, и воркеры обращаются к нему по API, как к любому другому внутреннему сервису. Задайте хост и порт переменными окружения воркера, чтобы в файле DAG адреса не было вовсе. В коде между контейнером и ноутбуком не меняется ничего.

А что с управляемым Airflow, который я не администрирую?

Ответ тот же, с одним дополнительным требованием. Хостовый планировщик работает вне вашей сети, поэтому решателю нужен адрес, до которого тот сможет достучаться, и стабильным это делает статический публичный IP. Включите проверку ключей и выдайте каждому окружению собственный ключ, чтобы утёкшее значение можно было отозвать, не трогая остальные.

Стоит ли выносить решение в отдельную задачу ради наблюдаемости?

Соблазнительно, но стоит вам надёжности. Отдельная задача означает, что токен едет через XCom и стареет, пока планировщик решает, что запускать дальше, а именно так токены и протухают по дороге. Держите их вместе, а наблюдаемость берите из логов и длительности задач: они показывают то же самое, но без гонки.

Могут ли несколько DAG использовать один экземпляр решателя?

Да, и так обычно и делают. Один экземпляр в режиме Server обслуживает каждый воркер, который до него дотягивается, и делить между ними нечего, потому что платы за решение нет. Используйте пул Airflow, чтобы ограничить общую параллельность между DAG, потому что вам важен лимит на число одновременных решений, а не на число пайплайнов, которым решение понадобилось.

Коротко

Поставьте решатель на машину, доступную воркерам, задайте хост в окружении и держите решение и отправку в одной задаче с парой повторов за пулом. Всё остальное представляет собой обычный DAG. Про сторону парсинга смотрите на странице сервиса распознавания капчи для парсинга, а описание клиентской части ищите на странице сервиса распознавания капч для Python. Регулярные задачи как раз и показывают, чем отличается модель оплаты. Ежечасный DAG, который решает капчу на каждом запуске, у поставщика с тарификацией превращается в реальный счёт, и в этом вся разница: распознавание капчи на железе, которым вы уже владеете, стоит одинаково, запускается ли DAG раз в сутки или раз в минуту.