核心在于稳态调度、弹性扩展、任务不丢、请求不撞墙;需用brpop或kafka acks=all保障投递,etcd心跳保节点感知,分层限速控并发,标准化url去重,指数退避防重试风暴。

golang 实现分布式爬虫,核心不在“能不能分”,而在“怎么分得稳、扩得开、不丢任务、不撞墙”。直接上 Kafka + etcd + colly 的组合能跑,但多数人卡在调度不均、节点失联后任务堆积、重试逻辑和限速冲突这几个点上。
任务队列选 Redis 还是 Kafka?看失败容忍度
Redis 的 LPop/LPush 简单快,适合中小规模(日均百万 URL 以内)、允许少量任务丢失的场景;Kafka 则必须用,如果你要保障「至少一次」投递、需回溯重放、或已有消息中间件基建。
- 用 Redis:务必配合
BRPop阻塞式取任务,避免轮询空耗 CPU;任务入队前序列化成 JSON,并带timestamp和retry_count字段 - 用 Kafka:别直接用
sarama.SyncProducer发完就不管——它默认不等 ack,网络抖动时会静默丢消息;必须设"acks": "all"并检查Send返回的 error - 共性陷阱:没做任务去重。URL 入队前必须标准化(
url.Parse+url.URL.EscapedPath()+ 去 fragment),否则https://a.com/1和https://a.com/1#top被当两个任务
worker 节点注册与心跳怎么不漏判?
etcd 或 ZooKeeper 不是摆设——节点上线/下线必须靠它驱动调度器决策。单纯靠 TCP 连接存活判断不可靠,因为网络分区时连接可能假活。
- 每个 worker 启动时,在 etcd 写临时 key:
/workers/{hostname}:{pid},TTL 设为 15 秒 - worker 每 5 秒续租一次:
client.KeepAlive();断连后 etcd 自动删 key,调度器监听/workers/目录变更即可感知 - 别用固定 IP 做节点标识:容器或云主机 IP 经常变,改用 hostname + 启动时间戳生成唯一 ID,比如
fmt.Sprintf("%s-%d", hostname, time.Now().UnixMilli())
colly 在分布式里怎么避免并发失控?
colly 默认的 Async = true 只是开启 goroutine,真正控并发的是 LimitRule。但在多 worker 场景下,单个 collector 的限流只管本机,跨节点完全没约束。
- 必须分层限速:节点级用
rate.Limiter控总 QPS(比如每秒最多 20 个请求),域名级再用 colly 的LimitRule控并行数(如Parallelism: 2) - 别把
rate.NewLimiter放在 OnRequest 回调里创建——每次请求都 new 一个,令牌桶状态不共享;应全局复用一个*rate.Limiter实例 - 随机延迟比固定 Delay 更有效:
c.Limit(&colly.LimitRule{Delay: 1 * time.Second, RandomDelay: 500 * time.Millisecond}),防被识别为机器流量模式
任务失败后重试,为什么越重试越慢?
重试不是简单 Visit(url) 再来一遍。没控制退避策略和最大重试次数,会导致失败任务反复抢占资源,挤占正常任务带宽。
- 重试必须带指数退避:
time.Second * (1 ,且上限封顶(比如最多 3 次,第 3 次间隔 ≥ 8 秒) - 失败任务不要立刻回推队列头部——否则所有 worker 都抢着处理这个失败 URL;应推到独立的
retry_queue,并加 delay(如 60 秒后才可被消费) - 403/429 错误优先走 UA 轮换或代理切换,而不是无脑重试;这类错误大概率是反爬触发,重试只会加速被封
真实系统里最常被忽略的,是任务状态的跨节点可见性。比如某个 URL 已被 worker A 标记为「正在处理」,worker B 却因网络延迟没收到通知,又取了一次——这不是代码写错,而是没引入分布式锁或乐观更新机制。别指望 etcd 的单 key watch 能覆盖所有竞争路径,关键状态变更得配 version 字段 + CAS 操作。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











