如何在 Apache Airflow DAG 中识别验证码(Python)

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

在 Airflow 里,一个验证码步骤就是一个普通的 Python 任务。你调用识别工具,拿回一个 token,在同一个任务里把它用掉。没有 operator 要装,也没有插件要写。真正把人绊倒的是位置问题:你的 worker 跑在调度器把它放到的任何地方,而一个绑定在你笔记本回环地址上的识别工具,另一台主机上的容器根本访问不到。先把这一块弄对,剩下的 DAG 也就十五行。

你需要什么

  • Airflow 2.x 或 3.x,并且 CapSkip SDK 已装在你的 worker 所使用的那个镜像或虚拟环境里。
  • 一个真的会返回挑战的目标:一个挡在小组件后面的页面,而不是一个永远都能通过的测试密钥。
  • CapSkip 以 Server 模式运行在你的 worker 能访问到的机器上;如果 worker 和识别工具就是同一台机器,也可以用 Local 模式。两者都记录在 连接设置.
# Install into the worker environment, not just the scheduler.
pip install capskip

识别工具必须放在哪里

这就是全部问题所在,所以放在最前面讲。任务跑在一个 worker 进程内部,而在大多数真实部署里,这个 worker 是一个容器,和你手动管理的任何机器都不在同一台主机上。那个容器里的回环地址指的就是容器自己,所以把 SDK 指向它,什么也找不到。

模式监听地址适用场景
本地127.0.0.1,仅限该设备只有一个 worker,并且和识别工具在同一台机器上
服务器你的内网地址或公网 IP容器、worker 集群、VPS 或者托管的 Airflow

对几乎所有 Airflow 部署来说,答案都是 Server 模式。你在应用里换掉监听地址,把 SDK 的 host 指向那台机器,集群里的每个 worker 就共用同一个识别工具。如果调用方位于你自己网络之外,建议使用静态公网 IP。这并不改变产品本身是什么:它依然是你自己的硬件,依然不按量计费,所以把它从回环地址上挪开,挪动的只是它运行的位置,别的什么都没变。CapSkip 是 Windows 应用,所以实际落地就是整个集群去调用一台 Windows 机器。

第 1 步:把 host 配置一次

不要把地址硬编码进 DAG 文件。SDK 会从环境里读取 CAPSKIP_HOST、CAPSKIP_PORT 和 CAPSKIP_API_KEY,这对 Airflow 来说是最干净的方式,因为你本来就有办法给 worker 设置环境变量。如果你更愿意把它放在元数据库里,用 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 文件,凡是在导入时创建的东西,每解析一次就会被创建一次,调度器上和 worker 上都一样。

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

整个集成就这么多。识别调用是同步的,它会替你轮询,退避从四分之一秒起步,最后返回一个字典,token 就放在它的 code 字段里。

第 3 步:不要在任务之间传递 token

这个错误值得点名,因为 Airflow 会让它显得非常自然。TaskFlow 的返回值会变成 XCom,于是一个识别任务返回 token、一个下游任务消费它,看起来像是很好的设计。但它是一场竞速。一个 reCAPTCHA token 大约只有两分钟有效期,而两个 Airflow 任务之间的间隔是调度器的决定,你控制不了。再加上一次排队延迟、一个已经占满的 pool,或者一次 worker 重启,token 就在半路上过期了。

这种失败是间歇性的,也是最难缠的一种:DAG 在测试里好好的,到了生产环境有百分之几的概率失败,表现为表单被拒,而识别工具那边一个错误都没有。 关于 reCAPTCHA token 有效期有多长的那篇文章 讲清了这些时间。在 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 参数是真的在干活。pool 限制的是整个部署里同时运行的任务实例数量,这样一次两百个运行的回填,就不会对着同一台机器同时发起两百次识别。在界面里创建一个名为 captcha 的 pool,给它你想要的并行识别数量,之后每个指名用它的任务都会排在这个上限后面。

帮忙而不是帮倒忙的重试

在任务上设置 retries,让识别和提交整体重来一遍。由于 token 是在任务体内部取的,重试会自动拿到一个新的,这正是你想要的行为,也是这两个步骤要待在一起的原因。

给重试留一点延迟。对着一个刚刚挑战过你的站点立刻重试,通常会再被挑战一次,而在一条定时流水线里,等上几分钟什么代价都没有。两次重试通常就够了:第三次还失败,一般是 sitekey 不对,或者识别工具访问不到,这两种情况都不是靠等能解决的。

常见错误

你所看到的原因修复
每个任务都抛 NetworkExceptionworker 访问不到识别工具切换到 Server 模式,并在 worker 上设置 CAPSKIP_HOST
在调度器上正常,在 worker 上失败worker 镜像里没装 SDK要装在任务实际运行的地方,而不只是解析 DAG 的地方
负载上来后出现 TimeoutException同时进行的识别数超出了机器能承受的量把任务放进一个 pool,并限制槽位数量
表单拒绝了一个识别本身没问题的 tokentoken 在两个任务之间放老了把提交挪进做识别的那个任务里
ERROR_PAGEURL一个相对 URL 被送到了 API发送绝对 URL,带上协议头

所有错误码及其触发条件的完整清单见 CapSkip API 文档.

常见问题

我的 worker 跑在容器里,它们能访问到识别工具吗?

可以,用 Server 模式。识别工具监听的是一个网络地址而不是回环地址,worker 就像调用任何其他内部服务那样通过 API 调它。把 host 和端口设成 worker 的环境变量,这样 DAG 文件里就不会出现任何地址。从容器换到笔记本,代码一个字都不用改。

如果是我自己管不到的托管 Airflow 呢?

答案一样,只是多一个要求。托管的调度器跑在你的网络之外,所以识别工具需要一个它能路由到的地址,而静态公网 IP 才能让这件事稳定下来。打开密钥校验,给每个环境配一个自己的密钥,这样某个值泄露了,撤销它也不会波及其他环境。

为了可观测性,识别应该单独作为一个任务吗?

这么做很诱人,代价却是可靠性。单独一个任务意味着 token 要以 XCom 的形式传递,并且在调度器决定下一步跑什么的时候一直在变老,token 就是这样在半路上过期的。把它们放在一起,可观测性改从日志和任务耗时里拿,这些同样能看出问题,还没有竞速风险。

多个 DAG 能共用一个识别实例吗?

可以,而且这就是常见的做法。一个以 Server 模式运行的实例服务所有能访问到它的 worker,也没有什么按次成本需要分摊。用 Airflow 的 pool 来限定跨 DAG 的总并发,因为你真正在意的上限是同时有多少次识别在跑,而不是恰好有多少条流水线想要识别。

简短版结论

把识别工具放在 worker 能访问到的机器上,把 host 设在环境变量里,并让识别和提交待在同一个任务里,配上两次重试,再排在一个 pool 后面。其余部分就是一个普通的 DAG。想看爬虫那一侧的做法,请看 面向网络爬虫的验证码识别页面;至于客户端的接口面,请看 Python 验证码识别页面。定时任务正是定价模式显形的地方。一个每小时跑一次、每次都要识别的 DAG,在按量计费的供应商那里就是一张实打实的账单,而区别就在这里: 验证码识别工具 运行在你已经拥有的硬件上,无论 DAG 是一天触发一次还是一分钟触发一次,成本都一样。