ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

多源生活服务数据难聚合?Scrapy采集架构设计与落地全实战

多源生活服务数据难聚合?Scrapy采集架构设计与落地全实战 做过本地生活运营或者数据分析的朋友大概率碰过这个难题商家信息、团购套餐、用户评分散在各大平台格式不统一、更新还快靠人工整理效率低到离谱写几个单页脚本凑合用平台一改版就全崩。去年我给一家本地商业运营公司做数据支撑系统的时候就踩过这个坑。当时客户需要每天同步本地三个主流生活服务平台的全品类商家数据用来做竞品分析和商圈调研。最开始团队里的实习生写了一堆零散的采集脚本每天跑十几个小时还经常因为平台前端升级集体失效数据错漏能占到三成。后来我花了三周时间基于Scrapy重构了整套架构统一了多源数据的接入标准做了全链路的反爬适配和清洗流程上线后稳定运行了大半年采集效率提升了四倍日常维护只需要偶尔加个字段、调下规则。这篇文章就把整套架构的设计思路、核心模块实现和踩过的坑全部整理出来从项目结构到代码片段再到问题排查全是线上跑过的实战内容。一、本地化生活服务数据采集的核心痛点本地化生活服务的数据和普通资讯类数据完全不是一个维护难度核心问题集中在四点第一是数据源异构性极强。不同平台的页面结构、字段命名、数据格式千差万别同样是“门店地址”有的拆成省市区三级有的是一整段文本有的还带特殊符号统一处理的成本很高。第二是数据更新频率高。商家的营业状态、团购价格、评分每天都在变需要高频次的增量更新全量采集既慢又容易触发反爬。第三是反爬机制多样。主流生活服务平台都有完善的反爬体系IP封禁、UA校验、Cookie追踪、滑块验证应有尽有单靠代理池根本扛不住。第四是数据质量参差不齐。公开平台上的信息很多存在重复、缺失、错误的情况比如同一家门店在不同平台的名称、电话不一致如果不做清洗直接入库后续根本没法用。二、整体架构设计这套架构的核心思路是分层解耦、可插拔扩展每一层都可以独立替换和升级不用因为某个平台改版就动整个系统。整体分为五层从上到下分别是调度层、采集层、代理层、清洗层、存储层。调度层负责管理采集任务、控制采集频率、分配分布式节点基于Redis实现任务队列和去重集支持全量和增量两种采集模式。采集层基于Scrapy实现封装了统一的请求模板和解析规则每个数据源对应一个Spider字段映射统一到标准Item解决异构数据的接入问题。代理层独立的代理池服务支持多家代理供应商自动切换带健康检查和失败重试保证请求的成功率。清洗层也就是Pipeline处理链路依次做格式标准化、重复数据去重、内容合法性校验、缺失值补全输出统一格式的结构化数据。存储层结构化数据存入MySQL非结构化的详情页快照存入MongoDB已采集的URL和内容指纹存入Redis做去重。这种分层结构的好处是扩展性极强后面要加新的数据源只需要写一个新的Spider和对应的字段映射其他层完全不用动反爬策略升级也只需要改中间件不影响业务逻辑。三、前期准备与技术选型落地这套架构不需要太复杂的技术栈核心组件都是Python生态里成熟的方案采集框架Scrapy 2.11生态完善中间件和Pipeline机制非常适合做结构化采集分布式支持Scrapy-Redis用Redis做任务队列和去重快速实现分布式部署数据库MySQL 8.0存结构化业务数据Redis 7.0做缓存和去重MongoDB存原始页面快照辅助工具fake-useragent生成UArequests做代理健康检查Apscheduler做定时任务环境准备上建议Python版本用3.10以上兼容性和性能都比较均衡。依赖包不用一次性装全用到哪个装哪个就行。四、核心模块分步实现1. 项目结构与标准Item定义首先搭好Scrapy项目的基础结构在原生结构基础上扩展了字段映射、规则配置和工具类模块方便后续维护。我最开始图省事把所有逻辑混在一起后来加规则的时候越写越乱拆成模块化之后维护成本降了很多。项目目录大概是这样的life_data_collect/ ├── spiders/ # 各平台采集逻辑 │ ├── meituan.py │ ├── dianping.py │ └── eleme.py ├── items.py # 统一数据结构定义 ├── middlewares.py # 请求中间件反爬、代理、UA ├── pipelines.py # 数据清洗与存储链路 ├── settings.py # 全局配置 └── utils/ # 工具类哈希、代理、清洗第一步先定义统一的标准Item所有平台的采集结果都映射到这个结构上从源头解决数据异构的问题。这一步非常关键不然后续清洗会非常麻烦。importscrapyclassLifeShopItem(scrapy.Item):# 基础标识shop_idscrapy.Field()# 平台内店铺IDsourcescrapy.Field()# 数据来源平台# 店铺基本信息shop_namescrapy.Field()# 店铺名称shop_typescrapy.Field()# 店铺品类addressscrapy.Field()# 详细地址provincescrapy.Field()# 省份cityscrapy.Field()# 城市districtscrapy.Field()# 区县phonescrapy.Field()# 联系电话# 运营信息ratingscrapy.Field()# 评分review_countscrapy.Field()# 评论数avg_pricescrapy.Field()# 人均价格business_hoursscrapy.Field()# 营业时间# 采集信息crawl_timescrapy.Field()# 采集时间戳content_hashscrapy.Field()# 内容哈希用于增量判断定义的时候要注意所有字段都用最细的粒度拆分比如地址拆成省市区详细地址不要图省事只存一个总字段不然后续统计和筛选会很麻烦。2. 采集Spider的通用化封装每个平台的页面结构不一样但请求流程、异常处理、字段映射的逻辑是共通的。所以先写一个基类Spider把通用逻辑封装进去每个具体平台的Spider只需要写解析规则和字段映射就行能省很多重复代码。基类里主要做三件事统一请求头、统一异常重试、统一字段预处理。具体平台的Spider继承这个基类只需要实现parse_shop_detail方法把页面里的数据提取出来赋值给标准Item。核心的基类片段如下importscrapyimporttimefromhashlibimportmd5from..itemsimportLifeShopItemclassBaseShopSpider(scrapy.Spider):allowed_domains[]defstart_requests(self):# 从Redis读取任务列表分布式场景下直接继承RedisSpider即可forurlinself.task_list:yieldscrapy.Request(urlurl,callbackself.parse_shop_detail,errbackself.error_back,meta{retry_times:0})deferror_back(self,failure):# 失败重试逻辑最多重试3次retry_timesfailure.request.meta.get(retry_times,0)ifretry_times3:retry_times1new_requestfailure.request.copy()new_request.meta[retry_times]retry_timesyieldnew_requestelse:self.logger.error(f请求失败超过3次{failure.request.url})defparse_shop_detail(self,response):# 子类实现具体解析逻辑raiseNotImplementedError(子类必须实现parse_shop_detail方法)defbuild_item(self,data):# 统一构建Item生成内容哈希itemLifeShopItem()fork,vindata.items():ifkinitem.fields:item[k]v.strip()ifisinstance(v,str)elsev# 生成内容哈希用于增量去重content_strf{item[shop_name]}{item[phone]}{item[address]}item[content_hash]md5(content_str.encode(utf-8)).hexdigest()item[crawl_time]int(time.time())returnitem有了这个基类新增一个平台的采集逻辑只需要几十行代码维护成本大大降低。3. 反爬中间件设计反爬是数据采集里最头疼的部分这套架构里把所有反爬逻辑都封装在中间件里和业务逻辑完全解耦升级反爬策略不用动Spider。核心做了四层反爬防护第一层是UA动态轮换每个请求随机换一个真实浏览器UA避免被识别成程序访问。第二层是代理池自动切换每个请求分配不同的代理IP失败自动换代理重试。第三层是请求频率控制按域名维度控制并发数和请求间隔避免短时间请求过多被封。第四层是Cookie池管理针对需要登录的平台维护一批有效Cookie随机分配使用。这里给出最常用的代理中间件核心实现importrandomfrom.utils.proxy_poolimportProxyPoolclassProxyMiddleware:def__init__(self):self.proxy_poolProxyPool()defprocess_request(self,request,spider):# 跳过不需要代理的请求ifrequest.meta.get(no_proxy):return# 从代理池获取一个可用代理proxyself.proxy_pool.get_proxy()ifproxy:request.meta[proxy]fhttp://{proxy}defprocess_response(self,request,response,spider):# 状态码异常标记代理不可用重新请求ifresponse.statusin[403,429,503]:proxyrequest.meta.get(proxy)ifproxy:self.proxy_pool.mark_bad_proxy(proxy)new_requestrequest.copy()new_request.dont_filterTruereturnnew_requestreturnresponse代理池建议单独做成一个服务不要和采集逻辑写在一起方便后续更换代理供应商和做健康检查。4. 数据清洗Pipeline链路采集到的原始数据不能直接入库必须经过清洗。Pipeline采用链路式设计每个Pipeline负责一项清洗工作按顺序执行方便调整和扩展。清洗链路一共分四步格式标准化统一日期、价格、电话号码的格式去除特殊字符和多余空格。内容校验过滤掉关键字段缺失的无效数据比如没有店铺名称、没有地址的记录。去重判断根据内容哈希判断是否和已采集的数据重复重复的就跳过入库。数据入库校验通过的数据写入MySQL原始页面快照写入MongoDB。去重Pipeline的核心代码如下importredisfromscrapy.exceptionsimportDropItemclassDedupPipeline:def__init__(self,redis_host,redis_port):self.redis_connredis.Redis(hostredis_host,portredis_port,db0)self.dedup_keylife_shop:dedup_hashclassmethoddeffrom_crawler(cls,crawler):returncls(redis_hostcrawler.settings.get(REDIS_HOST),redis_portcrawler.settings.get(REDIS_PORT))defprocess_item(self,item,spider):content_hashitem.get(content_hash)ifnotcontent_hash:returnitem# 判断哈希是否已存在ifself.redis_conn.sismember(self.dedup_key,content_hash):spider.logger.info(f数据重复跳过{item[shop_name]})raiseDropItem(重复数据)self.redis_conn.sadd(self.dedup_key,content_hash)returnitem这里用Redis的集合做去重性能很高千万级别的数据判断也能在毫秒级返回。如果数据量特别大可以换成布隆过滤器能省很多内存。5. 增量采集机制全量采集既费时间又容易触发反爬日常运行肯定要用增量采集。这套方案用URL去重内容哈希校验的双层机制保证只采集有更新的内容。具体逻辑是每次采集任务启动时先拉取需要更新的店铺URL列表和Redis里的已采集URL做对比只采集新出现或者过期的URL。采集到内容后计算内容哈希和之前的哈希对比如果一样就说明内容没更新跳过入库减少数据库写入。对于评分、评论数这类经常变的字段设置一个更新周期比如7天强制更新一次保证数据的时效性。这种机制下日常增量采集只需要全量采集1/5的时间对目标平台的压力也小很多不容易被限制访问。五、踩坑与问题排查这套架构前后迭代了三个版本踩了不少坑这里列几个最常见的大家可以直接避开。页面结构频繁改版导致解析失效这是做第三方数据采集最常见的问题平台前端一发版选择器就失效了。我的解决方案是不要写太细的层级选择器尽量用特征属性定位比如data-type、class里的业务关键字不要依赖具体的DOM层级。另外加一个解析失败告警解析成功率低于阈值的时候自动发通知及时发现问题。代理质量参差不齐成功率忽高忽低不要只靠一家代理供应商建议同时接两到三家做代理分级。高匿代理用来采核心数据普通代理用来采列表页。另外一定要加代理健康检查连续失败的代理自动剔除不要等到影响采集任务了才发现。分布式任务重复采集用Scrapy-Redis的时候如果节点多了很容易出现多个节点抢同一个任务的情况。解决方法是给任务加一个分布式锁拿到锁的节点才能处理任务处理完再释放。或者用Redis的有序集合做任务队列用zrem的原子性保证任务只被消费一次。数据去重不彻底脏数据入库只靠内容哈希去重是不够的比如同一家店铺电话换了一个数字哈希就不一样了。建议再加一层业务去重比如根据“店铺名称地址”做模糊匹配相似度超过阈值就判定为同一家店。当然这会增加计算量可以每天离线跑一次不用实时做。采集速度和反爬的平衡很多人一上来就把并发拉满结果跑十分钟就被全平台限制。建议先从低并发开始慢慢往上加找到一个既不被封、速度又能接受的平衡点。另外可以加自适应限速一旦出现429状态码自动降低并发数恢复正常后再慢慢提上来。六、性能优化建议如果数据量特别大还可以从这几个方向做优化过滤非必要请求图片、CSS、JS这些资源全部禁止加载只请求HTML页面能省一半以上的流量和时间。启用HTTP缓存对于不经常变的页面启用本地缓存重复采集的时候直接读缓存不用发请求。异步化存储入库操作改成异步不要阻塞采集线程用消息队列做削峰采集速度会提升很多。节点横向扩展采集层是无状态的可以随便加节点任务量上来了多加几台机器就行线性提升采集能力。七、总结本地化生活服务数据聚合本质上是一个“采-洗-存”的标准化工程。核心不是写多少复杂的采集逻辑而是搭建一个可扩展、易维护的架构把可变的业务规则和稳定的基础能力分开。Scrapy本身提供了非常好的插件化机制配合中间件和Pipeline可以很方便地把反爬、清洗、存储这些能力做成可插拔的组件。再加上分布式改造和增量机制完全可以支撑中大型的数据聚合需求。
返回列表