如何在 Dagster 流水线中识别验证码(op 与 asset)

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

在 Dagster 里,一个验证码步骤就是一个带重试策略的 op,围绕它的整条流水线大约三十行。真正属于 Dagster 特有、而且没有别的编排工具会提醒你的,是这个 op 返回值的去向。Dagster 会通过 IO manager 持久化 op 和 asset 的输出,而默认的那个会把它们 pickle 到磁盘上。一个验证码 token 是只能用一次、寿命两分钟的凭据,磁盘是它最不该出现的地方。

你需要什么

  • 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

你的代码究竟跑在哪里

在选定 host 字符串之前先回答这个问题,因为它决定了整套网络配置。Dagster 开源版完全跑在你自己的机器上,所以这里没什么可讨论的。Dagster+ 有两种形态,就这件事而言,它们正好相反。

Hybrid 是符合你期望的那一种。你在自己的基础设施里跑一个 agent,由它向外连到控制平面。Dagster+ 没有进入你网络的入口,看不到你的代码,也不接触你的数据,所以从云端界面发起的一次运行,依然执行在你自己的硬件上。把 agent 和 CapSkip 放在同一台机器上,op 调用 127.0.0.1:8080 的方式,就和你电脑上的一个脚本完全一样。不需要隧道,也不需要公网地址。

Serverless 是例外。在那里,你的代码执行在 Dagster 的环境里而不是你的环境里,回环地址也就不再有任何意义:它指向 Dagster 的容器,而那里没有任何东西在监听。在这种配置下,识别工具需要 Server 模式,以及一个这次运行能访问到的地址。这算不上降级:它还是同一台硬件上的同一个识别工具,只是换了个绑定地址,而且依然不按次计费。

第 1 步:把识别工具包成一个 resource

在 Dagster 里,凡是外部依赖,惯用的做法都是 resource,而不是一个模块级的客户端。继承 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,配上一条重试策略

一次识别调用就是一次网络调用,而网络调用会失败。Dagster 的 RetryPolicy 是声明式的,写在装饰器上:最大次数、基础延迟、退避曲线和抖动模式。指数退避加正负抖动是合理的默认选择,因为它会把一批同时发生的失败摊开,而不是让它们一起涌回来。

# 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 返回的是什么:一个状态码,而不是一个 token。这正是下一节要讲的全部重点。

第 3 步:绝不让 token 跨越 op 的边界

这是 Dagster 特有的坑,而且很容易踩进去,因为最自然的设计就是一个 op 负责识别、另一个 op 负责提交。Dagster 不会在内存里把值从一个 op 传给另一个。每一个输出都要经过 IO manager,默认的那个是文件系统实现,它会把输出以 pickle 文件的形式存在本地磁盘上。于是一个从 op 返回的 token 会被写进文件,被下一个 op 读回来,之后还一直留在那里。

这件事错了两次。它把一个还有效的凭据写到磁盘上,而没有任何东西会去清理它。它还在识别和提交之间插进了一次存储往返,而这正是两分钟的有效期承受不起的延迟。到了重跑的时候情况更糟:Dagster 可能会加载上一次运行存下来的输出,而不是重新计算,于是表单拿到的是一个昨天就过期的 token。

解决办法不是写一个聪明的 IO manager,而是把识别和提交放在同一个 op 里,让 token 只活在一个局部变量里,压根不成为任何输出。上游和下游的一切都可以继续保持为独立的 op 或 asset。只有必须共用 token 的那一对才需要合并。

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

如果你用的是 asset 而不是 op,同样的规则依然成立,只是要多加一句:一个被物化的 asset 是一条永久记录,在界面上还留有历史,而一个一次性的密文完全没有理由成为这种东西。把识别建模成 graph backed asset 内部的一个 op,让 asset 去物化提交的结果。

完整可运行示例

一个完整的 job。抓取页面、从中取出 sitekey,然后在一个 op 里完成识别和提交。两个 op、一个 resource、一个 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 会把取值推迟到运行时,而不是把值固化进定义里,所以密钥既不会出现在你的代码仓库里,也不会出现在序列化快照里。

不要重试那些永远不会通过的失败

声明式的策略会对所有失败一视同仁地重试,于是一个格式不对的参数,每次都会以完全相同的方式失败,却白白耗掉三次尝试和一分钟的退避。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 没在运行,或者 host 配错了;TimeoutException 表示这次识别超过了 recaptchaTimeout,它默认是 300 秒。这两种都值得再试一次。ValidationException 表示参数有问题,ApiException 表示 API 拒绝了这个请求,这两种再试一次也不会有任何改善。四者都派生自一个共同基类 CapSkipError,所以如果你更想在一个地方统一处理,捕获它就行。

用并发 pool 给批量任务限速

把一个分区 job 扇出到两百个 URL 上,Dagster 会毫不犹豫地把它们全部跑起来,而这个负载超过了你愿意压向单个识别工具的量。并发 pool 就是控制手段。给 op 标上一个 pool 名,把上限设一次,超出上限的部分就会排队,而不是一拥而上。

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

上面例子里的 op 已经带上了 pool 参数。上限是跨所有运行生效的,而不是每次运行各算各的,当三个调度都指向同一台机器时,这正是你想要的。如果你更愿意在一个进程内部做扇出,而不是每个 URL 一个 op,Python SDK 自带一个真正的 asyncio 客户端,这种做法写在 并行识别验证码的指南.

把识别工具跑在另一台机器上

code location 容器、Kubernetes agent 和 Serverless 运行,都会把你的 op 挪离你的电脑。代码里除了 host 什么都不用改,而 resource 本来就是从配置里读它的。

CapSkip 有两种连接模式。Local 绑定 127.0.0.1,只响应该设备。Server 绑定你的内网地址或公网 IP,于是容器、虚拟机或一次 Serverless 运行都能通过 API 访问同一台 Windows 机器。静态公网 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 设置指南.

常见错误及其含义

你所看到的原因修复
表单拒绝了一个看起来没问题的 token它在两个 op 之间经过了 IO manager,到手时已经放老了在一个 op 里完成识别和提交,让 token 始终只是一个局部变量
重跑时提交了一个来自旧运行的 tokenDagster 加载了之前存下来的输出,而不是重新计算同一个修复。token 绝不能成为 op 或 asset 的输出
DAGSTER_HOME/storage 下出现了一个 pickle 文件默认的文件系统 IO manager 把你的 token 写到了磁盘上同一个修复,然后把文件删掉。把它当作一次凭据泄露来处理
在 Dagster+ Serverless 上抛出 NetworkException这次运行执行在 Dagster 的环境里,而不是你的环境里把识别工具切到 Server 模式,或者改用 Hybrid agent
在你自己的 agent 上抛出 NetworkExceptionCapSkip 没在运行,或者 host 填错了启动应用,或者把 resource 的 host 指向服务器地址
三次重试白白耗在一个错误的 sitekey 上每次尝试都以同样的确定性方式失败只针对 NetworkException 和 TimeoutException 抛出 RetryRequested
回填期间识别工具过载分区默认是并发运行的把这个 op 放进一个 pool,并给它设上限
ValidationExceptionsitekey 或页面 URL 缺失,或者格式不对在调用前把两者都打印出来,并确认 sitekey 是线上正在用的那个

常见问题

识别应该做成一个 asset 吗?

不应该。asset 是一个有物化历史的持久对象,而一个两分钟就过期的 token 恰恰相反。把你真正产出的东西建模成 asset,无论那是提交后的记录,还是你终于被允许读到的那个页面,并把识别保留为它内部的一个 op。

Dagster+ 的运行真的能访问到 127.0.0.1 吗?

在 Hybrid 上可以,因为真正执行代码的是你自己基础设施里的那个 agent。那里的回环地址指的是 agent 所在的机器,所以跑在那台机器上的识别工具会正常响应。在 Serverless 上,运行发生在 Dagster 的环境里,你需要 Server 模式,以及一个能访问到的地址。

这和在 Airflow 里做有什么不同?

主要差别在值的流动方式上。Airflow 通过 XCom 传递小的值,你只要不把 token 放进去就行。Dagster 默认会通过 IO manager 持久化每一个输出,所以同样的错误在这里写出的是一个 pickle 文件,而不是一行数据库记录。识别工具那一侧完全一样,Airflow 的版本写在 Airflow 验证码指南.

一次很慢的识别会阻塞这次运行的其余部分吗?

这个 op 在轮询期间是被阻塞的,但 Dagster 的多进程执行器会继续跑那些不依赖它的其他 op。pool 的上限只限制同时进行的识别数量,不会让别的东西停下来。在你自己的硬件上,一次慢识别唯一的代价就是墙上时钟时间。

简短版结论

把识别工具包进一个 ConfigurableResource,并用 EnvVar 读取密钥。给 op 配一条带指数退避的 RetryPolicy,当你想跳过那些不可能成功的重试时,就自己抛出 RetryRequested。最重要的是,把识别和提交留在同一个 op 里,因为 Dagster 默认会持久化 op 的输出,而 token 不是该被持久化的东西。

这一切在 Python 那一侧的内容写在 Python 验证码识别页面。这种验证码类型的细节则在 reCAPTCHA v2 识别页面。同样的三个调用在 Node.js、PHP 和 C# 里也都有,它们列在 验证码识别 SDK 页面.

在把它排到一次大规模回填上之前,还有一件事要知道。CapSkip 是一个 验证码识别工具 运行在你已经拥有的硬件上,所以一个要识别五万行的 job,和一个只识别五十行的 job 花的钱一模一样。