Một loạt 500 tác vụ CAPTCHA hoàn thành không đồng đều — một số giải quyết trong 8 giây, một số khác mất 45 giây. Việc chờ đợi mọi tác vụ hoàn thành trước khi xử lý kết quả sẽ lãng phí thời gian giữa giải pháp đầu tiên và giải pháp cuối cùng. Tính năng phát trực tuyến cho phép quy trình xuôi dòng của bạn tiêu thụ từng kết quả ngay khi nó đến.
Truyền trực tuyến so với xử lý hàng loạt
| Cách tiếp cận | Thời gian đến kết quả đầu tiên | Bộ nhớ | Độ trễ đường ống |
|---|---|---|---|
| Đợi tất cả | Sau nhiệm vụ chậm nhất | Tất cả kết quả trong bộ nhớ | Cao |
| Truyền phát như đã giải quyết | Sau nhiệm vụ nhanh nhất | Mỗi lần một kết quả | Thấp |
| Lô nhỏ (khối 10) | Sau đoạn đầu tiên | 10 kết quả cùng một lúc | Trung bình |
Python: Trình tạo Async cho kết quả phát trực tuyến
Sử dụng asyncio và aiohttp, mỗi giải pháp mang lại kết quả ngay lập tức thông qua trình tạo không đồng bộ:
import asyncio
import aiohttp
import time
API_KEY = "YOUR_API_KEY"
SUBMIT_URL = "https://ocr.captchaai.com/in.php"
RESULT_URL = "https://ocr.captchaai.com/res.php"
async def submit_task(session, task_data):
"""Submit a single CAPTCHA task."""
params = {
"key": API_KEY,
"method": task_data.get("method", "userrecaptcha"),
"json": 1,
}
if params["method"] == "userrecaptcha":
params["googlekey"] = task_data["sitekey"]
params["pageurl"] = task_data["pageurl"]
elif params["method"] == "turnstile":
params["sitekey"] = task_data["sitekey"]
params["pageurl"] = task_data["pageurl"]
async with session.post(SUBMIT_URL, data=params) as resp:
result = await resp.json(content_type=None)
if result.get("status") != 1:
return None, result.get("request", "unknown")
return result["request"], None
async def poll_task(session, task_id, timeout=300):
"""Poll until solved or timeout."""
start = time.monotonic()
while time.monotonic() - start < timeout:
await asyncio.sleep(5)
params = {"key": API_KEY, "action": "get", "id": task_id, "json": 1}
async with session.get(RESULT_URL, params=params) as resp:
result = await resp.json(content_type=None)
if result.get("request") == "CAPCHA_NOT_READY":
continue
if result.get("status") == 1:
return result["request"], None
return None, result.get("request", "unknown")
return None, "TIMEOUT"
async def solve_one(session, index, task_data, semaphore):
"""Solve a single task within concurrency limits."""
async with semaphore:
start = time.monotonic()
task_id, error = await submit_task(session, task_data)
if error:
return {"index": index, "status": "failed", "error": error, "time": 0}
token, error = await poll_task(session, task_id)
elapsed = time.monotonic() - start
if token:
return {"index": index, "status": "solved", "token": token, "time": round(elapsed, 1)}
return {"index": index, "status": "failed", "error": error, "time": round(elapsed, 1)}
async def stream_results(tasks, max_concurrent=20):
"""
Async generator that yields each result as it completes.
Results arrive in completion order, not submission order.
"""
semaphore = asyncio.Semaphore(max_concurrent)
async with aiohttp.ClientSession() as session:
pending = set()
for i, task in enumerate(tasks):
coro = solve_one(session, i, task, semaphore)
pending.add(asyncio.ensure_future(coro))
while pending:
done, pending = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED)
for future in done:
yield future.result()
async def main():
tasks = [
{"sitekey": "SITE_KEY", "pageurl": f"https://example.com/page{i}"}
for i in range(50)
]
solved = 0
failed = 0
async for result in stream_results(tasks, max_concurrent=15):
# Process each result immediately
if result["status"] == "solved":
solved += 1
print(f" [{solved + failed}/{len(tasks)}] Task {result['index']} SOLVED in {result['time']}s")
# Use token immediately — don't wait for batch
# await submit_form(result["token"])
# await save_to_database(result)
else:
failed += 1
print(f" [{solved + failed}/{len(tasks)}] Task {result['index']} FAILED: {result['error']}")
print(f"\nDone: {solved} solved, {failed} failed")
asyncio.run(main())
Cài đặt phụ thuộc:
pip install aiohttp
JavaScript: Mẫu phát trực tuyến EventEmitter
Node.js sử dụng cách tiếp cận theo hướng sự kiện - đưa ra từng kết quả khi nó giải quyết:
const { EventEmitter } = require("events");
const API_KEY = "YOUR_API_KEY";
const SUBMIT_URL = "https://ocr.captchaai.com/in.php";
const RESULT_URL = "https://ocr.captchaai.com/res.php";
class CaptchaStream extends EventEmitter {
constructor(maxConcurrent = 15) {
super();
this.maxConcurrent = maxConcurrent;
this.active = 0;
this.queue = [];
this.total = 0;
this.completed = 0;
}
async submitAndPoll(index, taskData) {
const params = new URLSearchParams({
key: API_KEY,
method: taskData.method || "userrecaptcha",
googlekey: taskData.sitekey,
pageurl: taskData.pageurl,
json: "1",
});
const start = Date.now();
const submitResp = await (await fetch(SUBMIT_URL, { method: "POST", body: params })).json();
if (submitResp.status !== 1) {
return { index, status: "failed", error: submitResp.request, time: 0 };
}
const taskId = submitResp.request;
for (let i = 0; i < 60; i++) {
await new Promise((r) => setTimeout(r, 5000));
const url = `${RESULT_URL}?key=${API_KEY}&action=get&id=${taskId}&json=1`;
const poll = await (await fetch(url)).json();
if (poll.request === "CAPCHA_NOT_READY") continue;
const elapsed = ((Date.now() - start) / 1000).toFixed(1);
if (poll.status === 1) return { index, status: "solved", token: poll.request, time: elapsed };
return { index, status: "failed", error: poll.request, time: elapsed };
}
return { index, status: "failed", error: "TIMEOUT", time: ((Date.now() - start) / 1000).toFixed(1) };
}
async processNext() {
if (this.queue.length === 0 || this.active >= this.maxConcurrent) return;
const { index, taskData } = this.queue.shift();
this.active++;
try {
const result = await this.submitAndPoll(index, taskData);
this.emit("result", result);
} catch (err) {
this.emit("result", { index, status: "failed", error: err.message });
} finally {
this.active--;
this.completed++;
if (this.completed === this.total) {
this.emit("done");
} else {
this.processNext();
}
}
}
start(tasks) {
this.total = tasks.length;
this.queue = tasks.map((taskData, index) => ({ index, taskData }));
// Launch initial batch
const initial = Math.min(this.maxConcurrent, tasks.length);
for (let i = 0; i < initial; i++) {
this.processNext();
}
return this;
}
}
// Usage
const tasks = Array.from({ length: 50 }, (_, i) => ({
sitekey: "SITE_KEY",
pageurl: `https://example.com/page${i}`,
}));
const stream = new CaptchaStream(15);
let solved = 0, failed = 0;
stream.on("result", (result) => {
if (result.status === "solved") {
solved++;
console.log(`[${solved + failed}/${tasks.length}] Task ${result.index} SOLVED (${result.time}s)`);
// Use token immediately
// submitForm(result.token);
} else {
failed++;
console.log(`[${solved + failed}/${tasks.length}] Task ${result.index} FAILED: ${result.error}`);
}
});
stream.on("done", () => {
console.log(`\nComplete: ${solved} solved, ${failed} failed`);
});
stream.start(tasks);
Khi nào nên sử dụng Phát trực tuyến so với Thu thập tất cả
| Kịch bản | Cách tiếp cận |
|---|---|
| Gửi biểu mẫu bằng cách sử dụng mã thông báo | Luồng - gửi từng biểu mẫu ngay khi mã thông báo đến |
| Xuất CSV tất cả kết quả | Thu thập tất cả - viết một lần khi đợt hoàn thành |
| Bảng điều khiển với tiến trình trực tiếp | Phát trực tiếp - cập nhật giao diện người dùng trên mỗi sự kiện kết quả |
| Hàng loạt với sự phụ thuộc giữa các tác vụ | Thu thập tất cả – xử lý theo thứ tự sau khi hoàn thành |
| Lô lớn (1.000+) | Truyền phát - giảm mức sử dụng bộ nhớ cao điểm |
Khắc phục sự cố
| Vấn đề | Nguyên nhân | Cách xử lý |
|---|---|---|
| Kết quả đến theo thứ tự ngẫu nhiên | Bình thường - phát trực tuyến mang lại kết quả nhanh nhất trước tiên | Sử dụng result.index để ánh xạ trở lại tác vụ ban đầu |
| Bộ nhớ vẫn tăng trong khi truyền phát | Lưu trữ tất cả kết quả trong mảng | Xử lý và loại bỏ kết quả trong trình xử lý |
| Kết quả đầu tiên mất quá nhiều thời gian | Tất cả các nhiệm vụ được gửi đồng thời | Gửi xen kẽ với giới hạn semaphore hoặc đồng thời |
| Cảnh báo EventEmitter: MaxListenersExceeded | Quá nhiều người nghe trên luồng | Sử dụng setMaxListeners() hoặc đảm bảo một người nghe cho mỗi loại sự kiện |
| Trình tạo Async bị treo | Nhiệm vụ chưa được giải quyết trong nhóm đang chờ xử lý | Thêm thời gian chờ vào poll_task; đảm bảo tất cả các hợp đồng tương lai hoàn thành hoặc có lỗi |
Câu hỏi thường gặp
Tính năng phát trực tuyến có làm tăng lệnh gọi API so với đợt không?
Không - theo cả hai cách, số lượng lệnh gọi gửi và thăm dò ý kiến đều như nhau. Việc phát trực tuyến chỉ thay đổi khi ứng dụng của bạn xử lý từng kết quả chứ không phải số lượng lệnh gọi API được thực hiện.
Làm cách nào để duy trì thứ tự nhiệm vụ khi phát trực tuyến?
Mỗi kết quả đều mang index ban đầu của nó. Nếu thứ tự quan trọng đối với quá trình xử lý xuôi dòng, bộ đệm sẽ tạo ra cấu trúc được sắp xếp và thực hiện các lần chạy liền kề (như tập hợp lại gói TCP).
Tôi có thể kết hợp phát trực tuyến với điểm kiểm tra không?
Vâng. Nối từng kết quả vào một tệp điểm kiểm tra khi nó đến. Khi tiếp tục, tải điểm kiểm tra, lọc ra các chỉ mục đã hoàn thành và chỉ xử lý lại các tác vụ còn lại.
bài viết liên quan
Các bước tiếp theo
Xử lý các giải pháp CAPTCHA ngay khi chúng xuất hiện —lấy khóa API CaptchaAI của bạnvà xây dựng đường ống truyền phát.
Hướng dẫn liên quan: