微服务爬虫架构设计:解耦采集/解析/存储,支持百万级数据并发
2023年我所在的公司要做一个公共数据采集平台,需要对接300多个数据源,每天采集100万+条数据,支持多租户使用,最开始我用的是传统的Scrapy-Redis单体架构,跑了不到一个月就出现了各种问题:模块耦合严重,改一个功能影响其他功能;扩容困难,采集和解析耦合在一起,没办法单独扩容;故障影响范围大,一个数据源出问题导致整个平台都挂了。
后来我花了2个月时间把整个架构改成了微服务架构,四层解耦,K8s容器化部署,支持弹性扩缩容,现在平台每天可以稳定采集200万+条数据,支持100+租户同时使用,资源利用率提升了60%,成本降低了40%,稳定性达到了99.9%。
这篇文章我就把完整的微服务爬虫架构设计分享出来,从架构设计、每个服务的实现、K8s部署到性能优化,全部是实战经验,适合中大型爬虫项目,看完你也能搭建出支持百万级并发的爬虫平台。
一、传统单体爬虫架构的痛点
传统的Scrapy-Redis单体架构虽然简单,但是对于中大型项目来说,有很多无法解决的痛点:
- 模块耦合严重:采集、解析、存储逻辑都在同一个项目里,改一个功能可能影响其他功能,维护难度大
- 无法单独扩容:采集是IO密集型,解析是CPU密集型,耦合在一起的话,要么浪费CPU资源,要么浪费IO资源,没办法根据需求单独扩容
- 故障影响范围大:一个数据源的爬虫出问题,可能导致整个进程崩溃,影响所有数据源的爬取
- 迭代速度慢:不同的数据源开发进度不一样,耦合在一起的话,上线的时候要一起上线,影响迭代速度
- 多租户支持差:没办法隔离不同租户的任务和资源,容易出现资源竞争的情况
- 资源利用率低:峰谷明显,高峰期资源不够用,低峰期资源浪费严重
这些痛点在项目小的时候不明显,当数据源超过100个,日采集量超过10万条的时候,就会非常突出,严重影响项目的发展。
二、微服务爬虫架构核心设计
我设计的微服务爬虫架构是四层解耦的云原生架构,完全解决了单体架构的痛点:
┌─────────────────────────────────────────────────────────────────┐ │ API 网关层 │ │ 负责权限认证、流量控制、请求路由、负载均衡 │ └───────────────────┬─────────────────────────┬───────────────────┘ │ │ ┌───────────────────▼───────────────────┐ ┌───▼───────────────────┐ │ 调度服务 │ │ 管理后台 │ │ 负责任务调度、分片、优先级控制、超时管理 │ │ 任务管理、数据查看、监控 │ └───────────────────┬───────────────────┘ └───────────────────────┘ │ ┌───────────────────▼───────────────────┐ │ 消息队列 │ │ Kafka/RocketMQ,任务分发、削峰填谷、解耦 │ └─────────┬───────────────┬─────────────┘ │ │ ┌─────────▼───────┐ ┌─────▼─────────┐ ┌──────────────┐ │ 采集服务集群 │ │ 解析服务集群 │ │ 存储服务集群 │ │ IO密集型,多副本 │ │ CPU密集型,多副本│ │ 多数据源适配 │ └─────────────────┘ └───────────────┘ └──────────────┘整个架构的核心特点是完全解耦、弹性伸缩、高可用、易扩展,每个服务都可以独立开发、独立部署、独立扩容,互不影响。
2.1 API网关层
API网关是整个平台的入口,负责:
- 权限认证:所有请求都要经过网关认证,校验租户的API Key是否有效,权限是否足够
- 流量控制:对每个租户的请求做限流、熔断,避免单个租户的请求量太大影响其他租户
- 请求路由:把不同的请求路由到对应的服务,比如任务管理的请求路由到调度服务,数据查询的请求路由到存储服务
- 负载均衡:把请求均匀的分发到各个服务的多个副本上,提高可用性
- 日志收集:收集所有请求的日志,方便排查问题和统计调用量
我用的是APISIX做API网关,性能非常高,每秒可以处理几万次请求,完全满足需求。
2.2 调度服务
调度服务是整个平台的大脑,负责:
- 任务管理:接收用户提交的爬取任务,管理任务的生命周期,包括待调度、运行中、完成、失败等状态
- 任务分片:把大的任务拆分成多个小的子任务,比如爬取100万条数据,拆成1000个任务,每个任务爬1000条,分发到不同的采集节点
- 优先级控制:支持任务优先级,高优先级的任务先执行,比如实时爬取的任务优先级高于离线爬取的任务
- 超时管理:监控任务的执行时间,超过超时时间的任务自动重新调度,避免任务卡住
- 失败重试:任务执行失败的时候自动重试,最多重试3次,重试失败的任务进入死信队列,通知人工处理
- 增量爬取:支持增量爬取,自动判断哪些数据需要更新,不需要全量重新爬
调度服务是无状态的,可以部署多个副本,提高可用性,我用的是Go语言开发,性能非常高,每秒可以调度上万个任务。
2.3 消息队列层
消息队列是整个架构的核心,用来解耦各个服务,削峰填谷:
- 解耦:采集服务、解析服务、存储服务之间通过消息队列通信,不需要互相调用,修改一个服务不会影响其他服务
- 削峰填谷:高峰期任务量很大的时候,消息队列可以缓存任务,避免服务被打垮,低峰期的时候慢慢消费,提高资源利用率
- 保证任务不丢失:消息队列支持持久化,就算服务挂了,消息也不会丢失,重启之后可以继续消费
- 顺序消费:支持任务的顺序消费,比如同一个网站的任务按顺序执行,避免请求太频繁被封
我用的是Kafka集群,3节点,多副本,吞吐量非常高,每秒可以处理几十万条消息,完全满足百万级并发的需求。
2.4 采集服务层
采集服务是IO密集型服务,负责从各个数据源采集数据:
- 多数据源适配:支持网页、API、数据库、FTP等多种数据源的采集,每种数据源对应一个采集插件,新的数据源只要开发对应的插件就可以,不需要修改核心代码
- 无状态设计:采集服务是无状态的,不存任何本地数据,可以随意扩容缩容
- 反爬对抗内置:内置代理池、UA池、指纹混淆、验证码识别等反爬对抗功能,所有采集任务都可以复用这些功能,不需要重复开发
- 可观测性:每个采集任务的状态、耗时、成功率都有监控,方便排查问题
采集服务我用的是Python+Golang混合开发,简单的采集任务用Python开发,效率高,性能要求高的任务用Golang开发,速度快。
2.5 解析服务层
解析服务是CPU密集型服务,负责把采集到的原始数据解析成结构化数据:
- 多解析方式支持:支持XPath、CSS、正则、大模型等多种解析方式,用户可以根据需求选择
- 解析模板管理:用户可以自定义解析模板,不需要写代码,通过可视化配置就能解析数据
- 自动适配改版:内置大模型解析,网页改版之后不需要修改模板,自动适配新的结构
- 数据清洗:内置数据清洗、去重、标准化功能,解析之后直接输出干净的结构化数据
解析服务是无状态的,可以根据负载自动扩容,高峰期可以扩容到几十个副本,低峰期缩减到几个副本,节省成本。
2.6 存储服务层
存储服务负责把解析好的结构化数据存储到对应的存储介质里:
- 多存储介质适配:支持MySQL、PostgreSQL、MongoDB、Elasticsearch、OSS等多种存储介质,用户可以根据需求选择
- 数据同步:支持数据同步到其他系统,比如用户的数据库、数仓、大数据平台
- 查询接口:提供统一的查询接口,用户可以通过API查询采集到的数据
- 数据质量管理:内置数据质量校验,错误的数据自动告警,保证数据的准确性
2.7 管理后台
管理后台给用户提供可视化的操作界面:
- 任务管理:用户可以在后台创建、修改、删除爬取任务,查看任务的运行状态
- 数据查询:用户可以查询、导出采集到的数据
- 解析模板配置:可视化配置解析模板,不需要写代码
- 监控面板:查看任务的运行情况、资源使用情况、成功率、耗时等指标
- 租户管理:管理员可以管理租户,分配权限,设置配额
三、核心功能实现
我把几个核心功能的实现思路分享出来,大家可以参考。
3.1 任务调度实现
调度服务用的是Celery+Redis的分布式任务调度,支持定时任务、一次性任务、周期任务:
fromceleryimportCeleryimportredis app=Celery('crawler',broker='redis://redis:6379/0',backend='redis://redis:6379/0')@app.task(bind=True,max_retries=3)defcrawl_task(self,task_id,url,config):"""爬取任务"""try:# 执行爬取任务result=crawler.crawl(url,config)# 把结果发送到解析队列send_to_kafka('parse_topic',{'task_id':task_id,'raw_data':result,'parse_config':config['parse_config']})return{'status':'success','task_id':task_id}exceptExceptionase:# 重试self.retry(exc=e,countdown=60)任务执行失败的时候自动重试,最多重试3次,重试失败的任务进入死信队列。
3.2 弹性扩缩容实现
我用的是K8s的HPA(Horizontal Pod Autoscaler)实现自动扩缩容,根据服务的CPU、内存使用率或者队列的消息积压量自动调整副本数:
apiVersion:autoscaling/v2kind:HorizontalPodAutoscalermetadata:name:crawler-hpaspec:scaleTargetRef:apiVersion:apps/v1kind:Deploymentname:crawler-serviceminReplicas:2maxReplicas:50metrics:-type:Resourceresource:name:cputarget:type:UtilizationaverageUtilization:70-type:Podspods:metric:name:kafka_consumergroup_lagtarget:type:AverageValueaverageValue:1000这个配置表示:当CPU使用率超过70%,或者Kafka消费组的滞后量超过1000的时候,自动扩容采集服务的副本数,最多扩容到50个,低峰期的时候自动缩容到2个,大大提高了资源利用率,降低了成本。
3.3 多租户隔离实现
为了避免不同租户的任务互相影响,我做了三层隔离:
- 资源隔离:每个租户的任务有独立的资源配额,CPU、内存、请求量都有限制,不会占用其他租户的资源
- 队列隔离:每个租户的任务发到独立的Kafka Topic,不会和其他租户的任务混在一起
- 数据隔离:每个租户的数据存在独立的数据库Schema或者表,互相不可见,保证数据安全
这样就算某个租户的任务出问题,也不会影响其他租户的正常使用。
3.4 可观测性实现
整个架构的所有服务都有完善的监控:
- 指标监控:用Prometheus采集所有服务的CPU、内存、QPS、延迟、成功率等指标,用Grafana做可视化面板
- 日志收集:用ELK stack收集所有服务的日志,统一存储,方便排查问题
- 链路追踪:用Jaeger做分布式链路追踪,一个请求从网关到调度服务到采集服务到解析服务,整个链路都可以追踪,出了问题可以快速定位到是哪个服务出了问题
- 告警:用Alertmanager做告警,指标异常的时候自动发送钉钉、邮件通知,比如CPU使用率过高、Kafka消息积压、任务成功率过低等
四、百万级日采集量性能测试
我对整个架构做了压力测试,测试结果如下:
| 指标 | 测试结果 |
|---|---|
| 日最大采集量 | 320万条 |
| 峰值QPS | 3500次/秒 |
| 任务调度延迟 | <50ms |
| 采集成功率 | 98.5% |
| 解析成功率 | 97% |
| 平均响应时间 | 200ms |
| 服务可用性 | 99.92% |
测试的时候,采集服务的副本数从2个自动扩容到35个,解析服务从1个扩容到15个,测试完成之后自动缩容到原来的数量,整个过程不需要人工干预,非常方便。
成本对比
和原来的单体架构相比,微服务架构的成本降低了40%:
- 原来的单体架构需要20台4核8G的服务器,每个月成本6000元
- 现在的微服务架构,高峰期最多用到35个采集副本+15个解析副本,但是大部分时间只有5+2个副本,每个月成本3600元,降低了40%
- 而且性能提升了3倍,从原来的每天100万条提升到320万条
五、成本优化方案
微服务架构虽然性能高,但是如果不优化的话,成本也会很高,我总结了几个成本优化的技巧:
5.1 弹性扩缩容
这是成本优化的核心,根据实际的负载自动调整副本数,高峰期扩容,低峰期缩容,我用了弹性扩缩容之后,成本直接降了一半。
- 对于非实时的任务,可以安排在晚上或者凌晨运行,这时候服务器的成本更低,还可以利用闲置资源
- 用抢占式实例,价格是普通实例的1/3,就算被回收了也没关系,因为任务是无状态的,重新调度就可以了
5.2 冷热数据分离
- 热数据(最近7天的数据)存在性能好的SSD盘,查询速度快
- 冷数据(超过7天的数据)存在便宜的对象存储或者归档存储,成本只有SSD的1/10
- 定期把冷数据备份到对象存储,节省磁盘成本
5.3 资源复用
- 多个服务可以部署在同一个节点上,只要资源足够,提高资源利用率
- 用K8s的资源配额和限制,避免服务占用过多的资源
- 对于非核心的服务,可以用低配置的节点,降低成本
5.4 开源组件优先
所有的组件都用开源的,不需要付费license:
- API网关用APISIX,比Kong便宜,性能更好
- 消息队列用Kafka,开源免费,性能高
- 监控用Prometheus+Grafana,完全免费
- 存储用MySQL、MongoDB、Elasticsearch,都是开源的
这样一年可以节省十几万的软件 license 费用。
六、踩坑经验总结
微服务爬虫架构虽然好处很多,但是开发和运维的复杂度也比单体架构高很多,我踩了很多坑,分享给大家:
坑1:服务之间的调用关系复杂,排查问题困难
最开始没有做链路追踪,出了问题不知道是哪个服务的问题,排查非常困难。
解决方案:
- 加入分布式链路追踪,每个请求都有唯一的trace ID,贯穿整个调用链路,出了问题可以直接用trace ID查到整个链路的日志
- 完善监控,每个服务的指标都采集,有异常立刻告警
- 规范日志格式,所有服务的日志都包含trace ID、用户ID、任务ID等关键字段,方便排查
坑2:消息队列积压,服务被打垮
最开始没有做限流和熔断,高峰期任务量太大,解析服务处理不过来,消息队列积压了几百万条消息,导致整个集群都卡住了。
解决方案:
- 加入限流熔断机制,每个服务的消费速度超过阈值的时候,自动限制生产速度,避免消息积压
- 动态扩缩容,消息队列积压的时候自动扩容消费服务的副本数,加快消费速度
- 消息设置过期时间,超过24小时还没消费的消息自动丢弃,避免无限积压
坑3:分布式事务问题,数据不一致
最开始没有考虑分布式事务的问题,采集成功了,但是解析或者存储失败了,导致数据丢失,数据不一致。
解决方案:
- 用消息队列的ACK机制,只有消息处理成功了才提交ACK,失败的话会重新消费
- 加入对账机制,每天统计采集的任务数、解析成功的数量、存储成功的数量,对账不一致的时候自动补全
- 重要的数据做双写,同时存到两个存储介质,避免数据丢失
坑4:运维复杂度高,需要专门的运维人员
微服务架构的服务很多,部署、监控、运维都比单体架构复杂很多,最开始需要专门的运维人员维护,成本很高。
解决方案:
- 用K8s做容器编排,自动部署、自动扩容、自动恢复,大大降低运维复杂度
- 用DevOps流程,代码提交之后自动构建、自动测试、自动部署,不需要人工干预
- 做完善的监控和告警,大部分问题可以自动恢复,不需要人工处理
七、考点/技巧提炼
1. 什么样的项目适合用微服务爬虫架构?
- 中大型爬虫项目,数据源超过50个,日采集量超过10万条
- 需要支持多租户,多个团队或者用户使用
- 迭代速度快,需要频繁上线新功能
- 对可用性要求高,不能经常出故障
- 峰谷明显,需要弹性扩缩容降低成本
如果是小型的爬虫项目,数据源少,采集量小,用单体架构就足够了,不需要用微服务,增加复杂度。
2. 微服务爬虫架构的优势是什么?
- 解耦:各个服务独立开发、独立部署、独立运维,互不影响
- 弹性扩容:每个服务可以根据负载单独扩容,资源利用率高,成本低
- 高可用:没有单点故障,一个服务挂了不会影响其他服务
- 易扩展:新的数据源、新的功能只要开发对应的插件就可以,不需要修改核心代码
- 多租户支持:可以隔离不同租户的资源和数据,适合做SaaS化的爬虫平台
3. 微服务架构是不是比单体架构性能高?
不一定,如果是小型项目,微服务架构因为有网络调用的开销,性能反而不如单体架构;但是对于中大型项目,微服务架构可以水平扩展,性能可以做到比单体架构高很多。
4. 微服务爬虫架构的技术选型需要注意什么?
- 优先选云原生的组件,支持容器化部署,方便在K8s上运行
- 优先选开源、社区活跃的组件,遇到问题容易解决,不需要自己踩坑
- 各个服务之间的通信协议要统一,比如用HTTP或者gRPC,方便调用
- 要考虑整个团队的技术栈,不要选团队不熟悉的技术,增加学习成本
总结
微服务爬虫架构是中大型爬虫项目的最佳实践,虽然开发和运维的复杂度比单体架构高,但是带来的好处也是非常明显的,解耦、弹性扩缩容、高可用、易扩展,可以支撑百万级甚至千万级的日采集量。
我这个架构已经在公司的生产环境跑了一年多了,非常稳定,支撑了几十个业务方的爬虫需求,大大降低了开发成本,提高了资源利用率。当然微服务架构也不是银弹,要根据自己的项目规模和需求选择合适的架构,不要为了用微服务而用微服务,适合自己的才是最好的。
