Как решать капчу в пайплайне Dagster (ops и assets)

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

Шаг с капчей в Dagster укладывается в один op с политикой повторов, а весь пайплайн вокруг него занимает около тридцати строк. Специфичная для Dagster часть, о которой не предупредит ни один другой оркестратор, касается того, что происходит со значением, которое возвращает op. Dagster сохраняет выходные значения op и asset через IO manager, а тот, что стоит по умолчанию, пишет их на диск через pickle. Токен капчи представляет собой одноразовый секрет со сроком жизни в две минуты, поэтому диск оказывается последним местом, где ему стоит осесть.

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

  • Dagster 1.9 или новее и Python 3.10 или новее, плюс CapSkip Python SDK.
  • Dagster с открытым исходным кодом или развёртывание Dagster+. Код в обоих случаях одинаковый.
  • URL страницы с защищённой формой и её sitekey.
  • Запущенный CapSkip в режиме Local, когда код и решатель работают на одной машине, или в режиме Server, когда это не так. Оба режима описаны в разделе Настройки подключения.
# pip install capskip
pip install dagster dagster-webserver capskip

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

Где на самом деле выполняется ваш код

Ответьте на этот вопрос до того, как выберете строку хоста, потому что от него зависит вся сетевая настройка. Dagster с открытым исходным кодом целиком работает на ваших собственных машинах, так что обсуждать тут нечего. У Dagster+ есть две формы, и для нашей задачи они противоположны.

Hybrid ведёт себя так, как вам и хотелось бы. Вы запускаете агент внутри собственной инфраструктуры, и он сам подключается наружу к control plane. У Dagster+ нет входа в вашу сеть, он не видит ваш код и не трогает ваши данные, поэтому запуск, начатый из облачного интерфейса, всё равно выполняется на вашем оборудовании. Поставьте агент на ту же машину, где стоит CapSkip, и op будет обращаться к 127.0.0.1:8080 ровно так же, как скрипт на вашем компьютере. Никаких туннелей и публичных адресов.

Исключение составляет Serverless. Там ваш код выполняется в среде Dagster, а не в вашей, и loopback перестаёт значить что-либо полезное: он указывает на контейнер Dagster, где никто не слушает. В такой конфигурации решателю нужен режим Server и адрес, до которого запуск сможет достучаться. Это не понижение уровня: тот же решатель на том же оборудовании, только с другим адресом привязки, и платы за каждое решение по-прежнему нет.

Шаг 1: заверните решатель в ресурс

Для всего внешнего в Dagster принят ресурс, а не клиент уровня модуля. Унаследуйтесь от ConfigurableResource, объявите поля подключения, и тогда в интерфейсе появится запись для них в launchpad, конфигурация будет проверена до старта запуска, а тесты смогут подменить всё это заглушкой. Ключ читайте из окружения, а не зашивайте в код.

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

Один метод покрывает reCAPTCHA v2, Invisible, Enterprise и v3, потому что варианты задаются именованными опциями, а не отдельными вызовами: invisible со значением 1, enterprise со значением 1 или version со значением v3 и параметром action. У Turnstile и GeeTest есть свои методы такой же формы. Полный список параметров есть в документации CapSkip API.

Шаг 2: один op с политикой повторов

Вызов решателя выполняется по сети, а сетевые вызовы падают. RetryPolicy в Dagster описывается декларативно и вешается на декоратор: максимальное число попыток, базовая задержка, кривая backoff и режим jitter. По умолчанию разумно взять экспоненциальный backoff с jitter в обе стороны, потому что он растягивает пачку одновременных сбоев во времени вместо того, чтобы возвращать их все разом.

# 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: код статуса, а не токен. Именно об этом весь следующий раздел.

Шаг 3: никогда не выпускайте токен за границу op

Вот специфичная для Dagster ловушка, и попасть в неё легко, потому что очевидная схема выглядит так: один op решает капчу, а второй отправляет форму. Dagster не передаёт значения между ops в памяти. Каждое выходное значение проходит через IO manager, а по умолчанию используется файловый, который сохраняет выходы как pickle-файлы на локальном диске. Поэтому токен, возвращённый из op, записывается в файл, читается следующим op и остаётся там лежать.

Плохо это сразу дважды. Действующий секрет записывается на диск, где его никто не удаляет. И между решением капчи и отправкой появляется поход в хранилище, а именно такой задержки не выдерживает срок жизни в две минуты. При повторном выполнении становится хуже: Dagster может загрузить сохранённый выход прошлого запуска вместо пересчёта, и тогда форме достанется токен, просроченный ещё вчера.

Решение не в хитром IO manager. Держите решение капчи и отправку внутри одного op, чтобы токен жил в локальной переменной и вообще не становился выходным значением. Всё, что выше и ниже по потоку, может остаться отдельными ops или assets. Объединяется только та пара, которой приходится делить токен.

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

Если вы работаете с assets, а не с ops, правило то же самое, с одной оговоркой: материализованный asset представляет собой постоянную запись с историей в интерфейсе, а одноразовому секрету там не место. Опишите решение капчи как op внутри asset на основе графа, а сам asset пусть материализует результат отправки.

Полный рабочий пример

Готовое задание. Получить страницу, вытащить из неё sitekey, затем решить капчу и отправить форму в одном op. Два ops, один ресурс, один объект Definitions.

# 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 откладывает получение значения до момента запуска и не запекает его в определение, поэтому ключ никогда не попадает ни в репозиторий, ни в сериализованный снимок.

Не повторяйте те сбои, которые всё равно не пройдут

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

# 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 не запущен или указан неверный хост, а TimeoutException означает, что решение не уложилось в recaptchaTimeout, который по умолчанию равен 300 секундам. Оба заслуживают повторной попытки. ValidationException означает неверный аргумент, а ApiException означает, что API отклонил запрос, и ни то, ни другое со второй попытки не исправится. Все четыре наследуются от общего базового класса CapSkipError, так что вы можете ловить его, если предпочитаете обрабатывать всё в одном месте.

Ограничение пачки через concurrency pool

Разверните партиционированное задание на двести URL, и Dagster с удовольствием попытается выполнить их все, а это большая нагрузка, чем вы хотели бы направить на один решатель. Управлять этим позволяют concurrency pools. Пометьте op именем pool, задайте лимит один раз, и всё, что выходит за предел, встанет в очередь вместо того, чтобы наваливаться сразу.

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

Аргумент pool уже стоит на op в примере выше. Лимиты действуют на все запуски сразу, а не на каждый по отдельности, и это именно то, что нужно, когда три расписания смотрят на одну машину. Если вам удобнее разворачивать веер внутри одного процесса, а не заводить по одному op на URL, в Python SDK есть настоящий клиент на asyncio, и такой подход разобран в руководстве по параллельному решению капчи.

Запуск решателя на другой машине

Контейнеры code location, агенты Kubernetes и запуски Serverless уводят ваш op с вашего компьютера. В коде не меняется ничего, кроме хоста, а ресурс и так читает его из конфигурации.

У CapSkip два режима подключения. Local слушает 127.0.0.1 и отвечает только этому устройству. Server слушает ваш сетевой или публичный IP, поэтому контейнер, виртуальная машина или запуск Serverless доберутся до той же машины с Windows по API. Статический публичный IP делает адрес постоянным. Это по-прежнему ваше оборудование, и оплата за решения не взимается в любом случае, поэтому стоимость напряжённого дня от режима не зависит.

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

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

Частые ошибки и что они означают

Что вы видитеПричинаИсправить
Форма отклоняет токен, который выглядит правильнымОн прошёл через IO manager между двумя ops и приехал устаревшимРешайте капчу и отправляйте форму внутри одного op, чтобы токен остался локальной переменной
Повторное выполнение отправляет токен из старого запускаDagster загрузил ранее сохранённый выход вместо пересчётаРешение то же. Токен никогда не должен быть выходом op или asset
В DAGSTER_HOME/storage появляется pickle-файлФайловый IO manager по умолчанию записал ваш токен на дискРешение то же, затем удалите файл. Считайте его утёкшим секретом
NetworkException в Dagster+ ServerlessЗапуск выполнился в среде Dagster, а не в вашейПереведите решатель в режим Server или используйте агент Hybrid
NetworkException на вашем собственном агентеCapSkip не запущен или указан неверный хостЗапустите приложение или направьте хост в ресурсе на адрес сервера
Три повторные попытки, сожжённые на неверном sitekeyКаждая попытка падает одинаково и детерминированноПоднимайте RetryRequested только для NetworkException и TimeoutException
Решатель перегружен во время backfillПартиции по умолчанию выполняются параллельноПоместите op в pool и задайте для него лимит
ValidationExceptionОтсутствующий или некорректный sitekey либо URL страницыВыведите оба значения в лог перед вызовом и проверьте, что sitekey действующий

FAQ

Должно ли решение капчи быть asset?

Нет. Asset представляет собой долговечный объект с историей материализаций, а токен, который истекает через две минуты, является полной его противоположностью. Опишите как asset то, что вы на самом деле произвели, будь то отправленная запись или страница, которую вам наконец разрешили прочитать, а решение капчи оставьте как op внутри него.

Может ли запуск Dagster+ действительно достучаться до 127.0.0.1?

На Hybrid да, потому что код выполняет агент в вашей собственной инфраструктуре. Loopback там означает машину этого агента, поэтому решатель на ней отвечает как обычно. На Serverless запуск происходит в среде Dagster, и вам нужен режим Server с доступным адресом.

Чем это отличается от того же самого в Airflow?

В основном тем, как перемещаются значения. Airflow прогоняет небольшие значения через XCom, и вы просто не кладёте туда токен. Dagster по умолчанию сохраняет каждый выход через IO manager, поэтому та же ошибка пишет pickle-файл вместо строки в базе. Со стороны решателя всё одинаково, а версия для Airflow описана в руководстве по капче в Airflow.

Блокирует ли долгое решение капчи остальную часть запуска?

Сам op заблокирован, пока опрашивает результат, но исполнитель multiprocess в Dagster продолжает выполнять другие ops, которые от него не зависят. Лимит pool ограничивает, сколько решений идёт одновременно, и при этом ничего больше не останавливает. На вашем собственном оборудовании единственной ценой медленного решения остаётся потраченное время.

Коротко

Заверните решатель в ConfigurableResource и читайте ключ из EnvVar. Дайте op политику RetryPolicy с экспоненциальным backoff и поднимайте RetryRequested самостоятельно, когда хотите пропустить повторы, которые всё равно не сработают. И самое главное, держите решение капчи и отправку в одном op, потому что Dagster по умолчанию сохраняет выходы ops, а токен сохранять не нужно.

Сторона Python во всём этом разобрана на странице сервиса распознавания капч для Python. Подробности по этому типу капчи находятся на странице сервиса распознавания reCAPTCHA v2. Те же три вызова есть и в Node.js, PHP и C#, и они перечислены на на странице SDK для распознавания капчи.

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