单机爬虫的瓶颈与分布式队列的价值
R语言中虽然可以使用 parallel、foreach 等包实现多核并行抓取,但这种并行仍然局限于单台机器。假如一个新闻站点需要采集几十万条详情页,单机循环逐个请求,大量时间消耗在等待响应上,CPU反而空闲。即便开了四个核心同时请求,IP出口不变,目标站点一旦限频,所有任务都会受到影响。

更棘手的是任务状态管理。单机脚本在运行过程中如果中断,需要靠额外文件记录已经抓取到哪一页、哪些URL等待处理,恢复起来非常麻烦。把待抓取地址放进Redis队列后,任务数据可以从R进程转移出来,任何一台机器上的R脚本只要能连上Redis,就能领取下一个任务。这种生产者-消费者模型让任务分配从进程内逻辑升级为独立的中间件能力。
Redis在分布式任务队列中的优势很明显。它所有数据都存在内存中,读写延迟在亚毫秒级别,命令执行天然原子化,不会出现多个R进程同时拿到同一个URL的情况。哪怕在单机上跑四个R进程,Redis也能充当线程安全的任务池,替代复杂的共享内存方案。
Redis列表结构与任务入队消费
列表是Redis最基础的数据结构之一,常用来做FIFO队列。生产者使用 LPUSH 或 RPUSH 写入URL,消费者使用 RPOP 或 LPOP 读取并移除。R语言里可以用 redux 包连接Redis,下面是连接与入队示例:
library(redux)
conn = hiredis()
conn$RPUSH("crawl:queue", "https://ipipp.com/article/1")
conn$RPUSH("crawl:queue", "https://ipipp.com/article/2")
conn$LLEN("crawl:queue")
这里 RPUSH 表示从右侧插入,LLEN 返回队列长度。多个R进程同时执行 RPUSH 不会产生竞争,Redis会按顺序追加。消费端如果使用 LPOP,队列里没有任务时返回NULL,脚本需要自行轮询,这会造成CPU空转。更合适的做法是使用 BRPOP,它会在队列为空时阻塞等待,直到有新的任务或超时时间到达。
不过单纯使用列表只能保证任务不重复分发,无法直接做去重。如果多个生产者可能写入相同URL,可以在入队前用集合判断。下面代码展示了去重入队:
enqueue_url = function(conn, queue, url) {
is_new = conn$SADD("crawl:seen", url)
if (is_new == 1) {
conn$RPUSH(queue, url)
return(TRUE)
} else {
return(FALSE)
}
}
SADD 命令在成员已经存在时返回0,第一次加入返回1。利用这个特性可以避免重复抓取。但需要注意,集合会一直增长,如果URL数量非常大,需要定期清理或使用布隆过滤器降低内存占用。
可靠队列与失败重试方案
使用 BRPOP 消费时,任务在弹出后立即从队列中消失。如果R进程在处理过程中崩溃,这条URL就永久丢失了。对于需要保证抓取完整性的场景,应该引入可靠队列模式,把任务从等待队列转移到处理中列表,确认成功后再移除。
Redis提供了 RPOPLPUSH 命令,可以从源列表弹出元素并同时压入目标列表,整个过程原子执行。R脚本可以这样改造消费者:
task = conn$RPOPLPUSH("crawl:queue", "crawl:processing")
if (!is.null(task)) {
result = tryCatch(
crawl_page(task),
error = function(e) NULL
)
if (!is.null(result)) {
conn$LREM("crawl:processing", 1, task)
conn$RPUSH("crawl:done", task)
} else {
conn$LREM("crawl:processing", 1, task)
conn$RPUSH("crawl:failed", task)
}
}
这里 RPOPLPUSH 把任务放入 crawl:processing 列表,处理成功后从处理中列表删除,并加入完成列表。如果处理失败,任务进入失败列表,可以由后续脚本重试或人工检查。对于进程直接崩溃的情况,任务会滞留在 crawl:processing 中,需要一个定时任务检查该列表,把超时未完成的任务重新放回等待队列。
超时回收可以用带时间戳的有序集合实现。每个任务进入处理中列表时,同时把URL和当前时间写入有序集合 crawl:processing_zset。定时程序读取所有未完成记录,如果当前时间减去入队时间超过阈值,就执行 LREM 从处理中列表移除,并重新 RPUSH 到等待队列。这样即使worker崩溃,也能保证任务最终被再次领取。
反爬限流与结果落盘
分布式爬虫虽然提升了抓取速度,但对目标站点来说,多个worker同时请求可能触发风控。解决方案通常有两类:一是给每个worker配置独立代理,让外部看来请求来自不同出口;二是在Redis中记录每个域名的请求时间,实现共享限流。R语言使用 httr 包时可以通过 use_proxy 设置代理,代理池可以从Redis列表中循环读取。
滑动窗口限流的思路是,每次请求前在Redis中记录时间戳,然后统计最近一段时间内的请求次数。可以用Redis列表保存时间戳,每次请求前删除超时的旧记录,再检查长度是否超过阈值。代码如下:
rate_limit = function(conn, domain, limit, window) {
now = as.numeric(Sys.time())
key = paste0("rate:", domain)
conn$ZREMRANGEBYSCORE(key, "-inf", now - window)
count = conn$ZCARD(key)
if (count < limit) {
conn$ZADD(key, now, paste0(now, "-", sample.int(10000, 1)))
return(TRUE)
} else {
return(FALSE)
}
}
这段代码使用有序集合,score为时间戳,成员用随机后缀避免重复覆盖。返回TRUE表示允许请求,FALSE表示需要等待。多个R进程共享同一个Redis限流集合,就不会出现各自计时导致总请求量超限的问题。
抓取结果可以直接写入Redis,也可以落库到PostgreSQL或SQLite。分布式场景下,每个worker如果都直接写数据库,连接数可能成为瓶颈,所以往往会先把结果写入Redis的列表或哈希中,再由单独的入库进程批量消费。R中可以使用 redis 包或 redux 包把解析出来的数据序列化为JSON后 RPUSH 到结果队列。
完整示例:R语言分布式爬虫
下面给出一个可运行的简化版本。它包含一个生产者脚本,将种子URL写入队列,以及一个消费者脚本,循环地从Redis领取任务、抓取网页标题并写入结果列表。实际生产环境可以增加更多异常处理和重试逻辑。
library(redux)
library(httr)
library(rvest)
# 连接Redis
conn = hiredis()
# 种子URL入队
seed_urls = c(
"https://ipipp.com/page/1",
"https://ipipp.com/page/2",
"https://ipipp.com/page/3"
)
for (url in seed_urls) {
if (conn$SADD("crawl:seen", url) == 1) {
conn$RPUSH("crawl:queue", url)
}
}
# 消费任务并抓取标题
while (TRUE) {
task = conn$BRPOP("crawl:queue", 5)
if (is.null(task)) {
next
}
url = task[2]
page = tryCatch(
httr::GET(url, httr::timeout(10)),
error = function(e) NULL
)
if (!is.null(page)) {
if (httr::status_code(page) == 200) {
title = rvest::html_text(rvest::html_node(rvest::read_html(page), "title"))
conn$RPUSH("crawl:results", title)
conn$LREM("crawl:processing", 1, url)
} else {
conn$RPUSH("crawl:failed", url)
}
} else {
conn$RPUSH("crawl:failed", url)
}
Sys.sleep(1)
}
这个示例中 BRPOP 返回的是两元素向量,第一个元素是队列名,第二个是实际URL,所以用 task[2] 获取。抓取成功后标题被推入 crawl:results 列表,失败任务进入 crawl:failed 列表。生产者与消费者可以放在不同机器上运行,只要它们连接同一个Redis实例,就能自然地实现任务分配。
需要说明的是,这一示例假设所有worker都可以访问外网,并且目标站点允许抓取。真实项目中还应当遵守robots协议,设置合理的抓取间隔,并对响应内容做编码检测。Redis队列可以继续扩展为优先级队列、分片队列或多级队列,但核心思想始终一致:把任务状态外置,让执行进程无状态,这样分布式扩展才会简单可靠。