ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

京东商品数据采集到分析可视化:Selenium+Django+Spark+Hadoop全链路实战

京东商品数据采集到分析可视化:Selenium+Django+Spark+Hadoop全链路实战 做电商数据采集分析这个方向说难不难说简单也绝对不简单。市面上的教程大多只讲了某一个环节要么教你怎么爬要么教你怎么画图表很少能把这整条链路串起来讲透。这次要分享的这套京东商品采集与分析平台正好把“采集—存储—计算—展示”完整拉通了而且技术栈选得特别接地气——selenium爬虫做数据入口Django做可视化后端Spark和Hadoop负责数据仓库和分析计算。无论你是想学爬虫、刚接触大数据生态还是想把数据分析的课程项目做成一个能演示、有说服力的毕业设计或简历项目这套源码都能提供一个相当完整的参考。下面我把这个项目的实现思路、核心代码逻辑以及我在实操过程中踩过的坑按我的理解依次拆开讲。1. 项目整体设计与技术选型1.1 这个项目到底解决什么问题很多人以为做电商数据分析不就是爬个数据画几个图其实真正落地的时候你会发现事情远没有这么简单。单说数据量一个商品品类的搜索结果页就有几千条商品每条商品又有标题、价格、店铺、评论数、好评率、销量等多个字段如果再把不同关键词、不同时间点的采集结果叠加起来单机用Pandas处理就已经很吃力了。所以这个项目的第一个核心诉求是把采集到的数据存进可扩展的数据仓库再用分布式计算引擎去做统计最后通过一个可交互的可视化页面把结论呈现出来。项目的价值就在这“一条龙”上。它覆盖了从数据生产端到消费端的完整链路selenium负责把京东这类动态加载页面的商品数据抓下来数据经过清洗后进入数据仓库Hive承担表结构的组织Spark负责跑复杂的分析任务Django在Web端把分析结果渲染成图表。你去看市面上大多数课程项目和开源代码很少有同时覆盖这四层的大部分是爬虫和Django两两组合或者Spark和Hive纯离线处理不涉及网页展示。1.2 技术栈选择的逻辑选技术栈不是一个一个比谁名气大而是要看每层任务的实际需求。这个项目的设计者选型我觉得很合理我把各层任务和对应技术的关系整理了一张表项目层级核心任务技术选型选择理由对比其他方案数据采集获取京东商品列表页和详情页的动态数据selenium ChromeDriver京东页面高度依赖JavaScript渲染requests直接抓取拿不到完整结构selenium可以模拟真实浏览器行为数据落地原始数据保存 结构化存储CSV MySQL HDFSCSV方便抽样检查MySQL适合管理维度表HDFS是海量原始数据的最终归宿数据仓库统一数据模型分层管理数据Hive基于HadoopHive把SQL能力移植到分布式存储上可以复用大量SQL经验比纯MapReduce开发效率高太多数据分析统计商品价格、评论、店铺分布等指标Spark SQL DataFrame内存计算比MapReduce快得多尤其适合多轮迭代的统计分析任务可视化展示图表展示分析结果支持交互Django ECharts前端模板Django自带ORM和模板系统和MySQL、Hive结果集能快速衔接开发Web端最快可以看到这套组合里没有多余的技术每一层都能找到明确的不可替代性。唯一可以商榷的地方是爬虫这块如果你只采集几百条数据用requests加解析库确实更轻但selenium的优势在于稳定面对复杂页面结构时不用频繁断点调试拿不到数据的页面。再说到Spark和Hadoop可能有人觉得一个课程级别的项目用不到这么重的框架但恰恰是这两个组件把项目的技术档次从“入门爬虫脚本”拉升到了“生产级数据平台”的架构水平面试的时候能聊的东西也就多了。2. 数据采集层Selenium爬虫的落地细节2.1 为什么选了Selenium而不直接上requests我接触过不少朋友一上来习惯性就写requests结果发现京东商品列表页返回的HTML里根本没有商品数据这是因为数据是通过AJAX异步加载的而且部分字段还是混在script变量里的。这时候再回去补selenium就比较费时了。selenium的做法是直接驱动一个真实的浏览器内核去渲染页面等所有网络请求都完成、DOM稳定之后再去提取信息从原理上绕开了异步加载的坑。selenium的适用场景有几个典型特征页面大量依赖JS渲染、需要模拟点击和滚动操作、反爬策略较严格需要真实浏览器指纹。京东恰好一个不落全占了。所以选它属于“对症下药”。2.2 环境配置与ChromeDriver版本坑先讲环境准备这块是最容易出问题的地方。你需要装好Python 3.8以上版本然后安装两个关键的库pip install selenium pip install webdriver-manager第二个库强烈建议装它能自动匹配当前Chrome浏览器对应的驱动版本省去了手动去下载chromedriver的麻烦。以前我不小心在Windows上配了个Linux版本的chromedriver折腾了半小时才发现是系统不匹配。用webdriver-manager之后代码里只需要写from selenium import webdriver from selenium.webdriver.chrome.service import Service from webdriver_manager.chrome import ChromeDriverManager options webdriver.ChromeOptions() options.add_argument(--headless) # 无头模式不弹出浏览器窗口 options.add_argument(--disable-gpu) options.add_argument(--no-sandbox) options.add_argument(--disable-blink-featuresAutomationControlled) # 隐藏自动化特征 driver webdriver.Chrome( serviceService(ChromeDriverManager().install()), optionsoptions )这里有一个很实用的参数要展开说--disable-blink-featuresAutomationControlled。如果不加这个参数网站很容易通过navigator.webdriver属性识别出你是自动化脚本然后直接拒绝返回数据。加上它之后浏览器会隐藏掉webdriver的标识被拦截的几率大幅下降。2.3 采集流程与核心代码实现采集的逻辑我拆成四步搜索关键词、渲染等待、滚动加载、解析保存。第一步是打开京东搜索页并输入关键词注意在输入之前要给搜索框一个显式等待等页面元素真正加载出来再操作否则可能会报ElementNotInteractableExceptionfrom selenium.webdriver.common.by import By from selenium.webdriver.support.ui import WebDriverWait from selenium.webdriver.support import expected_conditions as EC search_url https://search.jd.com/Search?keyword蓝牙耳机 driver.get(search_url) # 显式等待最多等待10秒直到搜索框和搜索按钮都出现 wait WebDriverWait(driver, 10) wait.until(EC.presence_of_element_located((By.ID, key)))第二步是关键点京东的商品列表不是一次性渲染完的而是随着滚动条下拉逐屏加载。如果不做滚动处理你只能拿到前二三十条数据。常见的做法是模拟滚轮动作import time for i in range(8): # 模拟滚动到页面底部触发懒加载 driver.execute_script(window.scrollTo(0, document.body.scrollHeight);) time.sleep(2) # 等待新内容渲染滚动频率要控制好太快会导致后续脚本被识别太慢几十页商品光等滚动就要十分钟。实测下来每两秒一次是比较平衡的节奏。第三步是定位商品条目并提取字段。京东搜索结果页每一个商品项都在li.gl-item节点下通过CSS选择器能很稳定地定位到。核心解析逻辑我封装成了一个函数def parse_products(driver): items driver.find_elements(By.CSS_SELECTOR, li.gl-item) products [] for item in items: try: title item.find_element(By.CSS_SELECTOR, .p-name em).text.strip() price item.find_element(By.CSS_SELECTOR, .p-price strong i).text shop item.find_element(By.CSS_SELECTOR, .p-shop a).text.strip() comment item.find_element(By.CSS_SELECTOR, .p-commit a).text.strip() products.append({ title: title, price: float(price) if price else 0.0, shop: shop, comment_count: comment, keyword: current_keyword, crawl_time: time.strftime(%Y-%m-%d %H:%M:%S), }) except Exception as e: # 个别商品可能缺少某个字段跳过而不是中断整个任务 continue return products这里做了一个很重要的容错处理单个商品解析失败时用continue跳过而不是让整个爬虫崩溃。实际操作中我发现京东偶尔会在列表里插入一些促销位或广告卡片它们和普通商品的DOM结构不一样解析的时候必然会抛异常这种情况下直接跳过是最合理的方案。最后是保存策略。我的建议是原始数据先落一份CSV方便自己随时抽查数据质量然后再写入MySQL。CSV的文件名记得带上时间戳避免重复运行覆盖了历史数据import csv with open(fproducts_{time.strftime(%Y%m%d_%H%M%S)}.csv, w, newline, encodingutf-8-sig) as f: writer csv.DictWriter(f, fieldnames[title, price, shop, comment_count, keyword, crawl_time]) writer.writeheader() writer.writerows(products)2.4 反爬应对与数据质量保障说句实在话爬虫这门技术必须在合规的框架内使用。我在这个项目里遵循几条原则只采集公开的商品列表和详情信息不碰用户个人数据控制请求频率不做高并发轮询采集逻辑保持单线程顺序执行抓到数据就走。除了浏览器自动化特征伪装还有几个细节对数据质量和账号安全影响很大。第一个是避免同一个关键词短时间内反复翻页连续采集3到4页之后最好sleep 8到10秒给页面和服务端一个缓冲第二个是随机化浏览器窗口尺寸不同的窗口宽高会影响页面渲染出的数据条数这个会影响采集结果的稳定性第三个是如果采集过程中发现driver响应超时不要无脑重试先检查一下网络和页面结构是否有变化。数据预处理方面我最常遇到的问题是价格字段格式不统一有的页面拿到的价格是“199.00”有的直接是字符串“199.00”清洗时统一做strip和类型转换就好。还有评论数有时候是“52万”“2.3万”这种缩写需要写个正则或者判断逻辑统一转成整型数字import re def parse_comment_count(text): 把 2.3万 / 5200 / 1200 统一转为整数 if not text: return 0 text text.replace(, ).strip() if 万 in text: return int(float(text.replace(万, )) * 10000) try: return int(text) except ValueError: return 03. 数据存储与仓库建设从MySQL到Hive3.1 数据分层设计思路这个项目的数据仓库设计保持了经典的ODS层、DWD层、ADS层三层架构。可能有人觉得课程项目没必要搞这么复杂但既然引入了Hive和Hadoop那就应该按数据仓库的规范来做否则Spark没有用武之地整个项目也就名不副实了。三层架构各司其职ODS层原始数据层爬虫采集到的数据原样落地不做过多的清洗和改动主要用途是保留原始事实方便后续核对。DWD层明细数据层对ODS层做清洗和规范化比如统一价格单位、格式化时间字段、去除缺失值严重的记录保证数据质量。ADS层应用数据层面向具体分析主题的汇总结果比如“各价格区间的商品数量”“评论数TOP10的商品”“店铺商品数排行”等Django可视化页面直接读这一层的数据。这套分层模型在真实企业里几乎是通用的面试的时候如果能把这个逻辑讲清楚比单纯说“我用Hive建了个表”要有说服力得多。3.2 MySQL表结构与Hive建表语句数据采集下来之后先放进MySQL作为Django可视化部分快速查询的数据源之一同时作为离线数据进入HDFS的中转站。商品明细表的设计如下CREATE TABLE product_info ( id INT AUTO_INCREMENT PRIMARY KEY, title VARCHAR(255) NOT NULL, price DOUBLE, shop VARCHAR(128), comment_count INT, keyword VARCHAR(64), crawl_time DATETIME, INDEX idx_price (price), INDEX idx_keyword (keyword) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;索引加到price和keyword上是有讲究的后面Spark分析完的结果也会写回MySQL的表Django端做范围查询和关键词筛选时没有索引很容易在百万级数据上把查询拖到秒级。Hive端建表则指向HDFS上的数据文件。我通常配合Sqoop或直接DataFrame写入的方式把数据传到HDFS然后建一张外部表挂载上去CREATE EXTERNAL TABLE ods_product_info ( title STRING, price DOUBLE, shop STRING, comment_count BIGINT, keyword STRING, crawl_time TIMESTAMP ) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION /user/hive/warehouse/ods/product_info;使用外部表的好处是删除表结构不会影响HDFS上的原始数据这在数据仓库初建阶段特别重要因为表结构肯定要反复迭代内部表一旦删了数据也跟着没了外部表就没有这个顾虑。3.3 数据由MySQL同步进HDFS的实测方案从MySQL到HDFS的同步我实测下来有两条路径比较顺手。一条是用Sqoop一条是用Spark的DataFrame直接读写。Sqoop的写法比较成熟sqoop import \ --connect jdbc:mysql://localhost:3306/jd_spider \ --username root \ --password 123456 \ --table product_info \ --target-dir /user/hive/warehouse/ods/product_info \ --fields-terminated-by , \ --m 1要注意--m 1这个参数如果数据量不是特别大强制使用单并行度就够了既不会对MySQL造成压力也不会在HDFS上产生太多小文件。Spark写入的方式则适合还要同时做一轮格式转换的场景。下面这个代码片段是把MySQL整表读出来直接把数据写进HDFSfrom pyspark.sql import SparkSession spark SparkSession.builder.appName(sync_mysql_to_hdfs).getOrCreate() df spark.read \ .format(jdbc) \ .option(url, jdbc:mysql://localhost:3306/jd_spider) \ .option(dbtable, product_info) \ .option(user, root) \ .option(password, 123456) \ .option(driver, com.mysql.cj.jdbc.Driver) \ .load() df.write.csv(/user/hive/warehouse/ods/product_info, modeoverwrite)两条路径各有用武之地Sqoop适合纯同步场景Spark适合“同步预处理”一肩挑的场景。我当时是先用Sqoop把历史数据补了一次后面增量数据统一走Spark因为Spark写完可以直接做一轮DWD层的清洗逻辑减少了一次任务调度。4. 数据分析层Spark处理核心逻辑4.1 为什么“点名”Spark而不是继续用Pandas数据分析任务放单机上用Pandas也能跑但一旦数据量上了百万级内存占用会迅速飙升而且很多聚合操作写起来远不如Spark SQL直观。更关键的是这个项目引入了Hadoop和HiveSpark天然能直接跑在HDFS上从Hive表读数据然后计算整个链路是无缝衔接的。从架构上看Spark的定位就是替代MapReduce做通用计算引擎。MapReduce虽然稳定但每个Map和Reduce阶段都要落盘读写磁盘迭代次数一多磁盘IO就成了瓶颈。Spark采用基于内存的RDD和DataFrame抽象shuffle过程中的中间结果优先放内存遇到需要反复迭代的统计任务性能优势极其明显。4.2 核心分析任务拆解与PySpark实现项目里的分析任务我设计成了四大类商品价格分布统计、品牌/店铺集中度分析、评论数与价格相关性分析、关键词效益对比。先看价格分布统计。这里用Spark SQL从DWD层读取数据按价格区间分桶聚合from pyspark.sql import SparkSession from pyspark.sql.functions import col, when spark SparkSession.builder \ .appName(jd_price_analysis) \ .config(spark.sql.warehouse.dir, hdfs://localhost:9000/user/hive/warehouse) \ .enableHiveSupport() \ .getOrCreate() df spark.sql(SELECT price, comment_count, shop, keyword FROM dwd_product_info) price_dist df.withColumn( price_range, when(col(price) 50, 0-50元) .when(col(price) 100, 50-100元) .when(col(price) 200, 100-200元) .when(col(price) 500, 200-500元) .otherwise(500元以上) ).groupBy(price_range).count().toPandas() print(price_dist)这段代码的核心思路是使用when构建分段逻辑再groupBy计数。之所以用withColumn新增一列而不是直接改price列是为了保留原始字段后续其他分析还要用到。生成的结果toPandas之后可以方便地转成JSON接口供Django读取。再看店铺集中度分析这个分析对电商选品很有价值可以看到一个品类是不是被少数头部卖家垄断shop_top df.groupBy(shop) \ .agg( count(*).alias(product_count), round(avg(price), 2).alias(avg_price), round(sum(comment_count), 0).alias(total_comments) ) \ .orderBy(col(product_count).desc()) \ .limit(20)这里有个容易被忽视的细节agg里的列别名最好显式加上不然Spark生成的列名是count(1)这种后面写SQL查询的时候还要额外处理列名映射非常麻烦。评论数与价格的相关性分析我直接用了Spark DataFrame的统计函数来做虽然Spark的corr函数只能两两计算但配合selectExpr可以一次算多个组合corr_result df.selectExpr( corr(price, comment_count) as price_comment_corr, corr(price, cast(comment_count as double)) as price_comment_corr2 ).collect()[0] print(价格与评论数相关系数:, corr_result[price_comment_corr])4.3 分析结果写回MySQL供Django读取Spark分析完的结果要展示到Web端最简单的做法是把聚合结果写入MySQL的一张结果表Django只管做查询和渲染。代码如下result_df.write \ .format(jdbc) \ .option(url, jdbc:mysql://localhost:3306/jd_spider) \ .option(dbtable, ads_price_range_stats) \ .option(user, root) \ .option(password, 123456) \ .option(driver, com.mysql.cj.jdbc.Driver) \ .mode(overwrite) \ .save()写结果表的时候我建议每个分析任务单独建一张ADS表不要所有结果混在一张表里。这样带来的好处是Django端查询逻辑简单清晰每个图表对应一个接口对应一张表出问题的时候定位特别快。这里还要提醒一个MySQL驱动包的问题com.mysql.cj.jdbc.Driver这个类名需要mysql-connector-java 8.0以上版本如果用老版本要换成com.mysql.jdbc.Driver。我一开始用pyspark提交任务因为驱动没放进$SPARK_HOME/jars目录一直报“ClassNotFound”排查了半天才发现是依赖缺失。5. Django可视化分析与前端展示5.1 Django项目结构与数据模型Django这部分在项目中承担的角色是后端服务和页面渲染。我们不需要在Django里部署Spark任务只需要把Spark写好的结果表取出来渲染成可视化图表。项目结构我大概切成这样jd_analysis_web/ ├── manage.py ├── analysis/ # 分析应用 │ ├── models.py # 数据模型 │ ├── views.py # 视图函数API 页面渲染 │ ├── urls.py # 路由配置 │ └── templates/ │ └── dashboard.html # 可视化大屏模板 ├── static/ │ ├── echarts.min.js │ └── style.css └── config/ ├── settings.py └── urls.pyDjango模型层直接映射MySQL里的ADS结果表比如价格分布表可以定义为# analysis/models.py from django.db import models class AdsPriceRangeStats(models.Model): price_range models.CharField(max_length32, verbose_name价格区间) cnt models.IntegerField(verbose_name商品数量) class Meta: db_table ads_price_range_stats verbose_name 价格区间统计注意Meta里设置db_table为Spark写入的物理表名Django的ORM虽然有自己的建表能力但在这种场景下我们期望让Django去读已有表而不是反向建表覆盖掉Spark写的结果。5.2 视图设计与API接口实现视图层我采用“页面渲染数据接口分离”的方案。页面路由返回HTML模板数据接口返回JSON格式的图表数据。这样前后端职责清晰后面如果要改成前后端分离项目数据接口直接复用即可。# analysis/views.py import json from django.shortcuts import render from django.http import JsonResponse from .models import AdsPriceRangeStats def dashboard(request): return render(request, dashboard.html) def price_range_api(request): 价格分布柱状图数据接口 rows AdsPriceRangeStats.objects.all().order_by(price_range) data { categories: [r.price_range for r in rows], values: [r.cnt for r in rows], } return JsonResponse(data)路由配置也比较简单# analysis/urls.py from django.urls import path from . import views urlpatterns [ path(, views.dashboard, namedashboard), path(api/price-range/, views.price_range_api, nameprice_range_api), ]5.3 ECharts图表展示与页面模板Django模板里使用ECharts的方式非常直接引入echarts.min.js之后通过fetch调用上面的API接口拿到数据后再初始化图表。这样一个页面可以同时挂多个图表每个图表对应一个独立的API图表之间互不影响谁挂了也不影响其它模块展示。!-- templates/dashboard.html -- !DOCTYPE html html langzh-CN head meta charsetUTF-8 title京东商品数据分析平台/title script src/static/echarts.min.js/script /head body div idprice-chart stylewidth: 100%; height: 400px;/div script fetch(/api/price-range/) .then(res res.json()) .then(data { const chart echarts.init(document.getElementById(price-chart)); chart.setOption({ title: { text: 商品价格区间分布 }, tooltip: {}, xAxis: { data: data.categories }, yAxis: {}, series: [{ name: 商品数量, type: bar, data: data.values }] }); }); /script /body /htmlECharts的图表类型选择也有讲究价格分布这种有序的连续型数据用柱状图最直观店铺集中度这种带排序和占比的用饼图或者横向条形图更好关键词效益对比不同关键词的平均价格和总评论数则适合双轴折线图加柱状图混合展示。图表选得好整个可视化大屏的观感会提升一个档次。6. 常见问题与排查技巧实录6.1 Selenium采集不到数据的三个原因采集不到数据是这个项目里出现频率最高的问题我总结下来多半是这三个原因之一。第一个是页面元素定位失效京东前端改版过几次li.gl-item这个选择器目前是稳定的但如果哪天页面改版需要重新审视页面结构F12打开开发者工具重新定位元素。第二个是滚轮滑动速度过快页面还没来得及把后续商品渲染出来就执行了解析逻辑自然拿不全。解决办法是每滚动一次就秒级延迟并且滚动周期和解析周期错开。第三个是自动化特征被识别返回的HTML内容和正常浏览器不一致。这种情况要么页面弹验证码要么返回空页面最有效的应对方案是加载页面后先判断关键元素是否存在如果不存在主动刷新页面重试一次。6.2 Spark任务内存溢出的实际排查Spark跑大数据量的聚合任务内存溢出是最常见的报错。我在这个项目里有两次遇到java.lang.OutOfMemoryError的实战教训。第一次是因为HDFS文件块数量太碎一堆几KB的小文件导致读取时task数量爆炸。解决办法是生成HDFS文件之后做一次合并df.coalesce(4).write.csv(/user/hive/warehouse/dwd/product_info)coalesce(4)把输出文件控制在4个以内大幅降低后续读取的task数量。第二次是executor内存设置太小默认1G完全不够跑百万级数据。我在提交命令里做了调整spark-submit \ --master local[*] \ --executor-memory 4g \ --driver-memory 2g \ analyze_jd.py如果你在集群模式下跑还要注意spark.executor.cores和spark.executor.instances的配合。我当时在本地学习机上的经验是内存参数宁高勿低但也不要一台8G内存的机器就分配6G给executor要留出系统和其他服务的内存余量。6.3 中文乱码与数据库时区问题中文乱码这个问题在数据落进MySQL时最容易出现。我用了charsetutf8mb4建表之后正常写入中文没问题。但如果你的Python连接串里没有指定charset参数写入时依然可能默认用latin1导致乱码。正确姿势是在连接串里显式加上import pymysql conn pymysql.connect( hostlocalhost, userroot, password123456, databasejd_spider, charsetutf8mb4 )时区问题则主要出在JDBC连接上。如果你看到写入MySQL的时间字段少了8个小时多半是因为MySQL连接串里缺了serverTimezoneAsia/Shanghai加上就好url jdbc:mysql://localhost:3306/jd_spider?useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf86.4 常见问题速查表问题现象可能原因解决方案selenium找不到商品元素页面未加载完 / 页面结构改版增加显式等待检查CSS选择器采集到的价格为空字符串页面价格异步加载较慢滚动后多等1-2秒再解析Spark写入MySQL报ClassNotFoundJDBC驱动未放入SPARK_HOME/jars下载mysql-connector-java.jar放进jarsDjango页面图表空白API接口报错 / ECharts路径错误浏览器F12看请求状态码Hive表数据量对不上MySQLHDFS有重复文件用dfs -ls检查增量同步是否重复追加写在最后的话这套京东商品采集与分析平台我是从年初开始动手写的前前后后迭代了差不多两个月。最大的感触是项目里每个单独的技术点拿出来都有大量现成教程但真正把它们串到一条链路上遇到的问题是教程里没写过的。比如MySQL驱动版本不匹配、HDFS小文件导致Spark卡顿、Django读取中文表名编码报错等等这些才是项目经验最值钱的部分。如果你正在做一个大数据方向的课程设计或者毕设强烈建议不只是跑通代码而是把每一层的数据流画出来把每个技术组件“为什么存在”搞清楚面试时这就是你区别于其他人的亮点。最后再分享一个小技巧跑通整条链路后把Spark分析任务的调度脚本写好用crontab定期执行一次这样你的数据仓库会自动更新Django展示的图表也会跟着刷新整个平台就会呈现出“活”的状态演示效果拔群。
返回列表