如何在 Temporal 工作流的 Activity 中识别验证码

在 Temporal 里,验证码识别要放进 Activity。不是放在 workflow 方法里,也不是放在 workflow 调用的辅助函数里,就是放进 Activity。workflow 代码在每次工作流恢复时都会从历史记录中重放,所以它必须是确定性的:不能发网络请求,不能有随机性,不能读时钟。而一次识别同时占了这三条。放进 Activity 之后,剩下的就是一份重试策略和一点时序上的纪律,整件事大约四十行代码。
你需要什么
- Python 3.10 或更新版本、Temporal Python SDK 和 CapSkip Python SDK。
- 一个可连接的 Temporal Service。本地开发服务器或 Temporal Cloud 都可以。
- 受保护表单的页面 URL,以及它的 sitekey。
- worker 和识别工具在同一台机器上时,CapSkip 用 Local mode(本地模式),不在同一台时用 Server mode(服务器模式)。两种模式都说明在 连接设置.
# pip install temporalio pip install -U temporalio capskip httpx # A local service to develop against. temporal server start-dev
为什么识别不能放在 workflow 代码里
在 worker 重启、部署或睡了一周之后,Temporal 会重放工作流的历史记录来重建状态。要让两次重放得出同样的结果,workflow 代码必须是确定性的。SDK 明确列出了被排除的东西:不能有网络 IO,不能用线程,不能有随机性,不能调用外部进程,不能修改全局状态。它甚至把 workflow 代码放在一个沙箱里运行,每次运行都重新导入模块,这也是 activity 的导入语句要包在一个直通块里的原因。
一次验证码识别把这条规则违反了三遍。它是网络请求,每次返回的 token 都不一样,耗时还取决于机器。把它放进 workflow 方法,开发时看起来一切正常,但 worker 第一次在运行中途重启时就会报非确定性错误。
好消息是这条约束也给了你东西。Activity 由 Temporal 自己重试,用的是你声明的策略,而不是你手写的循环,而且它的结果会被记进历史。所以一次成功的识别在重放时不会再做一遍,这正是你对一次性答案所期望的行为。
第 1 步:识别 activity
CapSkip 的 Python SDK 提供的是真正的 asyncio 客户端,不是一层别名,所以 async activity 是最自然的选择,也不需要线程池。把 sitekey 和页面 URL 作为参数传入,返回 token。
# pip install capskip
import os
from temporalio import activity
from capskip import AsyncCapSkip
@activity.defn
async def solve_recaptcha(sitekey: str, page_url: str) -> str:
solver = AsyncCapSkip(
host=os.environ.get("CAPSKIP_HOST", "127.0.0.1"),
port=8080,
apiKey=os.environ.get("CAPSKIP_API_KEY", "capskip"),
)
result = await solver.recaptcha(sitekey=sitekey, url=page_url)
return result["code"]一个方法就覆盖了 reCAPTCHA v2、Invisible、Enterprise 和 v3。这些变体是同一个调用上的选项,而不是不同的方法:invisible 设为 1、enterprise 设为 1,或者 version 设为 v3 并带一个 action。Turnstile 和极验(GeeTest)有各自的方法,形态一样,完整的参数列表见 CapSkip API 文档.
如果你更想用同步客户端,那么这个 activity 必须写成普通的 def,worker 还需要配置 activity_executor,因为 Temporal 会在线程池里跑同步 activity。上面的异步版本完全绕开了这一点,这也是 Python SDK 确实比其他语言更好用的少数几处之一。
第 2 步:超时要比识别工具自身的更长
每个 activity 都需要 start_to_close_timeout,人们往往在这里悄无声息地毁掉自己的识别。CapSkip 对 reCAPTCHA、Turnstile 和极验(GeeTest)最多轮询 300 秒,对图片验证码最多 120 秒。把 activity 超时设得比这更短,Temporal 就会在识别工具还在干活时取消这次尝试,然后重试,于是一个表单就有两次识别同时在跑。
给 activity 留出比识别工具自身上限更多的余量。识别工具上限是五分钟时,activity 设六分钟比较宽裕。
# Longer than recaptchaTimeout, which defaults to 300s.
from datetime import timedelta
from temporalio.common import RetryPolicy
SOLVE_TIMEOUT = timedelta(minutes=6)
SOLVE_RETRIES = RetryPolicy(
initial_interval=timedelta(seconds=5),
backoff_coefficient=2.0,
maximum_attempts=4,
# These fail identically every time. Do not burn attempts.
non_retryable_error_types=["ValidationException", "ApiException"],
)不可重试列表是按异常类名匹配的,值得认真填写。ValidationException 表示参数缺失或格式不对,ApiException 表示 API 拒绝了请求,通常是 sitekey 与页面 URL 不匹配。这两种再试一次也不会变好。真正值得重试的是 NetworkException 和 TimeoutException:前者说明识别工具没在运行或 host 填错了,后者说明识别耗时超过了轮询超时。这四个异常都继承自同一个基类,所以如果你更愿意在一个地方统一处理,捕获 CapSkipError 就够了。
第 3 步:最后再识别,不要一开始就识别
持久化执行让 token 过期这件事在这里比在任何其他编排工具里都更容易出错。Temporal 工作流可以等待一个信号、睡上一天再继续,而它记录下来的 activity 结果会原封不动地从历史里取回。所以一个先识别、再等待审批、最后才提交的工作流,重放出来的会是昨天生成的 token。
reCAPTCHA token 只能用一次,大约两分钟后过期。安排工作流的顺序,让识别紧挨在提交之前,中间不要有任何可能阻塞的步骤。这个问题的整体形态见指南 reCAPTCHA token 过期.
完整可运行示例
三个 activity,加上一个按顺序调用它们的 workflow。读取 sitekey,识别,提交。每个 activity 有自己的超时和自己的重试策略,并且在 Temporal UI 里各自单独显示,所以哪一步慢了一眼就能看出来。
# activities.py
import os, re, httpx
from temporalio import activity
from capskip import AsyncCapSkip
@activity.defn
async def read_sitekey(page_url: str) -> str:
async with httpx.AsyncClient(timeout=30) as client:
html = (await client.get(page_url)).text
found = re.search(r'data-sitekey=["\']([^"\']+)', html)
if not found:
raise RuntimeError("No data-sitekey on the page.")
return found.group(1)
@activity.defn
async def solve_recaptcha(sitekey: str, page_url: str) -> str:
solver = AsyncCapSkip(host=os.environ.get("CAPSKIP_HOST", "127.0.0.1"))
return (await solver.recaptcha(sitekey=sitekey, url=page_url))["code"]
@activity.defn
async def submit_form(page_url: str, token: str) -> int:
async with httpx.AsyncClient(timeout=30) as client:
reply = await client.post(
page_url, data={"g-recaptcha-response": token}
)
return reply.status_codeworkflow 本身除了顺序之外不含任何逻辑。这正是重点:所有可能失败的东西都在 activity 里,而 workflow 是能扛过重放的那个确定性部分。
# workflow.py
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy
with workflow.unsafe.imports_passed_through():
from activities import read_sitekey, solve_recaptcha, submit_form
@workflow.defn
class SubmitProtectedForm:
@workflow.run
async def run(self, page_url: str) -> int:
sitekey = await workflow.execute_activity(
read_sitekey, page_url,
start_to_close_timeout=timedelta(seconds=60),
)
# Solve directly before the submit. Tokens go stale.
token = await workflow.execute_activity(
solve_recaptcha, args=[sitekey, page_url],
start_to_close_timeout=timedelta(minutes=6),
retry_policy=RetryPolicy(
maximum_attempts=4,
non_retryable_error_types=["ValidationException"],
),
)
return await workflow.execute_activity(
submit_form, args=[page_url, token],
start_to_close_timeout=timedelta(seconds=60),
)注意那两个带两个参数的 activity 上的 args 列表。单个位置参数可以直接传,但超过一个就必须通过 args 传递,写错这一点是这个 SDK 里最常见的第一个错误。
# worker.py - run this where CapSkip can be reached
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker
from activities import read_sitekey, solve_recaptcha, submit_form
from workflow import SubmitProtectedForm
async def main():
client = await Client.connect("localhost:7233")
worker = Worker(
client,
task_queue="captcha-queue",
workflows=[SubmitProtectedForm],
activities=[read_sitekey, solve_recaptcha, submit_form],
)
await worker.run()
asyncio.run(main())worker 在哪里运行,识别工具又在哪里运行
Temporal 把这两者切分得很干净,而且这对你有利。Temporal Service 负责调度任务并保存历史。它从不执行你的代码。你在自己基础设施里启动的 worker 会保持一条向外的长连接连到服务端,从队列里领取任务。你的网络上不需要开放任何入站连接。
所以和 CapSkip 在同一台 Windows 机器上的 worker,调用 127.0.0.1:8080 的方式和你桌上的脚本完全一样,在 Temporal Cloud 上也是如此。Temporal Cloud 里的云指的是编排层,不是计算层。
worker 一搬家,变的只有 host,别的都不变。跑在容器、Linux 虚拟机或 Kubernetes 里的 worker 无法通过回环地址访问 Windows 上的识别工具,所以识别工具要切到 Server mode(服务器模式)。Local mode(本地模式)绑定 127.0.0.1,只响应本机。Server mode 绑定你的内网或公网 IP,固定公网 IP 能让地址保持不变。两种模式下硬件都还是你自己的,也都不计次,所以一天跑一万次的工作流和只跑一次的成本一样。
# One environment variable, no code change. # CAPSKIP_HOST=10.0.0.12 on the worker. solver = AsyncCapSkip(host=os.environ["CAPSKIP_HOST"], port=8080)
一旦识别工具监听在网络地址上,就打开密钥校验,并给每一组 worker 集群分配各自的密钥,这样吊销其中一个不会影响其他。两种模式的完整步骤见 CapSkip 设置指南.
常见错误及其含义
| 你所看到的 | 原因 | 修复 |
|---|---|---|
| 重放时报非确定性错误 | 识别是从 workflow 代码里调用的 | 把它移进 activity,并用 execute_activity 调用 |
| 导入时报 RestrictedWorkflowAccessError | activity 模块被导入到了沙箱里 | 在 workflow.unsafe.imports_passed_through 里导入它 |
| activity 在识别中途被取消,然后重试 | start_to_close_timeout 比识别工具自身的超时还短 | 对 reCAPTCHA、Turnstile 和极验(GeeTest)把它设到 300 秒以上 |
| 表单拒绝了一个看起来没问题的 token | 工作流在识别和提交之间被阻塞了 | 让识别成为紧挨着提交的那一步 |
| 同一个失败白白重试了四次 | 一个确定性的错误被拿去重试了 | 把 ValidationException 和 ApiException 列为不可重试 |
| 每次尝试都报 NetworkException | worker 不在运行识别工具的那台机器上 | 把识别工具切到 Server mode,并设置 host |
| activity 报参数相关的 TypeError | 两个位置参数被直接传了进去 | 把它们作为列表通过 args 参数传入 |
| SDK 抛出 TimeoutException | 识别耗时超过了 recaptchaTimeout | 把它调到默认的 300 秒以上 |
常见问题
Temporal Cloud 上的工作流真的能调用 127.0.0.1 吗?
可以,因为 Temporal Cloud 并不运行你的代码。运行代码的是你的 worker,你在哪里启动它,它就在哪里,而且是向外连接到服务端。那台 worker 上的回环地址指的就是 worker 自己的机器,所以装在这台机器上的识别工具会正常响应。从开发服务器换到 Cloud,这一点不会有任何变化。
识别 activity 需要发心跳吗?
发不出有意义的心跳,因为 SDK 调用会一直阻塞到 token 返回,中间没有可以上报的位置。不如给 activity 设一个留足余量的 start_to_close_timeout,让失败的尝试交给重试策略处理。心跳适用于那些在你自己控制的任务上循环的 activity。
怎样批量识别又不把识别工具冲垮?
在 worker 上设置 max_concurrent_activities,或者给识别单独开一个任务队列和一个上限较低的 worker。这样限流发生在任务真正执行的地方,比试图把工作流的启动时间错开更可靠。在单个进程内,异步客户端会用 asyncio 并发展开,这部分内容见 并行识别验证码的指南.
这和在 Airflow 里做有什么不同?
Airflow 是一个带 DAG 的调度器,在构建图的过程中没什么能拦住你发网络请求,那是另一套坑。Temporal 是持久化执行,约束是确定性,答案永远是放进 activity。识别工具这一侧两边完全一样,Airflow 版本写在 Airflow 验证码指南.
简短版结论
把识别放进 activity,绝不要放在 workflow 方法里。给它一个比识别工具自身 300 秒上限更长的 start_to_close_timeout,把 ValidationException 和 ApiException 标为不可重试,并安排工作流的顺序,让识别成为提交之前的最后一步。把 worker 跑在 CapSkip 所在的机器上,host 就一直是 127.0.0.1。Python 的其余接口见 Python 验证码识别页面,同样的三个调用在 Node.js、PHP 和 C# 里都有,列在 验证码识别 SDK 页面。reCAPTCHA 本身的各项选项见 reCAPTCHA v2 识别页面.
在给它挂上定时任务之前,还有一点值得知道。CapSkip 是一款 验证码识别工具 ,跑在你已经拥有的硬件上,所以每分钟触发一次的工作流,和每周触发一次的成本相同。
