A distributed job can call CaptchaAI safely if one number governs it: the count of solves in flight across the whole cluster, capped at your plan's thread count as the threadsinfo endpoint reports it. Each engine has a primitive that enforces that cap: a Ray actor's max_concurrency, a fixed partition count times a per-partition thread pool in Spark, and Async I/O capacity times operator parallelism in Flink. The catch is that these engines also retry, speculate and replay work, and each of those re-runs a solve unless your records carry an idempotency key.
Below: one shared Python solve function, wired into Ray (core and Ray Data) and PySpark (batch and Structured Streaming), a Java version for Flink, and the cases where CAPTCHA work belongs outside the engine. The examples post tokens to forms you are authorized to automate, such as your own staging environment or a partner portal covered by a data agreement.
Plan threads versus cluster parallelism
CaptchaAI bills per concurrent thread, with unlimited solves per thread: BASIC ($15/month, 5 threads), STANDARD ($30/month, 15 threads), ADVANCE ($90/month, 50 threads) and PREMIUM ($170/month, 100 threads), with larger tiers on the pricing page. A thread is one CAPTCHA in flight. Your cluster thinks in cores instead: 64 Ray workers or a 200-task Spark stage will call the solver 64 or 200 times at once, and on a 15-thread plan most of those calls have nowhere to go.
The overflow does not arrive as a tidy rate-limit error. It comes back as ERROR_ZERO_BALANCE, which CaptchaAI's API reference describes as "You don't have free threads", or only as tasks sitting in CAPCHA_NOT_READY longer than usual. Either way, the fix is to never have more submissions in flight than you have threads.
So every job starts by asking the API. A POST to https://ocr.captchaai.com/res.php with key and action=threadsinfo returns a bare JSON object (no status wrapper): threads is your plan total as a string, working_threads the number busy right now.
{"threads": "600", "working_threads": 553}
Subtract a reserve for anything else on the same key (a second pipeline, a staging environment, teammates with the browser extension), and the result is your job cap. Then translate it into each engine's own knob:
| Engine | In-flight solves, kept at or below the cap | STANDARD with 3 reserved (cap 12) |
|---|---|---|
| Ray core | max_concurrency of one gate actor |
12 |
| Ray Data | ActorPoolStrategy(size=...) × threads per actor |
3 actors × 4 threads |
Spark batch or foreachBatch |
partitions × threads per partition | 4 partitions × 3 threads |
| Flink | Async I/O capacity × operator parallelism |
parallelism 4 × capacity 3 |
Workers beyond the cap buy nothing, because throughput is threads divided by solve time. CaptchaAI's published ceiling for reCAPTCHA v2 is under 60 seconds, so in the worst case a thread finishes about one solve a minute; Cloudflare Turnstile (under 10 seconds) cycles a thread much sooner. The arithmetic for a target volume is in how to process 10,000 CAPTCHA tasks per hour.
A shared solve function and its error classes
All three engines should call one client, so retry rules live in one place. The flow is the documented two-step one: submit to in.php with method=userrecaptcha, googlekey, pageurl and json=1 and get back {"status":1,"request":"<task id>"}; wait about 15 seconds, then call res.php with action=get, the id and json=1 every 5 seconds. CAPCHA_NOT_READY (no T) means pending, and status 1 carries the token in request. The client gives up after 120 seconds and fails the record rather than submitting again, because the unanswered task may still be holding a thread. What makes it safe inside an engine is how it sorts failures:
| Class | Codes | What the job does |
|---|---|---|
| Stop | ERROR_WRONG_USER_KEY, ERROR_KEY_DOES_NOT_EXIST, IP_BANNED |
Stop the whole job. Repeated bad-key requests earn an IP_BANNED that lasts 5 minutes |
| Threads | ERROR_ZERO_BALANCE |
Call threadsinfo. All threads busy: back off and retry. Threads idle: the account has no active plan, so stop |
| Transient | ERROR_SERVER_ERROR, ERROR_INTERNAL_SERVER_ERROR, an HTML error page |
At submit: retry after about 10 seconds, doubling, a bounded number of times. While polling: ask about the same ID again |
| Unsolvable | ERROR_CAPTCHA_UNSOLVABLE |
Stop polling that ID and allow one fresh submission |
| Bad record | ERROR_PAGEURL, ERROR_GOOGLEKEY, ERROR_WRONG_GOOGLEKEY, ERROR_BAD_TOKEN_OR_PAGEURL, ERROR_BAD_PARAMETERS, anything else |
Mark the record failed and never resend it unchanged |
"Stop the job" versus "fail the record" matters more in a cluster than on a laptop: a bad key treated as a per-record failure keeps every executor sending it, which is how a job earns an IP ban.
"""captcha_client.py: the CaptchaAI client every engine below imports."""
import os
import time
import requests
API = "https://ocr.captchaai.com"
API_KEY = os.environ.get("CAPTCHAAI_API_KEY", "YOUR_API_KEY")
STOP = {"ERROR_WRONG_USER_KEY", "ERROR_KEY_DOES_NOT_EXIST", "IP_BANNED"}
TRANSIENT = {"ERROR_SERVER_ERROR", "ERROR_INTERNAL_SERVER_ERROR"}
class StopJob(RuntimeError):
"""Key or plan problem: every other record would fail the same way."""
class BadTask(RuntimeError):
"""This record cannot be solved as sent; never resend it unchanged."""
class Unsolvable(BadTask):
"""ERROR_CAPTCHA_UNSOLVABLE: one fresh resubmission is allowed."""
class TransientError(RuntimeError):
"""Safe to retry after a back-off."""
class PollTimeout(TransientError):
"""No answer by the client deadline; the task may still hold a thread."""
def _post(path, data):
resp = requests.post(f"{API}/{path}", data=data, timeout=30)
try:
return resp.json()
except ValueError:
text = resp.text.strip()
if text.startswith("ERROR_") or text == "IP_BANNED":
return {"status": 0, "request": text} # plain-text code despite json=1
raise TransientError(f"unexpected {path} response: {text[:60]}")
def threads_info():
body = _post("res.php", {"key": API_KEY, "action": "threadsinfo"})
if "threads" not in body:
raise StopJob(f"threadsinfo rejected: {body.get('request')}")
return int(body["threads"]), int(body["working_threads"])
def job_cap(reserve=0):
"""In-flight cap for this job: plan threads minus what other consumers need."""
threads, working = threads_info()
cap = max(1, threads - reserve)
print(f"plan_threads={threads} busy_now={working} job_cap={cap}", flush=True)
return cap
def _raise_for(code):
if code in STOP:
raise StopJob(code)
if code in TRANSIENT:
raise TransientError(code)
if code == "ERROR_ZERO_BALANCE":
threads, working = threads_info()
if working >= threads:
raise TransientError("ERROR_ZERO_BALANCE: every plan thread is busy")
raise StopJob("ERROR_ZERO_BALANCE with idle threads: no active plan")
if code == "ERROR_CAPTCHA_UNSOLVABLE":
raise Unsolvable(code)
raise BadTask(code)
def _solve_once(sitekey, pageurl, max_wait):
sub = _post("in.php", {"key": API_KEY, "method": "userrecaptcha",
"googlekey": sitekey, "pageurl": pageurl, "json": 1})
if sub.get("status") != 1:
_raise_for(sub.get("request"))
task_id = str(sub["request"])
deadline = time.monotonic() + max_wait
time.sleep(15) # reCAPTCHA v2: no point polling earlier
while True:
try:
res = _post("res.php", {"key": API_KEY, "action": "get", "id": task_id, "json": 1})
except TransientError:
res = {} # an HTML error page: ask about the same ID again
if res.get("status") == 1:
return res["request"]
code = res.get("request")
if code not in (None, "CAPCHA_NOT_READY") and code not in TRANSIENT:
_raise_for(code)
if time.monotonic() > deadline:
raise PollTimeout(f"task {task_id} not ready after {max_wait}s")
time.sleep(5 if code == "CAPCHA_NOT_READY" else 10)
def solve_recaptcha_v2(sitekey, pageurl, max_wait=120):
"""One token. Submit-time transient errors get two more tries, UNSOLVABLE one."""
backoff, transient_left, unsolvable_left = 10, 2, 1
while True:
try:
return _solve_once(sitekey, pageurl, max_wait)
except PollTimeout:
raise # never stack a second task on one that may still be running
except Unsolvable:
if unsolvable_left == 0:
raise
unsolvable_left -= 1
except TransientError:
if transient_left == 0:
raise
transient_left -= 1
time.sleep(backoff)
backoff *= 2
def submit_record(record, solve=solve_recaptcha_v2):
"""Load the form, solve, post the token: one session, one record, no token kept."""
key = record["record_key"]
with requests.Session() as session:
try:
session.get(record["pageurl"], timeout=30).raise_for_status()
token = solve(record["sitekey"], record["pageurl"])
resp = session.post(record["submit_url"],
data={"g-recaptcha-response": token}, timeout=30)
except (BadTask, TransientError, requests.RequestException) as exc:
return key, "failed", str(exc)
return key, "submitted", str(resp.status_code)
submit_record encodes the rule that matters most here: solve and use the token in the same task and HTTP session, then drop it. Engines make it tempting to write tokens to a table for a later stage. Don't: CaptchaAI's docs say reCAPTCHA v3 tokens expire within minutes and Turnstile tokens are single-use, and a stage boundary can add minutes of queueing. Every record carries a record_key, which all the duplicate protection below hangs on.
Ray: one gate actor owns the cap
In Ray, the cleanest cluster-wide cap is a single actor that does the solving. Per the Ray docs on actor concurrency, each call to a threaded actor runs in a thread pool whose size is limited by max_concurrency, so creating the gate once with max_concurrency=cap means no more than cap solves run anywhere in the cluster, whichever tasks ask. Extra calls wait their turn at the actor.
"""ray_job.py: every Ray task shares one CaptchaAI cap held by a threaded actor."""
import json
import os
import sys
import ray
from captcha_client import job_cap, solve_recaptcha_v2, submit_record
@ray.remote
class SolverGate:
"""A threaded actor: at most max_concurrency solve() calls run at once."""
def solve(self, sitekey, pageurl):
return solve_recaptcha_v2(sitekey, pageurl)
@ray.remote(max_retries=1) # node or worker loss only; no retry_exceptions
def process_record(gate, record):
def gated_solve(sitekey, pageurl):
return ray.get(gate.solve.remote(sitekey, pageurl))
return submit_record(record, solve=gated_solve)
def main(jobs_path):
ray.init(namespace="captchaai", runtime_env={
"working_dir": ".", # ships captcha_client.py to every node
"pip": ["requests"],
"env_vars": {"CAPTCHAAI_API_KEY": os.environ["CAPTCHAAI_API_KEY"]},
})
cap = job_cap(reserve=int(os.environ.get("CAPTCHAAI_RESERVE", "0")))
gate = SolverGate.options(name="solver-gate", get_if_exists=True,
max_concurrency=cap).remote()
pending, window = [], 2 * cap # bound queued tasks instead of submitting all
with open(jobs_path, encoding="utf-8") as fh:
for line in fh:
if len(pending) >= window:
done, pending = ray.wait(pending, num_returns=1)
print(ray.get(done[0]), flush=True)
pending.append(process_record.remote(gate, json.loads(line)))
while pending:
done, pending = ray.wait(pending, num_returns=1)
print(ray.get(done[0]), flush=True)
if __name__ == "__main__":
main(sys.argv[1])
Details that do real work here:
- Threaded, not async. Non-async actors default to a concurrency of 1 and async actors to 1,000. Rewrite the gate with
async defmethods and Ray treats it as an async actor, so setmax_concurrencyexplicitly or the default is effectively no cap. - Errors cross the actor boundary intact. Ray re-raises an actor's exception at
ray.getas bothRayTaskErrorand your own class, sosubmit_recordstill catchesBadTask, whileStopJobreaches the driver and ends the run. - Retries stay out of the solve path. Ray reruns a task after worker or node loss (three times by default) but retries application exceptions only if you set
retry_exceptions. A rerun repeats the solve, so this job allows one system retry and keeps API retries insidesolve_recaptcha_v2. If you useretry_exceptions, list only exceptions that can occur before the solve. - The driver applies backpressure.
ray.wait(pending, num_returns=1)returns results as they complete, and the2 * capwindow follows Ray's own pattern for bounding the pending-task queue. env_varsinruntime_envdelivers the key. Variables already set on the cluster are visible to workers anyway.
Named actors live in a namespace, which is why ray.init sets one. For two drivers on one key, add lifetime="detached" so the gate outlives its creator. The first creator fixes max_concurrency, and detached actors are not garbage-collected, so ray.kill the gate when your plan changes. If records don't need a task of their own, ray.util.ActorPool over N plain actors that each call submit_record gives the same cap of N.
Ray Data: size the actor pool, not the batch
In Ray Data the knob is the compute strategy of map_batches (current docs deprecate the older concurrency= argument in its favor). An actor runs one batch at a time by default, so a fixed pool of size actors, each with a small thread pool, has at most size × THREADS_PER_ACTOR solves in flight:
"""ray_data_job.py: the same cap with Ray Data: pool size times threads per actor."""
import os
from concurrent.futures import ThreadPoolExecutor
import pandas as pd
import ray
from captcha_client import job_cap, submit_record
THREADS_PER_ACTOR = 4
class SolveBatch:
def __call__(self, batch: pd.DataFrame) -> pd.DataFrame:
with ThreadPoolExecutor(max_workers=THREADS_PER_ACTOR) as pool:
rows = list(pool.map(submit_record, batch.to_dict("records")))
return pd.DataFrame(rows, columns=["record_key", "outcome", "detail"])
ray.init(runtime_env={"working_dir": ".", "pip": ["requests"],
"env_vars": {"CAPTCHAAI_API_KEY": os.environ["CAPTCHAAI_API_KEY"]}})
cap = job_cap(reserve=int(os.environ.get("CAPTCHAAI_RESERVE", "0")))
ds = ray.data.read_json(os.environ["JOBS_PATH"])
ds.map_batches(SolveBatch, batch_format="pandas", batch_size=THREADS_PER_ACTOR,
compute=ray.data.ActorPoolStrategy(size=max(1, cap // THREADS_PER_ACTOR))
).write_parquet(os.environ["RESULTS_PATH"])
Use a fixed size, not a min_size/max_size range: an autoscaling pool grows with the data, which is exactly the number you are holding still. Leaving compute out is worse, because a callable class then gets an autoscaling pool from one to an unlimited number of actors.
Spark: fixed partitions times a bounded thread pool
Resist the per-row Python UDF. Its concurrency follows partition counts and executor slots that change with data size and scaling, and PySpark's UDF notes warn that a deterministic UDF may be invoked more times than it appears in the query; every extra invocation is a solve. Instead, repartition(k) fixes the number of tasks and each task runs its partition through a pool of m threads, so no more than k × m solves are ever in flight.
"""spark_batch.py: k partitions x m threads never exceeds the CaptchaAI cap."""
import sys
from concurrent.futures import ThreadPoolExecutor
from pyspark.sql import SparkSession
from captcha_client import job_cap, submit_record
THREADS_PER_TASK = 3
OUT_SCHEMA = "record_key: string, outcome: string, detail: string"
def solve_partition(rows):
with ThreadPoolExecutor(max_workers=THREADS_PER_TASK) as pool:
yield from pool.map(submit_record, rows)
def main(jobs_path, results_path):
spark = SparkSession.builder.appName("captcha-batch").getOrCreate()
cap = job_cap() # also rejects a bad key before any executor starts
partitions = max(1, cap // THREADS_PER_TASK)
jobs = spark.read.json(jobs_path) # record_key, sitekey, pageurl, submit_url
solved = jobs.repartition(partitions).rdd.mapPartitions(solve_partition)
spark.createDataFrame(solved, OUT_SCHEMA).write.mode("overwrite").parquet(results_path)
if __name__ == "__main__":
main(sys.argv[1], sys.argv[2])
Ship the client with spark-submit --py-files captcha_client.py,spark_batch.py and give executors the key with --conf spark.executorEnv.CAPTCHAAI_API_KEY=...; the driver needs it too, because job_cap() runs there. With fewer free task slots than partitions, real concurrency is lower than k × m, never higher.
Two settings decide how often a partition runs twice. spark.speculation defaults to false; keep it that way, because speculation relaunches slow tasks and CAPTCHA partitions are slow by design. spark.task.maxFailures defaults to 4, allowing three retries, and a retried partition re-solves every row it had finished. That is why submit_record turns per-record errors into outcome rows instead of failing the task. On Spark Connect, where the RDD API isn't available, put the same bounded pool inside DataFrame.mapInPandas, which also receives each partition as an iterator.
Structured Streaming: foreachBatch, maxOffsetsPerTrigger and a keyed sink
For a stream, run the same partition function inside foreachBatch. The Structured Streaming guide is explicit that foreachBatch gives at-least-once writes by default and suggests using the batch ID to deduplicate. After a failure, a micro-batch can run again, and without a guard every record in it is solved again. The guard below skips any record_key that already has an outcome row:
"""spark_stream.py: Kafka -> foreachBatch -> bounded solving -> keyed outcome table."""
import os
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, from_json, lit
from captcha_client import job_cap
from spark_batch import OUT_SCHEMA, THREADS_PER_TASK, solve_partition
TABLE = "captcha_outcomes"
JOB_SCHEMA = "record_key STRING, sitekey STRING, pageurl STRING, submit_url STRING"
spark = SparkSession.builder.appName("captcha-stream").getOrCreate()
CAP = job_cap(reserve=int(os.environ.get("CAPTCHAAI_RESERVE", "0")))
PARTITIONS = max(1, CAP // THREADS_PER_TASK)
def solve_batch(batch_df, batch_id):
session = batch_df.sparkSession
todo = batch_df.dropDuplicates(["record_key"])
if session.catalog.tableExists(TABLE): # a replayed batch skips finished records
todo = todo.join(session.table(TABLE).select("record_key"), "record_key", "left_anti")
solved = todo.repartition(PARTITIONS).rdd.mapPartitions(solve_partition)
(session.createDataFrame(solved, OUT_SCHEMA)
.withColumn("batch_id", lit(batch_id))
.write.mode("append").saveAsTable(TABLE))
jobs = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", os.environ["KAFKA_BOOTSTRAP"])
.option("subscribe", "captcha-jobs")
.option("maxOffsetsPerTrigger", str(CAP)) # worst case: one solve per thread per minute
.load()
.select(from_json(col("value").cast("string"), JOB_SCHEMA).alias("job"))
.select("job.*"))
(jobs.writeStream.foreachBatch(solve_batch)
.option("checkpointLocation", os.environ["CHECKPOINT_DIR"])
.trigger(processingTime="1 minute")
.start()
.awaitTermination())
maxOffsetsPerTrigger caps how many records one micro-batch pulls from Kafka, so a backlog doesn't become one enormous batch, and it bounds how much a replay can repeat. CAP records per one-minute trigger assumes the worst-case reCAPTCHA v2 time; raise it once you have measured real solve times. Micro-batches of one query run one after another (a late trigger fires as soon as the previous batch finishes), so the batch cap is the query's cap, but separate queries in one application run concurrently and each needs its own share of threads. The Kafka source needs the spark-sql-kafka-0-10 connector, and on a long-lived table a merge keyed on record_key beats a full anti-join; idempotent CAPTCHA solving covers key design.
Flink: Async I/O capacity times parallelism
Flink's Async I/O operator already has the knob: capacity is how many asynchronous requests may be in progress per parallel instance, so capacity × parallelism of that operator is your in-flight count. The docs also warn that asyncInvoke is not called multi-threaded, so a Thread.sleep for the 15-second wait would stall the subtask. The function below never blocks: it chains HttpClient.sendAsync calls and waits with CompletableFuture.delayedExecutor.
// CaptchaSolveJob.java: Flink Async I/O where capacity x parallelism <= plan threads.
import java.io.Serializable;
import java.net.CookieManager;
import java.net.URI;
import java.net.URLEncoder;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.Collection;
import java.util.Collections;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.TimeUnit;
import java.util.function.Predicate;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
import org.apache.flink.streaming.api.datastream.AsyncDataStream;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.async.AsyncFunction;
import org.apache.flink.streaming.api.functions.async.AsyncRetryStrategy;
import org.apache.flink.streaming.api.functions.async.ResultFuture;
import org.apache.flink.streaming.util.retryable.AsyncRetryStrategies;
public class CaptchaSolveJob {
static final String API = "https://ocr.captchaai.com/";
static final Set<String> STOP = Set.of("ERROR_WRONG_USER_KEY", "ERROR_KEY_DOES_NOT_EXIST", "IP_BANNED");
static final Set<String> RETRYABLE = Set.of("ERROR_SERVER_ERROR", "ERROR_INTERNAL_SERVER_ERROR",
"ERROR_ZERO_BALANCE", "ERROR_CAPTCHA_UNSOLVABLE");
/** One record to solve; a Flink POJO (public fields, no-arg constructor). */
public static class CaptchaJob {
public String recordKey;
public String sitekey;
public String pageurl;
public String submitUrl;
public CaptchaJob() {}
public CaptchaJob(String recordKey, String sitekey, String pageurl, String submitUrl) {
this.recordKey = recordKey;
this.sitekey = sitekey;
this.pageurl = pageurl;
this.submitUrl = submitUrl;
}
}
static class RetryableError extends RuntimeException {
RetryableError(String code) { super(code); }
}
static class StopJobError extends RuntimeException {
StopJobError(String code) { super(code); }
}
/** Retries a record whose outcome row says "retry"; once retries run out, Flink emits that row. */
static class NeedsRetry implements Predicate<Collection<String>>, Serializable {
@Override
public boolean test(Collection<String> rows) {
return rows.stream().anyMatch(row -> row.contains(",retry,"));
}
}
static class CaptchaAsyncFunction implements AsyncFunction<CaptchaJob, String> {
private final String apiKey;
private transient HttpClient api;
CaptchaAsyncFunction(String apiKey) { this.apiKey = apiKey; }
@Override
public void asyncInvoke(CaptchaJob job, ResultFuture<String> result) {
if (api == null) {
api = HttpClient.newHttpClient();
}
// A cookie jar per record: the page load and the form post share one session.
HttpClient site = HttpClient.newBuilder().cookieHandler(new CookieManager()).build();
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(120);
send(site, HttpRequest.newBuilder(URI.create(job.pageurl)).GET())
.thenCompose(page -> call("in.php", Map.of("key", apiKey, "method", "userrecaptcha",
"googlekey", job.sitekey, "pageurl", job.pageurl, "json", "1")))
.thenCompose(body -> "1".equals(field(body, "status"))
? poll(field(body, "request"), 15, deadline)
: CompletableFuture.<String>failedFuture(error(field(body, "request"))))
.thenCompose(token -> send(site, formPost(job.submitUrl, Map.of("g-recaptcha-response", token))))
.whenComplete((resp, err) -> {
Throwable cause = err instanceof CompletionException ? err.getCause() : err;
if (cause instanceof StopJobError) {
result.completeExceptionally(cause); // bad key: fail the job
} else if (cause instanceof RetryableError) {
result.complete(Collections.singleton(job.recordKey + ",retry," + cause.getMessage()));
} else if (cause != null) {
result.complete(Collections.singleton(job.recordKey + ",failed," + cause.getMessage()));
} else {
result.complete(Collections.singleton(job.recordKey + ",submitted," + resp.statusCode()));
}
});
}
@Override
public void timeout(CaptchaJob job, ResultFuture<String> result) {
// The default timeout() fails the task and restarts the job; record the miss instead.
result.complete(Collections.singleton(job.recordKey + ",failed,TIMEOUT"));
}
private CompletableFuture<String> poll(String id, long delaySeconds, long deadline) {
return CompletableFuture.supplyAsync(() -> id,
CompletableFuture.delayedExecutor(delaySeconds, TimeUnit.SECONDS))
.thenCompose(ignored -> call("res.php", Map.of("key", apiKey, "action", "get", "id", id, "json", "1")))
.thenCompose(body -> {
String code = field(body, "request");
if ("1".equals(field(body, "status"))) {
return CompletableFuture.completedFuture(code);
}
if (System.nanoTime() > deadline) {
return CompletableFuture.<String>failedFuture(new IllegalStateException("POLL_DEADLINE " + id));
}
if ("CAPCHA_NOT_READY".equals(code)) {
return poll(id, 5, deadline);
}
if ("ERROR_INTERNAL_SERVER_ERROR".equals(code) || code.startsWith("<")) {
return poll(id, 10, deadline); // server error or HTML error page: same ID again
}
return CompletableFuture.<String>failedFuture(error(code));
});
}
private CompletableFuture<String> call(String path, Map<String, String> params) {
return send(api, formPost(API + path, params)).thenApply(HttpResponse::body);
}
}
static HttpRequest.Builder formPost(String url, Map<String, String> params) {
String body = params.entrySet().stream()
.map(e -> e.getKey() + "=" + URLEncoder.encode(e.getValue(), StandardCharsets.UTF_8))
.collect(Collectors.joining("&"));
return HttpRequest.newBuilder(URI.create(url))
.header("Content-Type", "application/x-www-form-urlencoded")
.POST(HttpRequest.BodyPublishers.ofString(body));
}
static CompletableFuture<HttpResponse<String>> send(HttpClient client, HttpRequest.Builder request) {
return client.sendAsync(request.timeout(Duration.ofSeconds(20)).build(), HttpResponse.BodyHandlers.ofString());
}
/** Reads a field from {"status":1,"request":"..."}; a plain-text code comes back as-is. */
static String field(String body, String name) {
Matcher m = Pattern.compile("\"" + name + "\"\\s*:\\s*\"?([^\",}]*)").matcher(body);
return m.find() ? m.group(1) : body.trim();
}
static RuntimeException error(String code) {
if (STOP.contains(code)) {
return new StopJobError(code);
}
return RETRYABLE.contains(code) ? new RetryableError(code) : new IllegalStateException(code);
}
static int planThreads(String apiKey) throws Exception {
HttpResponse<String> resp = HttpClient.newHttpClient().send(
formPost(API + "res.php", Map.of("key", apiKey, "action", "threadsinfo")).build(),
HttpResponse.BodyHandlers.ofString());
String threads = field(resp.body(), "threads");
if (!threads.matches("\\d+")) {
throw new IllegalStateException("threadsinfo failed: " + resp.body());
}
System.out.println("threadsinfo " + resp.body());
return Integer.parseInt(threads);
}
public static void main(String[] args) throws Exception {
String apiKey = System.getenv().getOrDefault("CAPTCHAAI_API_KEY", "YOUR_API_KEY");
int reserve = Integer.parseInt(System.getenv().getOrDefault("CAPTCHAAI_RESERVE", "0"));
int cap = Math.max(1, planThreads(apiKey) - reserve);
int parallelism = Math.min(cap, 4);
int capacity = cap / parallelism; // capacity * parallelism <= cap
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<CaptchaJob> jobs = env.fromData(new CaptchaJob("order-1001",
System.getenv("TARGET_SITEKEY"), System.getenv("TARGET_PAGEURL"), System.getenv("TARGET_SUBMIT_URL")));
AsyncRetryStrategy<String> retry = new AsyncRetryStrategies.FixedDelayRetryStrategyBuilder<String>(1, 10_000L)
.ifResult(new NeedsRetry())
.build();
AsyncDataStream.unorderedWaitWithRetry(jobs, new CaptchaAsyncFunction(apiKey),
360, TimeUnit.SECONDS, capacity, retry)
.setParallelism(parallelism)
.print();
env.execute("captchaai-async-solve");
}
}
The numbers in main are deliberate:
- Parallelism is set on the async operator itself, so
capacity × parallelismdoesn't drift when the job default changes. - The timeout is 360 seconds. One attempt takes at most about 165 seconds: the 120-second poll deadline (which covers the 15-second initial wait, reCAPTCHA v2's under-60-second ceiling and polling overhead), a last delayed poll and the 20-second form post. Flink's docs define the timeout as running from the first invocation to final completion, retries included, so 360 seconds covers two full attempts plus the retry delay. The overridden
timeout()emits a failed outcome instead of failing the job. - Retries use
unorderedWaitWithRetry(Flink 1.16 and later) with a fixed 10-second delay and a result predicate, not an exception predicate. Server errors, busy threads and unsolvable CAPTCHAs complete the record with aretryrow, and a retry re-invokesasyncInvoke, so it reloads the page and submits a fresh task.maxAttemptsis 1, which in Flink's fixed-delay strategy means one retry: its attempt counter starts at 1 and it retries while the counter is at mostmaxAttempts. When retries run out, Flink emits the lastretryrow. Had the retryable path completed exceptionally instead, exhausted retries would fail the job, and the restart would replay the same poison record. Only stop-class codes fail the job. fromDataneeds Flink 1.19 or later; older releases usefromElements. With a bounded source like this one, Flink stops retrying once the input ends and emits whatever the last attempt returned, so test retries against an unbounded source. In production that is a Kafka topic, and the sink upserts onrecordKey, treating a finalretryrow as a failure.
The docs' fault-tolerance note is the duplicate warning in disguise: the operator stores in-flight requests in checkpoints and re-triggers them on recovery, so a restart resubmits up to capacity × parallelism solves. PyFlink gets the same operator from Flink 2.2 (AsyncDataStream.unordered_wait with an async def async_invoke) and the same arithmetic.
Where duplicate solves come from
On a thread plan a duplicate doesn't add a per-solve charge, but it holds a thread for up to a minute that a real record could have used. Each engine has its own ways to re-run a solve:
| Source | Engine and default | What runs again | Fix |
|---|---|---|---|
| Task retry | Ray max_retries (3, after worker or node loss), Spark spark.task.maxFailures (4 attempts, after any task failure) |
Every record in the task or partition, including finished ones | Check record_key before solving; keep per-record errors inside the task |
| Retrying on application exceptions | Ray retry_exceptions (off by default) |
The whole task, solve included | Retry only exceptions raised before the solve |
| Speculative execution | Spark spark.speculation (false) |
A second copy of slow partitions, which CAPTCHA partitions always are | Keep it off for solving stages |
| Micro-batch replay | Structured Streaming foreachBatch (at-least-once) |
A batch that did not finish committing | Anti-join or merge on record_key; bound batch size with maxOffsetsPerTrigger |
| Checkpoint recovery | Flink Async I/O | Requests that were in flight at the failure | Upsert on recordKey; expect up to capacity × parallelism repeats per restart |
| Client timeout, then resubmit | Any | The same widget, as a new task | Poll deadline above the type's ceiling; CaptchaAI documents no cancel action |
The last row catches people out: when your client gives up at 120 seconds, there is no call to cancel the task, so treat it as still occupying capacity for a while rather than firing a replacement straight away.
When to hand CAPTCHA work to a queue instead
Solving inside the engine fits when the CAPTCHA step is short relative to the job, one job owns the plan's threads, and you control the retry model. Move it out when several jobs or services share one key and need a single arbiter, when the protected action needs a browser session in a separate worker fleet, when solve latency stretches micro-batches or checkpoints, or when recovery replays re-trigger enough solves to matter.
The hand-off: the engine emits CAPTCHA jobs (record key, sitekey, page URL) to a Kafka topic, and a worker fleet sized so its total in-flight solves equal your cap consumes them, performs the protected action and publishes an outcome keyed by record, never a bare token. The thread cap then lives in exactly one place. Kafka and CaptchaAI for streaming CAPTCHA tasks shows the topic layout, backpressure in CAPTCHA solving queues covers a fleet that falls behind, and auto-scaling CAPTCHA solving workers covers sizing that fleet from queue depth. The cap itself is a semaphore; semaphore patterns for CAPTCHA concurrency compares the variants. If the pipeline collects web data, CAPTCHA handling for LLM training-data scraping covers staying within sources you may use.
Monitoring working_threads during a run
Poll threadsinfo from a small sidecar while the job runs; comparing working_threads with your cap shows whether the engine is doing what you configured:
"""watch_threads.py: log CaptchaAI thread use every 30 seconds during a job."""
import time
from captcha_client import threads_info
while True:
threads, working = threads_info()
print(f"{time.strftime('%H:%M:%S')} working={working}/{threads} ({working / threads:.0%})", flush=True)
time.sleep(30)
Read it against the cap logged at job start. If working_threads sits at threads while your job runs below its cap, something else is using the key and your reserve is too small. If it hovers well below the cap while records wait, the engine isn't reaching it: fewer free Spark slots than partitions, a Ray Data pool still starting, or a Flink source that can't feed the operator fast enough. A healthy run shows working_threads close to the cap and flat.
Troubleshooting
| Symptom | Likely cause in a distributed job | Fix |
|---|---|---|
Bursts of ERROR_ZERO_BALANCE right as a stage starts |
All tasks submit at once and k × m (or capacity × parallelism) exceeds the threads that are actually free |
Recompute the cap from threadsinfo with a reserve; check that no second query or job shares the key |
| Many records fail with "not ready after 120s" | Another consumer saturates the threads, so tasks queue on the API side | Watch working_threads; lower the cap or raise the reserve |
Ray gate ignores a new max_concurrency |
get_if_exists=True returned the existing (possibly detached) actor |
ray.kill it and let the next job recreate it |
| Flink job restarts on slow CAPTCHAs | Default timeout() completes the record exceptionally, which fails the task |
Override timeout() and set the timeout above all attempts |
| Flink job restarts on one unsolvable record | Retryable errors complete exceptionally, so exhausted retries fail the job | Emit a retry outcome row and retry with ifResult, as above |
Rows end in retry,ERROR_ZERO_BALANCE while working_threads is low |
Threads are idle, so the account has no active plan | Stop the job and check the plan; more retries won't help |
Every Spark task fails with StopJob |
Bad or missing key on the executors | Set spark.executorEnv.CAPTCHAAI_API_KEY; the driver's job_cap() call catches a bad key before executors start |
FAQ
Can I limit concurrency with executor cores or Ray num_cpus instead?
Only indirectly, and not reliably. Cores bound how many tasks run, not how many HTTP calls each task makes, and autoscaling changes the number under you. An explicit cap (gate actor, k × m, capacity × parallelism) stays correct whatever the cluster does.
Which error do I get when I exceed my plan's threads?
ERROR_ZERO_BALANCE, which CaptchaAI's API reference describes as having no free threads. Saturation can also appear only as longer CAPCHA_NOT_READY polling. Either way, diagnose with threadsinfo: busy threads mean back off, idle threads mean the account has no active plan.
Does this work in PyFlink or only in Java?
Flink documents a Python Async I/O API from release 2.2, built on AsyncDataStream.unordered_wait and an async def async_invoke, and the capacity × parallelism rule is identical. On older Flink versions, use the Java operator above or hand CAPTCHA jobs to a queue.
Should I add workers or buy more threads?
If working_threads sits at your cap while records wait, only more threads raise throughput; more workers would just queue. If it sits below the cap, fix the engine configuration first, then compare the tiers on the pricing page against solves per hour per thread for your CAPTCHA type.
Start with threadsinfo, write the cap into the one engine setting that enforces it, and give every record a key before its first solve.