爬虫集群把验证码任务量做到每小时数千次后,普通队列常在两处翻车:进程重启丢任务,消费者宕机无法重放。Apache Kafka 用持久化、分区有序、高吞吐的消息流解决这两点,把任务提交和识别结果处理彻底解耦。
架构设计
整体数据流如下:
[Scrapers] → Produce → [Kafka: captcha-tasks topic]
↓
[CAPTCHA Worker Group]
(consume tasks, solve via CaptchaAI)
↓
Produce → [Kafka: captcha-results topic]
↓
[Result Consumer Group]
(process solutions, update database)
| 主题 | 作用 |
|---|---|
captcha-tasks |
等待识别的验证码参数 |
captcha-results |
已识别完成、可供下游消费的 token |
什么时候真的需要 Kafka: 同步抓几十个页面,Redis 列表或线程池就够;到每小时上万次(比如大促比价爬虫),普通队列易丢任务、难扩容——正是 Kafka 分区 + 消费者组的强项。
环境准备
先装好依赖:
# Python
pip install kafka-python requests
# Node.js
npm install kafkajs axios
还需要可访问的 Kafka broker,默认 localhost:9092,生产环境换成集群地址。
国内网络建议加
-i指定镜像源(如清华 TUNA)加速安装。
步骤 1:创建 Kafka 主题
先建两个主题,一个装任务,一个装结果:
kafka-topics.sh --create --topic captcha-tasks \
--partitions 6 --replication-factor 1 \
--bootstrap-server localhost:9092
kafka-topics.sh --create --topic captcha-results \
--partitions 6 --replication-factor 1 \
--bootstrap-server localhost:9092
六个分区意味着同一消费者组最多六个 worker 并行消费。
步骤 2:编写任务生产者(爬虫端)
爬虫抓到验证码参数后直接扔进 captcha-tasks,不用等 CaptchaAI 处理完再抓下一页:
Python
import json
from kafka import KafkaProducer
producer = KafkaProducer(
bootstrap_servers=["localhost:9092"],
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
key_serializer=lambda k: k.encode("utf-8") if k else None,
acks="all", # Wait for all replicas to confirm
retries=3
)
def enqueue_captcha(task_id, sitekey, pageurl, captcha_type="userrecaptcha"):
"""Send a CAPTCHA task to Kafka."""
task = {
"task_id": task_id,
"method": captcha_type,
"sitekey": sitekey,
"pageurl": pageurl,
"submitted_at": __import__("time").time()
}
future = producer.send(
"captcha-tasks",
key=task_id, # Key ensures same task goes to same partition
value=task
)
future.get(timeout=10) # Block until confirmed
return task_id
# Submit tasks
enqueue_captcha("task_001", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
enqueue_captcha("task_002", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
producer.flush()
Node.js 版:
JavaScript
const { Kafka } = require("kafkajs");
const kafka = new Kafka({
clientId: "captcha-producer",
brokers: ["localhost:9092"],
});
const producer = kafka.producer();
async function enqueueCaptcha(taskId, sitekey, pageurl) {
await producer.connect();
const task = {
task_id: taskId,
method: "userrecaptcha",
sitekey: sitekey,
pageurl: pageurl,
submitted_at: Date.now(),
};
await producer.send({
topic: "captcha-tasks",
messages: [{ key: taskId, value: JSON.stringify(task) }],
});
}
(async () => {
await enqueueCaptcha(
"task_001",
"6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
"https://example.com"
);
await producer.disconnect();
})();
步骤 3:编写 CAPTCHA Worker(消费 + 识别)
worker 消费任务、调用 CaptchaAI 求解,写回 captcha-results,offset 手动提交:
Python
import json
import os
import time
import requests
from kafka import KafkaConsumer, KafkaProducer
API_KEY = os.environ["CAPTCHAAI_API_KEY"]
consumer = KafkaConsumer(
"captcha-tasks",
bootstrap_servers=["localhost:9092"],
group_id="captcha-workers",
value_deserializer=lambda m: json.loads(m.decode("utf-8")),
auto_offset_reset="earliest",
enable_auto_commit=False, # Manual commit after processing
max_poll_records=10
)
result_producer = KafkaProducer(
bootstrap_servers=["localhost:9092"],
value_serializer=lambda v: json.dumps(v).encode("utf-8")
)
def solve_captcha(task):
"""Submit to CaptchaAI and poll for result."""
# Submit
resp = requests.post("https://ocr.captchaai.com/in.php", data={
"key": API_KEY,
"method": task["method"],
"googlekey": task["sitekey"],
"pageurl": task["pageurl"],
"json": 1
})
data = resp.json()
if data.get("status") != 1:
return {"error": data.get("request")}
captcha_id = data["request"]
# Poll for result
for _ in range(60):
time.sleep(5)
result = requests.get("https://ocr.captchaai.com/res.php", params={
"key": API_KEY,
"action": "get",
"id": captcha_id,
"json": 1
}).json()
if result.get("status") == 1:
return {"solution": result["request"]}
if result.get("request") != "CAPCHA_NOT_READY":
return {"error": result.get("request")}
return {"error": "TIMEOUT"}
# Main consumer loop
print("CAPTCHA worker started. Waiting for tasks...")
for message in consumer:
task = message.value
print(f"Processing {task['task_id']}...")
result = solve_captcha(task)
result["task_id"] = task["task_id"]
result["solved_at"] = time.time()
# Publish result
result_producer.send("captcha-results", value=result)
result_producer.flush()
# Commit offset after successful processing
consumer.commit()
print(f" → {task['task_id']}: {'solved' if 'solution' in result else result.get('error')}")
Node.js 版:
JavaScript
const { Kafka } = require("kafkajs");
const axios = require("axios");
const API_KEY = process.env.CAPTCHAAI_API_KEY;
const kafka = new Kafka({
clientId: "captcha-worker",
brokers: ["localhost:9092"],
});
const consumer = kafka.consumer({ groupId: "captcha-workers" });
const producer = kafka.producer();
function sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
async function solveCaptcha(task) {
const submitResp = await axios.post(
"https://ocr.captchaai.com/in.php",
null,
{
params: {
key: API_KEY,
method: task.method,
googlekey: task.sitekey,
pageurl: task.pageurl,
json: 1,
},
}
);
if (submitResp.data.status !== 1) {
return { error: submitResp.data.request };
}
const captchaId = submitResp.data.request;
for (let i = 0; i < 60; i++) {
await sleep(5000);
const result = await axios.get("https://ocr.captchaai.com/res.php", {
params: { key: API_KEY, action: "get", id: captchaId, json: 1 },
});
if (result.data.status === 1) return { solution: result.data.request };
if (result.data.request !== "CAPCHA_NOT_READY")
return { error: result.data.request };
}
return { error: "TIMEOUT" };
}
async function run() {
await consumer.connect();
await producer.connect();
await consumer.subscribe({ topic: "captcha-tasks", fromBeginning: false });
await consumer.run({
eachMessage: async ({ message }) => {
const task = JSON.parse(message.value.toString());
console.log(`Processing ${task.task_id}...`);
const result = await solveCaptcha(task);
result.task_id = task.task_id;
result.solved_at = Date.now();
await producer.send({
topic: "captcha-results",
messages: [{ value: JSON.stringify(result) }],
});
console.log(
` → ${task.task_id}: ${result.solution ? "solved" : result.error}`
);
},
});
}
run();
扩容 Worker
Kafka 消费者组自动在 worker 间分配分区,加 worker 即触发一次重新平衡:
# 6 partitions, 3 workers → each worker gets 2 partitions
Worker-1: partitions 0, 1
Worker-2: partitions 2, 3
Worker-3: partitions 4, 5
# Add Worker-4 → rebalance
Worker-1: partitions 0, 1
Worker-2: partitions 2
Worker-3: partitions 3, 4
Worker-4: partition 5
扩容天花板就是分区数——想要更多 worker,先加分区。
生产者数据校验规则
- 缺求解器类型或路由信息的消息,写入前直接拒绝。
- 加幂等键,避免重试产生重复任务。
- 格式错误的消息转入死信路径,方便回放排查。
监控与关键指标
用 consumer lag 判断 worker 是否跟得上任务产出:
kafka-consumer-groups.sh --describe --group captcha-workers \
--bootstrap-server localhost:9092
| 指标 | 健康区间 | 预警信号 |
|---|---|---|
| Consumer lag | < 100 | > 1000(加 worker) |
| 入方向 msg/sec | 匹配爬虫速率 | 突增=采集高峰 |
| 出方向 msg/sec | 匹配入方向 | 落后=处理瓶颈 |
lag 长期偏高先加 worker,再查分区数够不够。
常见故障排查
常踩的四个坑:
| 问题 | 原因 | 处理方式 |
|---|---|---|
| Consumer lag 增长 | worker 跟不上产出 | 加 worker(上限是分区数) |
| 结果重复 | worker 提交前崩溃 | 按 task_id 做幂等校验 |
| 频繁 rebalance | worker 崩溃/重启 | 调大 session.timeout.ms;查 OOM |
| 任务分配不均 | key 选得不好 | 换随机 key 或加分区 |
常见问题
整理了四个高频问题:
什么时候该用 Kafka,而不是 Redis 队列?
量小逻辑简单用 Redis/RabbitMQ 更轻量;到每小时数万次、需要重放和消费者组扩容才值得上 Kafka。
Kafka 分区数要怎么定?
- 分区数就是并行度上限:6 分区最多 6 个 worker 同时消费。
- 按目标 QPS 估算并留余量——分区只能加不能减。
Consumer lag 一直涨,要立刻扩容吗?
短暂高峰后自行追平属正常抖动,持续超 1000 才需加 worker;同时看是任务量变大还是识别耗时变长。
CaptchaAI 请求失败要不要在 worker 里重试?
- 网络超时、5xx:重试几次即可。
- sitekey/pageurl 参数错误:重试无意义,直接标记错误发到
captcha-results或死信主题。
相关文章
- 流式批量验证码结果处理
下一步
现在搭建你的流式验证码管道——注册 CaptchaAI 获取 API Key,接入 Kafka 稳定跑通数千次识别任务。
相关指南
- 用于分布式处理的 Redis 队列
- RabbitMQ 消息队列集成
- 每小时解决 10,000 个任务