尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

亚太数据聚合平台实战:Scrapy-Redis分布式爬虫与ES数据清洗架构解析

亚太数据聚合平台实战:Scrapy-Redis分布式爬虫与ES数据清洗架构解析 1. 项目概述与核心价值最近在整理过往项目资料时翻到了一个2021年做的亚太地区数据收集网站项目。这个项目虽然过去几年了但其中的设计思路、技术选型和踩过的坑对于现在做类似区域化数据聚合或市场情报分析的朋友来说依然有很强的参考价值。简单来说这个网站的核心目标是构建一个能够自动化、结构化地收集、清洗并展示亚太地区特定行业公开数据的平台。它不是简单地爬取信息而是涉及多源异构数据处理、反爬策略应对、数据质量校验以及前端可视化呈现等一系列完整链条。当时做这个项目源于一个很实际的需求团队需要持续追踪亚太多个市场的动态但手动从各个国家不同的新闻网站、报告平台、公开数据库里找信息效率极低且数据格式混乱无法直接分析。市面上通用的数据工具要么太贵要么覆盖不全要么灵活性不够。于是我们就决定自己动手打造一个“专属情报站”。这个网站最终服务了内部的市场分析、竞品研究和战略规划其方法论也适用于很多需要对特定区域进行长期数据监测的场景比如跨境电商选品、区域经济研究、品牌舆情监控等。2. 项目整体架构与技术选型2.1 核心需求与设计思路拆解接到“亚太数据收集”这个需求首先要明确“数据”是什么以及“收集”的边界。我们聚焦的是公开的行业报告摘要、政策动态、市场规模数据和头部公司动向。目标网站覆盖了中、日、韩、东南亚、澳新等主要亚太经济体。设计思路上我们确立了几个原则稳定性优先数据收集是长期任务系统必须能7x24小时稳定运行具备一定的容错和自恢复能力。可扩展性数据源可能会增加或变更采集逻辑需要模块化便于维护和扩展。数据质量可控原始数据杂乱必须经过清洗、去重、结构化才能产生价值。合规与风险规避严格遵守各网站Robots协议控制请求频率避免对目标服务器造成压力这是红线。基于这些原则我们没有采用单一的“大爬虫”模式而是设计了一个微服务化的流水线架构。整个系统被拆分为调度中心、采集器集群、消息队列、清洗引擎、存储模块和前端应用六个部分各司其职通过API和消息进行松耦合通信。2.2 技术栈选型背后的考量技术选型是项目的地基每一个选择都经过了反复权衡。这里我详细说说当时的思考过程。后端与采集核心Python Scrapy RequestsPython是数据抓取和处理的绝对主力。我们以Scrapy框架为基础构建主要采集器因为它生态成熟异步处理能力强中间件机制灵活非常适合构建复杂的爬虫项目。但对于一些结构简单或反爬策略特殊的网站我们也会单独编写基于Requests BeautifulSoup的脚本作为Scrapy的补充。为什么不只用Scrapy因为有些时候“杀鸡不用牛刀”轻量脚本部署和调试更快。任务调度与协调Celery Redis定时采集和失败重试需要可靠的任务调度。Celery是Python领域分布式任务队列的事实标准搭配Redis作为Broker和结果后端实现了任务的异步执行、定时触发和状态监控。我们将每个数据源即一个目标网站的采集任务定义为一个Celery任务由调度中心统一管理。数据存储PostgreSQL Elasticsearch关系型数据库和搜索引擎的组合拳。PostgreSQL用于存储清洗后的结构化数据如文章标题、发布时间、来源、分类、摘要、关键指标如市场规模数字等利用其强大的JSONB字段也能灵活存储一些半结构化数据。Elasticsearch则用于全文检索和快速聚合分析前端复杂的筛选和关键词搜索都靠它支撑。这种“主存索引”的模式兼顾了数据管理的严谨性和查询性能的高效性。消息队列RabbitMQ在采集器和清洗引擎之间我们引入了RabbitMQ。采集器抓取到原始网页内容后并不直接处理而是包装成一个消息包含URL、原始HTML、采集时间等元数据投递到指定的队列。清洗引擎作为消费者从队列拉取消息进行处理。这样做的好处是解耦即使清洗引擎因升级或BUG暂时宕机采集到的数据也不会丢失会堆积在队列中等引擎恢复后继续消费。前端展示Vue.js ECharts为了给内部用户提供直观的仪表盘前端选用Vue.js框架因其灵活轻量生态丰富。数据可视化采用百度的ECharts它图表类型全文档友好能够轻松实现时间趋势图、地域分布图、词云等复杂图表完美满足数据呈现的需求。注意关于“反爬”与合规的严肃讨论这是所有数据收集项目必须正视的问题。我们的核心策略是“友好爬取”严格遵守robots.txt设置合理的下载延迟如3-10秒随机使用轮换的User-Agent池并大量使用HEAD请求先探测页面状态。对于必须登录才能查看的数据我们直接放弃。记住项目的长期价值建立在合法合规的基础上任何试图绕过正常访问限制的行为都风险极高得不偿失。3. 核心模块实现与实操细节3.1 分布式采集器集群的搭建采集器是系统的触手。我们采用Scrapy-Redis组件将Scrapy改造成了分布式爬虫。架构是这样的一个主节点Master运行Redis服务负责存储待爬取的URL队列Request Queue和去重集合Dupe Filter。多个从节点Slave运行Scrapy爬虫进程它们都连接到同一个Redis实例从公共队列中获取任务。具体搭建步骤环境准备所有节点安装Python、Scrapy、scrapy-redis库。pip install scrapy scrapy-redis修改Scrapy项目配置在项目的settings.py中进行关键配置。# 使用scrapy_redis的调度器 SCHEDULER scrapy_redis.scheduler.Scheduler # 使用scrapy_redis的去重组件 DUPEFILTER_CLASS scrapy_redis.dupefilter.RFPDupeFilter # 允许暂停/恢复爬虫Redis会记录状态 SCHEDULER_PERSIST True # 指定Redis连接信息 REDIS_HOST 你的Redis主节点IP REDIS_PORT 6379 # 可选的Redis密码 # REDIS_PARAMS {password: your_password}编写爬虫Spider这里以抓取一个新闻网站为例关键是要继承RedisSpider。import scrapy from scrapy_redis.spiders import RedisSpider class APACNewsSpider(RedisSpider): name apac_news # 爬虫名称 redis_key apac_news:start_urls # 从Redis的这个key中读取起始URL def parse(self, response): # 解析列表页提取文章链接和翻页链接 article_links response.css(div.article-list a::attr(href)).getall() for link in article_links: yield response.follow(link, callbackself.parse_article) # 翻页 next_page response.css(a.next-page::attr(href)).get() if next_page: yield response.follow(next_page, callbackself.parse) def parse_article(self, response): # 解析文章详情页 item {} item[title] response.css(h1.article-title::text).get().strip() item[publish_time] response.css(time::attr(datetime)).get() item[content] .join(response.css(article p::text).getall()) item[source_url] response.url item[region] self.determine_region(response) # 根据URL或内容判断地区 yield item def determine_region(self, response): # 一个简单的地区判断逻辑示例 url response.url if .jp in url: return Japan elif .kr in url: return Korea # ... 其他逻辑 return Asia-Pacific启动与任务投放首先在所有从节点运行爬虫命令scrapy crawl apac_news。此时爬虫会等待任务。然后在主节点向Redis的apac_news:start_urls列表类型为List中放入起始URL例如各国新闻网站的主页或栏目页。# 使用redis-cli redis-cli lpush apac_news:start_urls https://example-news.jp/economy redis-cli lpush apac_news:start_urls https://example-news.kr/business实操心得连接池管理大量爬虫节点同时连接Redis务必在Scrapy配置中启用CONCURRENT_REQUESTS控制并发并在Redis服务器端调整maxclients参数避免连接数耗尽。监控队列长度时刻关注Redis中apac_news:start_urls待爬队列和apac_news:dupefilter去重集合的大小。如果待爬队列空了爬虫会空闲如果去重集合过大会影响性能需要定期清理或使用布隆过滤器优化。节点差异化配置可以为不同地区的节点配置不同的代理IP池和User-Agent列表使其行为更贴近当地真实用户。3.2 基于消息队列的数据清洗流水线采集器吐出来的是原始的、脏数据我们的清洗引擎负责把它们变成干净的结构化信息。清洗逻辑通过订阅RabbitMQ队列来触发。清洗引擎的核心工作流消费消息清洗服务作为一个常驻进程使用Pika库RabbitMQ的Python客户端监听名为raw_html_queue的队列。HTML解析与正文提取使用lxml或html2text进行解析。但更关键的是正文提取。我们对比了readability-lxml、newspaper3k和goose3等库最终选择newspaper3k作为主力因为它对新闻类网站的正文、标题、发布时间提取效果较好且支持多语言。from newspaper import Article def extract_article(html, url): article Article(url) article.download(input_htmlhtml) article.parse() return { title: article.title, text: article.text, publish_date: article.publish_date, authors: article.authors, top_image: article.top_image }结构化信息抽取这是提升数据价值的关键。我们针对不同行业数据编写了特定的抽取规则。正则表达式用于抽取明确的模式如金额“$1.2B”, “¥50亿”、百分比、日期。命名实体识别NER使用spaCy库预训练的多语言模型识别文本中的人名、组织名、地名、产品名等。自定义规则与词典对于领域内特定术语如“5G基站”、“新能源汽车销量”我们维护了词典通过关键词匹配和上下文分析进行抽取。数据标准化与入库将抽取的信息映射到预定义的数据库Schema中。例如将所有货币统一为美元所有日期转换为ISO 8601格式将公司名与内部的标准名称库进行匹配。最后将结构化数据写入PostgreSQL并将用于搜索的字段标题、正文、摘要索引到Elasticsearch。一个清洗任务的代码骨架示例import pika import json from newspaper import Article import spacy import re # 加载spaCy模型 nlp spacy.load(en_core_web_sm) # 可根据需要加载多语言模型 def callback(ch, method, properties, body): message json.loads(body) url message[url] html message[html] source message[source] try: # 1. 提取正文 article_info extract_article(html, url) # 2. NER识别 doc nlp(article_info[text][:100000]) # 处理前10万字符防止过长 entities [(ent.text, ent.label_) for ent in doc.ents] # 3. 自定义规则抽取示例抽取市场规模数字 market_size_pattern r市场规模[约]?\s*([$¥€£]?\d(?:\.\d)?[亿万]?[美元欧元日元人民币]?) market_size_matches re.findall(market_size_pattern, article_info[text]) # 4. 构建最终数据对象 clean_data { source_url: url, title: article_info[title], content: article_info[text], publish_date: article_info[publish_date], entities: entities, extracted_data: { market_size: market_size_matches[0] if market_size_matches else None, # ... 其他抽取字段 }, region: infer_region_from_source(source), collection_time: message[collection_time] } # 5. 调用函数存入数据库和ES save_to_postgresql(clean_data) index_to_elasticsearch(clean_data) # 6. 确认消息消费成功 ch.basic_ack(delivery_tagmethod.delivery_tag) except Exception as e: print(fError processing {url}: {e}) # 可以选择将失败消息重新入队或放入死信队列供后续排查 ch.basic_nack(delivery_tagmethod.delivery_tag, requeueFalse) # 不重新入队避免死循环 # 连接RabbitMQ并开始消费 connection pika.BlockingConnection(pika.ConnectionParameters(localhost)) channel connection.channel() channel.queue_declare(queueraw_html_queue, durableTrue) channel.basic_consume(queueraw_html_queue, on_message_callbackcallback) channel.start_consuming()3.3 数据存储与检索设计PostgreSQL表设计要点我们设计了核心的articles表包含id,title,clean_content,publish_date,source_url,region,country,industry_tags,extracted_data(JSONB字段)created_at等字段。extracted_data字段灵活存储了清洗引擎抽出的结构化信息如{market_size: 100B USD, mentioned_companies: [Company A, Company B]}。这种设计平衡了结构的稳定性和需求的灵活性。Elasticsearch索引映射与优化在ES中我们更关注搜索和分析性能。映射Mapping设计时对title和clean_content字段使用text类型并进行分词我们采用了IK分词器处理中文默认标准分词器处理英文。对publish_date,region,country等字段使用keyword类型用于精确过滤和聚合。为了提高复杂查询的响应速度我们充分利用了ES的聚合Aggregation功能。例如用户想查看“日本2021年Q3关于半导体行业的报道中提及频率最高的公司”这个查询可以通过bool查询结合filter再对extracted_data.mentioned_companies字段进行terms聚合轻松实现。实操心得数据一致性这是一个分布式系统常见问题。我们的方案是清洗引擎在成功写入PostgreSQL后立即写入Elasticsearch。如果ES写入失败则记录日志并告警同时该条数据在PostgreSQL中标记为“未索引”。我们有一个定期的补偿任务扫描这些“未索引”记录并重试。虽然这不是严格的分布式事务但保证了最终一致性在实践中足够可靠。4. 前端展示与用户交互实现前端的目标是让非技术同事也能轻松查询和洞察数据。我们基于Vue.js和Element UI组件库快速搭建了管理后台。核心功能页面数据仪表盘首页是一个综合仪表盘使用ECharts展示关键指标的趋势图、地域分布热力图、热门话题词云。数据通过调用后端提供的聚合API后端从ES中查询获得。高级搜索页面这是使用频率最高的功能。提供了多条件筛选关键词全文搜索支持AND/OR逻辑按地区、国家、行业标签筛选按时间范围筛选按数据源筛选 搜索结果以列表形式展示支持按相关度或时间排序。列表的每一项可以展开查看详情包括原文链接、提取的结构化数据和高亮显示的关键词。数据源管理页面供管理员查看各数据源爬虫的运行状态、最近采集时间、成功/失败次数并可以手动触发或停止某个采集任务。前端与后端的交互前端通过Axios库调用后端RESTful API。后端API层使用Python的FastAPI框架开发它性能好自动生成交互式文档非常适合前后端分离的项目。一个典型的搜索API如下from fastapi import FastAPI, Query from typing import Optional from datetime import date app FastAPI() app.get(/api/search) async def search_articles( q: Optional[str] Query(None, description搜索关键词), region: Optional[str] Query(None, description地区), start_date: Optional[date] Query(None, description开始日期), end_date: Optional[date] Query(None, description结束日期), page: int Query(1, ge1), size: int Query(20, ge1, le100) ): # 构建Elasticsearch查询DSL query_body build_es_query(q, region, start_date, end_date, page, size) # 执行ES查询 es_response await es_client.search(indexarticles, bodyquery_body) # 处理并返回结果 return format_response(es_response)性能优化点API缓存对于仪表盘的聚合数据变化不频繁我们在后端使用Redis进行了缓存设置5-30分钟不等的过期时间大幅降低了ES和数据库的压力。前端虚拟滚动当搜索结果很多时列表渲染会卡顿。我们使用了虚拟滚动组件只渲染可视区域内的DOM元素流畅度提升显著。图表按需加载ECharts图表在组件mounted后再初始化并且监听视口变化只有图表进入视口时才加载数据。5. 部署、监控与运维实践5.1 容器化部署与编排我们使用Docker将每个组件采集器、清洗引擎、Web API、前端都容器化然后通过Docker Compose在测试环境一键部署。在生产环境则使用了Kubernetes进行编排。Dockerfile示例清洗引擎FROM python:3.9-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt -i https://pypi.tuna.tsinghua.edu.cn/simple COPY . . CMD [python, clean_worker.py]Kubernetes关键配置Deployment为每个服务定义Deployment设置合适的副本数如清洗引擎可以多副本并行消费。ConfigMap Secret将数据库连接字符串、API密钥、RabbitMQ地址等配置信息通过ConfigMap和Secret管理与镜像解耦。Service为前端和后端API创建Service提供稳定的内部访问端点。Ingress通过Ingress配置域名和SSL证书将前端服务暴露给外部用户访问。5.2 系统监控与日志收集一个看不见的系统是危险的。我们建立了多层监控基础设施监控使用Prometheus Grafana监控服务器Node Exporter和容器cAdvisor的CPU、内存、磁盘、网络指标。应用性能监控APM在关键Python服务中集成Sentry捕获运行时错误和异常。同时使用自定义的指标如每分钟处理的消息数、清洗成功率暴露给Prometheus。日志集中化所有容器的日志都通过Fluentd收集统一发送到Elasticsearch集群与业务ES隔离并通过Kibana进行查看和搜索。这为排查问题提供了巨大便利。业务健康度检查编写了一个简单的健康检查脚本定时访问网站首页、检查数据库连接、检查消息队列堆积情况并通过企业微信机器人发送日报。5.3 数据质量保障与更新策略数据项目质量是生命线。我们建立了几个机制采样复核每天随机采样1%已清洗的数据由人工复核标题、正文提取、关键信息抽取的准确性发现问题则回溯调整清洗规则。数据源健康度评分为每个数据源定义评分规则如可访问性、结构稳定性、内容更新频率自动计算分数。低分数据源会触发告警提示可能需要维护爬虫规则。增量更新与全量回溯大部分数据源采用增量采集只抓取新内容。但每周会对核心数据源进行一次全量采集在凌晨低峰期用于纠正可能因网站改版导致的增量采集遗漏并刷新去重集合。6. 典型问题排查与优化记录在项目运行中我们遇到了形形色色的问题这里记录几个最有代表性的。问题一爬虫被目标网站封禁IP。现象某个国家站点的采集任务连续失败返回403或429状态码。排查查看该爬虫节点的日志发现请求头中的User-Agent比较固定。检查目标网站的robots.txt确认没有禁止爬取目标路径。解决方案降低请求频率将Scrapy的DOWNLOAD_DELAY从2秒提高到5-10秒并加入随机延迟 (RANDOMIZE_DOWNLOAD_DELAY True)。丰富User-Agent池准备一个包含几十个主流浏览器标识的列表在中间件中随机选取。引入代理IP池对于反爬特别严格的站点使用付费的代理IP服务并在Scrapy中配置代理中间件实现IP轮换。设置自动恢复机制在Celery任务中捕获403/429异常将任务标记为失败并放入重试队列延迟一段时间如1小时后重试。问题二消息队列RabbitMQ消息堆积清洗延迟。现象Kibana监控显示raw_html_queue队列长度持续增长清洗引擎处理速度跟不上采集速度。排查检查清洗引擎的日志和资源监控发现CPU使用率不高但单条消息处理耗时很长10秒。分析耗时长的原因发现是NER模型spaCy加载和初始化在每次处理消息时都在进行并且处理超长文本时效率低下。解决方案预热与复用模型在清洗服务启动时就加载好所需的NLP模型并在整个进程生命周期内复用避免每次处理都重新加载。文本长度截断对于明显过长的文本如超过10万字符在提取正文后进行智能截断如取前N个字符或根据段落截断保证核心内容不丢失的同时提升处理速度。水平扩展最简单有效的方法增加清洗引擎的Pod副本数让多个消费者并行处理队列消息。问题三Elasticsearch查询超时。现象前端进行复杂聚合查询如按月份、地区、行业多维度分组统计时接口经常超时30秒。排查在Kibana的Dev Tools中执行同样的查询DSL发现确实很慢。使用_search?request_cachetrue查看是否命中缓存并使用Profile API分析查询各个阶段的耗时。解决方案优化索引映射确认用于聚合的字段如region,publish_date类型为keyword或date而不是text。text字段聚合性能极差。使用预聚合对于仪表盘上那些需要实时性不高但查询复杂的图表如“过去一年各区域每周发文量趋势”我们在后端使用Celery定时任务提前计算好结果存入Redis缓存。前端查询时直接读缓存。限制查询范围在前端界面和API层强制用户选择时间范围避免无限制的全量查询。默认只查询最近3个月的数据。升级硬件与调整配置适当增加ES集群的内存调整JVM堆大小并确保分片数量设置合理不是越多越好。问题四数据重复入库。现象数据库中出现标题和内容高度相似但来源URL略有不同的记录。排查检查Scrapy-Redis的去重规则发现默认是基于请求URL的指纹fingerprint去重。但有些网站同一篇文章可能有多个URL如带不同参数或者转载其他网站的内容。解决方案实现内容级别的去重。Simhash算法对清洗后的正文内容计算Simhash指纹。Simhash对细微改动如换行、个别字词不敏感适合检测相似内容。布隆过滤器Bloom Filter将Simhash值存入一个布隆过滤器。新的文章在入库前先计算Simhash然后查询布隆过滤器。如果可能存在相似文章则再进行一次精确的数据库查询比对如计算编辑距离或Jaccard相似度如果布隆过滤器说肯定不存在则直接入库。这在大数据量下能极大提升去重效率。去重时机在清洗引擎处理完数据准备写入数据库之前进行Simhash计算和比对。这个项目从构想到稳定运行花了团队近三个月的时间。回头看最大的收获不是做出了一个能跑的系统而是建立了一套应对“大规模、多源、异构公开数据收集与处理”的方法论和工程实践。技术细节会过时但解决问题的思路——分而治之、松耦合设计、重视监控和数据质量——是长期有效的。如果你也在筹划类似的项目希望这些踩过的坑和总结的经验能帮你少走些弯路。最后一个小建议在项目初期不要追求大而全先聚焦核心数据源跑通最小闭环End-to-End Pipeline快速验证数据价值和系统可行性之后再逐步扩展和优化这样更容易获得支持和持续投入。
返回列表