如何在 Apache Airflow DAG 中识别验证码(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 不对,或者识别工具访问不到,这两种情况都不是靠等能解决的。
常见错误
| 你所看到的 | 原因 | 修复 |
|---|---|---|
| 每个任务都抛 NetworkException | worker 访问不到识别工具 | 切换到 Server 模式,并在 worker 上设置 CAPSKIP_HOST |
| 在调度器上正常,在 worker 上失败 | worker 镜像里没装 SDK | 要装在任务实际运行的地方,而不只是解析 DAG 的地方 |
| 负载上来后出现 TimeoutException | 同时进行的识别数超出了机器能承受的量 | 把任务放进一个 pool,并限制槽位数量 |
| 表单拒绝了一个识别本身没问题的 token | token 在两个任务之间放老了 | 把提交挪进做识别的那个任务里 |
| 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 是一天触发一次还是一分钟触发一次,成本都一样。
