如何在 Prefect flow 中识别验证码并配置重试

Prefect 里的验证码环节就是一个带重试的 task,整个 flow 大约二十行。从其他托管平台过来的人会遇到一个意外,而且是个好意外:Prefect Cloud 从不运行你的代码。它只负责调度工作,由你自己机器上的一个 worker 通过向外的连接把工作取走。所以识别工具可以待在 127.0.0.1 上,flow 照样能访问到它,这在 Zapier 或 Make.com 上是做不到的。真正会出问题的只有两件事:重试行为和缓存,而且各自都只需要一个参数。
你需要什么
- Prefect 3 和 Python 3.10 或更高版本,外加 CapSkip Python SDK。
- 一个 Prefect Cloud 工作区,或者自托管的 Prefect server。就这件事而言两者用法完全一样。
- 受保护表单的页面 URL,以及它的 sitekey。
- 如果 worker 和识别工具在同一台机器上,就让 CapSkip 以本地模式运行;如果不在,就用服务器模式。两种模式都写在 连接设置.
# 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 存储元数据、协调运行,它不执行你的代码,也不需要任何进入你网络的入站访问。你在自己基础设施里启动的 worker 会向外轮询任务,并在本地提交运行。
这带来的实际后果值得直说。在与 CapSkip 同一台 Windows 机器上运行一个 Process worker,flow 调用 127.0.0.1:8080 的方式,和你桌上的脚本完全一样。不用隧道,不用公网地址,不用证书。这和纯云端自动化平台的情况正好相反,也是把识别工作放进编排器里让人省心的主要原因。
有一个例外,在选 work pool 之前你应该先知道。Prefect Managed work pool 会把你的 flow 跑在 Prefect 自己的基础设施上,而不是你的,这很方便,但也彻底去掉了回环地址这个选项。在那种配置下,识别工具需要服务器模式和一个可达的地址。Prefect 公布了 Managed 运行所用的六个固定出站地址,所以你可以在防火墙上只放行这几个到识别工具端口的访问,其余一律丢弃。Managed 运行还必须使用官方 Prefect 镜像,并且上限是 24 小时,这也是这里通常更适合用 Process 或 Docker work pool 的另一个原因。
第 1 步:把密钥放进 Secret block
不要把 API 密钥硬编码在 flow 文件里。Prefect 自带 Secret block,值在后端是静态加密的,加载它只要两行代码。在 Python shell 里保存一次就行。
# 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 接受固定延迟、一组延迟,或者一个指数退避辅助函数,再加上一个抖动系数把重试错开,这样一批失败不会齐刷刷地同时回来。
更重要的是另一个参数。验证码 token 是一次性的,几分钟内就会过期,所以绝不能从缓存里出来。Prefect 3 在没有开启结果持久化时不会缓存,也就是说大多数人是碰巧安全的。如果你的团队全局打开了持久化,而这样做的团队并不少,那么重试的 task 可能会把第一次产生的那个 token 再交出来,表单会拒绝它。把策略显式设定好,然后就别再想它了。
# 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,做识别,然后把 token 随表单提交回去。每一步都是一个 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())在提交之前紧挨着做识别,绝不要放在更早的某个调度步骤里。一个 token 在结果存储里躺十分钟等上游 task 跑完,等表单看到它时已经是个死 token 了。
批量识别又不把识别工具压垮
Prefect 默认通过线程池并发运行 task,所以一百次识别就是对 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 运行都生效,当三个调度都指向同一个识别工具时,这正是你想要的。如果你更想在单个进程里做扇出,Python SDK 的 AsyncCapSkip 是一个真正的 asyncio 客户端,那种做法写在 并行识别验证码的指南.
把识别工具跑在另一台机器上
worker 是会挪地方的。你桌上的一个 Process worker 会变成服务器上的 Docker worker,再变成 Kubernetes work pool,到某个时候 flow 就不在识别工具所在的机器上了。代码里除了主机地址之外什么都不用改。
CapSkip 有两种连接模式。本地模式绑定到 127.0.0.1,只有该设备本身能访问。服务器模式绑定到你的内网或公网 IP,这样虚拟机上的 worker、容器宿主机或者 Managed work pool,都能通过 API 调用同一台 Windows 机器。建议配一个固定公网 IP,好让这个地址保持稳定。不管用哪种模式,它都还是你自己的硬件,都还是不按量计费,所以忙碌一天的成本不会因为模式而变化。
# Same SDK, same call. Only the host moves. solver = CapSkip(host="10.0.0.12", port=8080, apiKey=key)
识别工具一旦监听在网络地址上,就把密钥校验打开,并且给每个 worker 发各自的密钥,这样吊销其中一个不会动到其他的。两种模式的完整说明见 CapSkip 设置指南.
常见错误及其含义
| 你所看到的 | 原因 | 修复 |
|---|---|---|
| 重试返回了同一个已过期的 token | 结果持久化开着,task 缓存了它的输出 | 在识别 task 上设置 cache_policy=NO_CACHE |
| 表单拒绝了一个看起来没问题的 token | 识别完成后过了好几分钟才提交 | 把识别放在紧挨着提交的那一步 |
| Managed work pool 上出现 NetworkException | flow 跑在 Prefect 的基础设施上,不是你的 | 把识别工具切到服务器模式,或者改用 Process worker |
| 自己的 worker 上出现 NetworkException | CapSkip 没在运行,或者 host 填错了 | 启动应用,或者把 host 指向服务器地址 |
| 三次重试白白耗在一个错误的 sitekey 上 | 每次尝试都以同样的确定性方式失败 | 加上 retry_condition_fn,遇到 ApiException 就快速失败 |
| 批量运行时识别工具过载 | task 默认并发运行 | 占用一个全局并发限制的槽位 |
| TimeoutException | 识别耗时超过了 recaptchaTimeout | 把它调到默认的 300 秒以上 |
| ValidationException | 参数缺失或格式错误 | 提交之前先检查 sitekey 和页面 URL |
常见问题
Prefect Cloud 的 flow 真的能调用 127.0.0.1 吗?
能,在混合式 work pool 上可以,因为代码跑在你的 worker 上,而不是 Prefect 的云端。那里的回环地址指的就是 worker 自己那台机器,所以只要 CapSkip 在那台机器上,调用就能成功。例外是 Managed work pool,那种情况下算力由 Prefect 提供,你需要服务器模式加上一个可达的地址。
识别应该单独做一个 task,还是并进一个更大的 task 里?
单独做一个 task。它是最容易出现瞬时失败的一步,它需要一套其他步骤并不想要的重试策略,而且单独拆出来意味着运行图能清楚告诉你识别有多少次是拖慢速度的那一环。把它紧挨着提交步骤放,这样用到 token 时它还是新鲜的。
这和在 Airflow 里做有什么不同?
主要是代码形态不同。Airflow 要的是一个 operator 和一个由你自己托管的调度器,重试配置挂在 task instance 上。Prefect 给你的是一个加了装饰器的函数,以及一个向外连接的 worker。识别工具这一侧两者完全相同,Airflow 版本写在 Airflow 验证码指南.
一次很慢的识别会算进我的 flow 运行时长里吗?
会,task 在轮询期间是阻塞的。在你自己的 worker 上这没关系,唯一的代价是墙上时钟时间。在 Managed work pool 上就有影响了,那里的算力按运行时长计费,而且上限是 24 小时。这也是把识别放在你已经拥有的硬件上的又一个理由。
简短版结论
把密钥存在 Secret block 里,把识别包进一个带重试和 NO_CACHE 的 task,并在紧挨着提交的那一步调用它。在 CapSkip 所在的机器上运行 worker,主机地址就一直是 127.0.0.1。在把批量任务扇出之前,先加一个全局并发限制。这一切的 Python 侧内容见 Python 验证码识别页面。同样的三个调用在 Node.js、PHP 和 C# 里也都有,它们列在 验证码识别 SDK 页面.
在把它设成每小时运行之前,有一件事值得知道。CapSkip 做的是 验证码绕过 ,跑在你已经拥有的硬件上,所以一个每天识别一万次的 flow,和一个每天只识别十次的 flow 花费相同。
