Как решать капчу во flow Prefect с повторными попытками

Шаг с капчей в Prefect представляет собой одну task с повторными попытками, а весь flow укладывается примерно в двадцать строк. Деталь, которая удивляет тех, кто пришёл с других облачных платформ, приятная: Prefect Cloud никогда не выполняет ваш код. Он планирует работу, а забирает её воркер на вашей собственной машине через исходящее соединение. Поэтому решатель может стоять на 127.0.0.1, и flow до него дотянется, чего не скажешь про Zapier или Make.com. По-настоящему ломаются две вещи: поведение при повторах и кэширование, и каждая лечится одним аргументом.
Что понадобится
- Prefect 3 и Python 3.10 или новее, плюс CapSkip Python SDK.
- Рабочее пространство Prefect Cloud или собственный сервер Prefect. Для этой задачи оба работают одинаково.
- URL страницы с защищённой формой и её sitekey.
- Запущенный CapSkip в режиме Local, когда воркер и решатель находятся на одной машине, или в режиме Server, когда нет. Оба описаны в разделе Настройки подключения.
# pip install prefect pip install -U prefect capskip # Point the CLI at your workspace, then start a worker # on the machine that should run the flows. prefect cloud login
Где на самом деле выполняется ваш flow
Именно этот вопрос определяет всю вашу сетевую схему, поэтому ответьте на него первым. По умолчанию Prefect использует гибридную модель: слой оркестрации размещён у них, а слой выполнения ваш. Prefect Cloud хранит метаданные и координирует запуски, он не выполняет ваш код и не требует входящего доступа в вашу сеть. Воркер, который вы запускаете внутри собственной инфраструктуры, сам опрашивает наружу и запускает задачи локально.
Практическое следствие стоит проговорить прямо. Запустите Process worker на той же машине с Windows, где стоит CapSkip, и flow обратится к 127.0.0.1:8080 ровно так же, как это сделал бы скрипт на вашем столе. Ни туннеля, ни публичного адреса, ни сертификата. Это противоположность тому, что происходит на чисто облачной платформе автоматизации, и это главная причина, по которой оркестратор оказывается удобным местом для работы с капчей.
Есть одно исключение, и знать о нём стоит до того, как вы выберете work pool. Managed work pool в Prefect выполняет ваш flow на инфраструктуре Prefect, а не на вашей, что удобно, но полностью убирает вариант с loopback. В такой конфигурации решателю нужен режим Server и доступный адрес. Prefect публикует шесть статических исходящих адресов, которые используют запуски Managed, поэтому вы можете пропустить через фаервол к порту решателя ровно их, а всё остальное отбросить. Запуски Managed вдобавок обязаны использовать официальный образ Prefect и ограничены 24 часами, и это ещё одна причина, по которой work pool типа Process или Docker обычно подходит здесь лучше.
Шаг 1: положите ключ в Secret block
Не прописывайте API-ключ прямо в файле flow. В Prefect есть Secret block, значения хранятся в бэкенде в зашифрованном виде, а загрузка занимает две строки. Сохраните его один раз из оболочки Python.
# pip install prefect
from prefect.blocks.system import Secret
secret = Secret(value="YOUR_API_KEY")
secret.save("capskip-api-key")
# Rotating it later needs overwrite, or the save is refused.
# secret.save("capskip-api-key", overwrite=True)Шаг 2: напишите task для решения капчи
Одна task, одно решение. Задайте ей повторные попытки, потому что вызов решателя это сетевой вызов, а сетевые вызовы падают. Prefect принимает фиксированную задержку, список задержек или помощник для экспоненциального backoff, а коэффициент jitter разносит повторы во времени, чтобы пачка сбоев не вернулась строем.
Важнее другой аргумент. Токен капчи одноразовый и истекает за пару минут, поэтому он ни в коем случае не должен приходить из кэша. В Prefect 3 кэширование выключено, пока не включено сохранение результатов, так что большинство людей в безопасности случайно. Если ваша команда включила сохранение результатов глобально, а так делают многие, повторённая task может отдать тот же токен, который она выдала в первый раз, и форма его отклонит. Задайте политику явно и забудьте об этом.
# pip install capskip
from prefect import task
from prefect.tasks import exponential_backoff
from prefect.cache_policies import NO_CACHE
from prefect.blocks.system import Secret
from capskip import CapSkip
# NO_CACHE matters: a token is valid once and expires fast.
@task(
retries=3,
retry_delay_seconds=exponential_backoff(backoff_factor=5),
retry_jitter_factor=0.5,
cache_policy=NO_CACHE,
)
def solve_recaptcha(sitekey: str, page_url: str) -> str:
key = Secret.load("capskip-api-key").get()
solver = CapSkip(host="127.0.0.1", port=8080, apiKey=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.
Шаг 3: не повторяйте ошибки, которые никогда не пройдут
Слепые повторы тратят время на детерминированные сбои. Некорректный аргумент падает одинаково при каждой попытке, и то же самое делает sitekey, который не принадлежит странице. Функция условия повтора получает состояние и принимает решение, а возврат False немедленно завершает task с исходным исключением.
# Retry the transient ones. Fail fast on the rest.
from capskip import ValidationException, ApiException
def worth_retrying(task, task_run, state) -> bool:
try:
state.result()
except (ValidationException, ApiException):
return False # bad arguments or a bad sitekey
except Exception:
return True # solver down, or a timeout
return TrueПередайте её как retry_condition_fn у task. NetworkException означает, что CapSkip не запущен или указан не тот хост, а TimeoutException означает, что решение не уложилось в recaptchaTimeout, который по умолчанию равен 300 секундам. Обе ошибки действительно стоят ещё одной попытки. Эти два исключения, а также ValidationException и ApiException, происходят от общего базового, поэтому перехват CapSkipError работает, если вам удобнее обрабатывать сбои в одном месте.
Полный рабочий пример
Весь flow целиком. Загрузите страницу, вытащите из неё sitekey, решите капчу, затем отправьте токен обратно вместе с формой. Каждый шаг представляет собой отдельную task, поэтому у каждого свои повторные попытки, свои логи и своя запись в графе запуска.
# pip install prefect capskip httpx
import re
import httpx
from prefect import flow, task
from prefect.tasks import exponential_backoff
from prefect.cache_policies import NO_CACHE
from prefect.blocks.system import Secret
from capskip import CapSkip
PAGE_URL = "https://example.com/page-with-recaptcha"
@task(retries=2, retry_delay_seconds=5)
def read_sitekey(page_url: str) -> str:
html = httpx.get(page_url, timeout=30).text
match = re.search(r'data-sitekey=["\']([^"\']+)', html)
if not match:
raise RuntimeError("No data-sitekey on the page.")
return match.group(1)
@task(
retries=3,
retry_delay_seconds=exponential_backoff(backoff_factor=5),
cache_policy=NO_CACHE,
)
def solve_recaptcha(sitekey: str, page_url: str) -> str:
key = Secret.load("capskip-api-key").get()
solver = CapSkip(host="127.0.0.1", port=8080, apiKey=key)
return solver.recaptcha(sitekey=sitekey, url=page_url)["code"]
@task(retries=2, cache_policy=NO_CACHE)
def submit_form(page_url: str, token: str) -> int:
reply = httpx.post(
page_url,
data={"g-recaptcha-response": token},
timeout=30,
)
return reply.status_code
@flow(name="captcha-protected-submit")
def run():
sitekey = read_sitekey(PAGE_URL)
token = solve_recaptcha(sitekey, PAGE_URL)
return submit_form(PAGE_URL, token)
if __name__ == "__main__":
print(run())Решайте капчу непосредственно перед отправкой, а не в каком-то более раннем шаге по расписанию. Токен, который десять минут лежит в хранилище результатов, пока доделывается вышестоящая task, к моменту встречи с формой уже мёртв.
Как решить пачку, не завалив решатель
По умолчанию Prefect выполняет задачи параллельно через пул потоков, поэтому сотня решений это генератор списка поверх submit и вообще никакой настройки task runner. Параллельности здесь больше, чем вы, скорее всего, захотите направить на одну машину.
Рычагом управления служит глобальный лимит параллельности. Создайте лимит один раз через CLI, затем занимайте слот внутри task, и всё, что выходит за предел, будет ждать, а не наваливаться сверху.
# Create the limit once. Six solves in flight at a time. prefect gcl create capskip --limit 6
# The limit is enforced across every flow run, not per flow.
from prefect import flow, task
from prefect.cache_policies import NO_CACHE
from prefect.concurrency.sync import concurrency
from prefect.futures import wait
from capskip import CapSkip
@task(retries=3, cache_policy=NO_CACHE)
def solve_one(sitekey: str, page_url: str) -> str:
with concurrency("capskip", occupy=1):
solver = CapSkip(host="127.0.0.1", port=8080)
return solver.recaptcha(sitekey=sitekey, url=page_url)["code"]
@flow
def solve_many(sitekey: str, urls):
# submit, not map: map would iterate the sitekey string.
futures = [solve_one.submit(sitekey, u) for u in urls]
wait(futures)Лимит действует на все запуски flow в рабочем пространстве, а это именно то, что нужно, когда на один решатель нацелены три расписания. Если вам удобнее разворачивать веер внутри одного процесса, AsyncCapSkip из Python SDK представляет собой настоящий asyncio-клиент, и этот подход разобран в руководстве по параллельному решению капчи.
Запуск решателя на другой машине
Воркеры переезжают. Process worker на вашем столе становится Docker worker на сервере, затем Kubernetes work pool, и в какой-то момент flow оказывается уже не на той машине, где работает решатель. В коде не меняется ничего, кроме хоста.
У CapSkip два режима подключения. Local привязывается к 127.0.0.1 и отвечает только этому устройству. Server привязывается к вашему сетевому или публичному IP, поэтому воркер на виртуальной машине, хост контейнеров или Managed work pool обращаются к той же машине с Windows через API. Статический публичный IP сохраняет этот адрес неизменным. В обоих случаях это по-прежнему ваше оборудование и по-прежнему без тарификации, поэтому стоимость загруженного дня от режима не зависит.
# Same SDK, same call. Only the host moves. solver = CapSkip(host="10.0.0.12", port=8080, apiKey=key)
Как только решатель начинает слушать сетевой адрес, включите проверку ключей и выдайте каждому воркеру собственный ключ, чтобы один можно было отозвать, не трогая остальные. Оба режима подробно разобраны в разделе руководстве по настройке CapSkip.
Частые ошибки и что они означают
| Что вы видите | Причина | Исправить |
|---|---|---|
| Повтор возвращает тот же просроченный токен | Включено сохранение результатов, поэтому task закэшировала свой вывод | Задайте cache_policy=NO_CACHE у task с решением капчи |
| Форма отклоняет токен, который выглядит нормальным | Капча была решена за несколько минут до отправки | Решайте на шаге непосредственно перед отправкой |
| NetworkException на Managed work pool | Flow выполнился на инфраструктуре Prefect, а не вашей | Переключите решатель в режим Server или используйте Process worker |
| NetworkException на вашем собственном воркере | CapSkip не запущен или указан неверный хост | Запустите приложение или направьте host на адрес сервера |
| Три повторные попытки, сожжённые на неверном sitekey | Каждая попытка падает одинаково и детерминированно | Добавьте retry_condition_fn и завершайте сразу на ApiException |
| Решатель перегружен во время пакетной обработки | Задачи по умолчанию выполняются параллельно | Занимайте слот в глобальном лимите параллельности |
| TimeoutException | Решение заняло больше, чем recaptchaTimeout | Поднимите его выше значения по умолчанию в 300 секунд |
| ValidationException | Отсутствующий или некорректный аргумент | Проверьте sitekey и URL страницы перед отправкой |
FAQ
Может ли flow в Prefect Cloud действительно обращаться к 127.0.0.1?
Да, на гибридном work pool, потому что код выполняется на вашем воркере, а не в облаке Prefect. Loopback там означает машину самого воркера, поэтому, если CapSkip стоит на этой машине, вызов проходит. Исключение составляет Managed work pool, где вычислительные ресурсы даёт Prefect и вам нужен режим Server с доступным адресом.
Должно ли решение капчи быть отдельной task или частью более крупной?
Отдельной task. Это шаг, который чаще других падает по временным причинам, ему нужна политика повторов, не нужная остальным шагам, а отдельное существование означает, что граф запуска покажет вам, как часто именно решение оказывается медленным местом. Держите его рядом с шагом отправки, чтобы токен был свежим в момент использования.
Чем это отличается от того же самого в Airflow?
В основном формой кода. Airflow хочет оператор и планировщик, который вы размещаете сами, а настройки повторов живут на экземпляре задачи. Prefect даёт вам декорированную функцию и воркер, который подключается наружу. Со стороны решателя всё одинаково, а версия для Airflow описана в руководстве по капче в Airflow.
Учитывается ли долгое решение капчи во времени выполнения моего flow?
Да, task заблокирована, пока идёт опрос. На вашем собственном воркере это не страшно, там единственная цена это потраченное время. На Managed work pool это уже важно, потому что вычисления тарифицируются по длительности запуска и ограничены 24 часами. Ещё одна причина держать решение капчи на оборудовании, которое у вас уже есть.
Коротко
Храните ключ в Secret block, заверните решение капчи в task с повторными попытками и NO_CACHE и вызывайте её на шаге прямо перед отправкой. Запустите воркер на той машине, где стоит CapSkip, и хост останется 127.0.0.1. Перед тем как разворачивать пачку веером, добавьте глобальный лимит параллельности. Сторона Python во всём этом разобрана на странице сервиса распознавания капч для Python. Те же три вызова есть и в Node.js, PHP и C#, и они перечислены на на странице SDK для распознавания капчи.
Одну вещь стоит знать, прежде чем ставить это в расписание на каждый час. CapSkip умеет обход капчи на оборудовании, которое у вас уже есть, поэтому flow, который решает десять тысяч капч в день, стоит столько же, сколько тот, что решает десять.
